{"version":3,"sources":["../../src/streaming/reader.ts"],"sourcesContent":["import Definitions from '../core/definitions';\r\nimport Schema from '../schema/schema';\r\nimport parse from '../parser/index';\r\nimport { ChunkDecoder, normalizeNewlines, splitLinesKeepRemainder, updateStringState } from './text';\r\nimport { toAsyncIterable } from './source';\r\nimport { IOStreamSource, StreamReaderOptions, StreamItem } from './types';\r\nimport IODocument from '../core/document';\r\n\r\n/**\r\n * A streaming reader for Internet Object data.\r\n * Reads chunked data from a source and parses it into IO objects.\r\n */\r\nexport class IOStreamReader implements AsyncIterable<StreamItem> {\r\n  private readonly source: AsyncIterable<any>;\r\n  private readonly definitions: Definitions | null;\r\n  private readonly options: StreamReaderOptions;\r\n\r\n  constructor(source: IOStreamSource, definitions?: Definitions | null, options?: StreamReaderOptions) {\r\n    this.source = toAsyncIterable(source);\r\n    this.definitions = definitions ?? null;\r\n    this.options = options || {};\r\n  }\r\n\r\n  /**\r\n   * Reads all items from the stream and returns them as an array.\r\n   * WARNING: This buffers the entire stream into memory.\r\n   */\r\n  async collect(): Promise<StreamItem[]> {\r\n    const items: StreamItem[] = [];\r\n    for await (const item of this) {\r\n      items.push(item);\r\n    }\r\n    return items;\r\n  }\r\n\r\n  /**\r\n   * Iterates over the stream, yielding parsed items (data, error, schemaName, etc.)\r\n   */\r\n  async *[Symbol.asyncIterator](): AsyncIterator<StreamItem> {\r\n    const maxBufferedChars = this.options.maxBufferedChars ?? 2_000_000;\r\n\r\n    let streamIndex = 0;\r\n    let headerLines: string[] = [];\r\n    let headerDone = false;\r\n\r\n    let defs: Definitions | null = this.definitions;\r\n    let defaultSchemaName: string | undefined = this.options.defaultSchema;\r\n\r\n    // Data parsing state\r\n    let currentSectionHeaderLine: string | null = null; // e.g. '--- $order'\r\n    let pendingLines: string[] = [];\r\n    let pendingSize = 0;\r\n    let inString: string | null = null;\r\n    let remainder = '';\r\n    const decoder = new ChunkDecoder();\r\n\r\n    // Helper to flush pending lines as a parsed section\r\n    const flushPending = async (): Promise<StreamItem[]> => {\r\n      if (pendingLines.length === 0) return [];\r\n\r\n      const sectionText = [\r\n        // If no explicit section header has been seen, omit it (default section)\r\n        currentSectionHeaderLine ? `${currentSectionHeaderLine}\\n` : '',\r\n        pendingLines.join('\\n'),\r\n        '\\n'\r\n      ].join('');\r\n\r\n      pendingLines = [];\r\n      pendingSize = 0;\r\n\r\n      const errors: Error[] = [];\r\n      let doc: IODocument;\r\n      try {\r\n        // parse() will merge external defs into header.definitions\r\n        doc = defs ? parse(sectionText, defs, errors) : parse(sectionText, null, errors);\r\n      } catch (err: any) {\r\n        // If parse throws (e.g. syntax error), yield an error item\r\n        return [{\r\n          data: null,\r\n          schemaName: currentSectionHeaderLine ? (currentSectionHeaderLine.replace('---', '').trim() || '$schema') : (defaultSchemaName ?? '$schema'),\r\n          index: streamIndex++,\r\n          error: err\r\n        }];\r\n      }\r\n\r\n      const out: StreamItem[] = [];\r\n      const sections = doc.sections;\r\n\r\n      // If parse produced errors but no sections (or empty sections), yield error\r\n      if ((!sections || sections.length === 0) && errors.length > 0) {\r\n         return [{\r\n          data: null,\r\n          schemaName: currentSectionHeaderLine ? (currentSectionHeaderLine.replace('---', '').trim() || '$schema') : (defaultSchemaName ?? '$schema'),\r\n          index: streamIndex++,\r\n          error: errors[0]\r\n        }];\r\n      }\r\n\r\n      if (!sections || sections.length === 0) return out;\r\n\r\n      for (let si = 0; si < sections.length; si++) {\r\n        const section = sections.get(si);\r\n        if (!section) continue;\r\n\r\n        const schemaName = (section.schemaName\r\n          ? (section.schemaName.startsWith('$') ? section.schemaName : `$${section.schemaName}`)\r\n          : (defaultSchemaName ?? '$schema'));\r\n\r\n        const data = section.data as any;\r\n\r\n        // For collections, yield each item. For single objects, yield once.\r\n        if (data && typeof (data as any)[Symbol.iterator] === 'function' && typeof data.toJSON === 'function') {\r\n          // Likely IOCollection\r\n          let idxInSection = 0;\r\n          for (const item of data as any) {\r\n            if (item && (item as any).__error) {\r\n              out.push({\r\n                data: null,\r\n                schemaName,\r\n                index: streamIndex++,\r\n                error: item,\r\n              });\r\n            } else {\r\n              out.push({\r\n                data: item,\r\n                schemaName,\r\n                index: streamIndex++,\r\n                error: undefined,\r\n              });\r\n            }\r\n            idxInSection++;\r\n          }\r\n        } else if (data) {\r\n          if ((data as any).__error) {\r\n            out.push({ data: null, schemaName, index: streamIndex++, error: data });\r\n          } else {\r\n            out.push({ data, schemaName, index: streamIndex++, error: undefined });\r\n          }\r\n        }\r\n      }\r\n\r\n      // If parse produced errors, attach them to the last yielded item as a signal.\r\n      if (errors.length > 0 && out.length > 0) {\r\n        out[out.length - 1].error = errors[0];\r\n      }\r\n\r\n      return out;\r\n    };\r\n\r\n    // --- Main Loop ---\r\n    for await (const chunk of this.source) {\r\n      const text = normalizeNewlines(decoder.decode(chunk));\r\n      remainder += text;\r\n\r\n      if (remainder.length > maxBufferedChars) {\r\n        throw new Error(`Stream reader exceeded maxBufferedChars (${maxBufferedChars}).`);\r\n      }\r\n\r\n      const { lines, remainder: newRemainder } = splitLinesKeepRemainder(remainder);\r\n      remainder = newRemainder;\r\n\r\n      for (const rawLine of lines) {\r\n        const line = rawLine;\r\n        const trimmed = line.trim();\r\n\r\n        if (!headerDone) {\r\n            if (trimmed.startsWith('---')) {\r\n            headerDone = true;\r\n\r\n            const headerText = headerLines.length ? `${headerLines.join('\\n')}\\n---\\n` : `---\\n`;\r\n            headerLines = [];\r\n\r\n            const headerErrors: Error[] = [];\r\n            const headerDoc = defs ? parse(headerText, defs, headerErrors) : parse(headerText, null, headerErrors);\r\n\r\n            defs = headerDoc.header?.definitions ?? defs;\r\n            const schema = headerDoc.header?.schema;\r\n            if (!defaultSchemaName && schema instanceof Schema) {\r\n                defaultSchemaName = '$schema';\r\n            }\r\n\r\n            currentSectionHeaderLine = trimmed === '---' ? null : trimmed;\r\n            continue;\r\n            }\r\n\r\n            headerLines.push(line);\r\n            continue;\r\n        }\r\n\r\n        if (trimmed.startsWith('---')) {\r\n            const flushed = await flushPending();\r\n            for (const item of flushed) yield item;\r\n\r\n            currentSectionHeaderLine = trimmed === '---' ? null : trimmed;\r\n            inString = null;\r\n            continue;\r\n        }\r\n\r\n        if (trimmed.length === 0) {\r\n            if (inString !== null) {\r\n            pendingLines.push(line);\r\n            }\r\n            continue;\r\n        }\r\n\r\n        if (inString === null && (trimmed.startsWith('~') || trimmed.startsWith('#'))) {\r\n            const flushed = await flushPending();\r\n            for (const item of flushed) yield item;\r\n        }\r\n\r\n        inString = updateStringState(line, inString);\r\n\r\n        pendingLines.push(line);\r\n        pendingSize += line.length;\r\n        if (pendingSize > maxBufferedChars) {\r\n            throw new Error(`Stream reader exceeded maxBufferedChars (${maxBufferedChars}) in pending lines.`);\r\n        }\r\n      }\r\n    }\r\n\r\n    if (remainder.trim().length > 0) {\r\n      if (!headerDone) {\r\n        headerLines.push(remainder);\r\n      } else {\r\n        pendingLines.push(remainder);\r\n      }\r\n    }\r\n\r\n    if (!headerDone && headerLines.length > 0) {\r\n      pendingLines = headerLines;\r\n      headerLines = [];\r\n      headerDone = true;\r\n    }\r\n\r\n    const flushed = await flushPending();\r\n    for (const item of flushed) yield item;\r\n  }\r\n}\r\n\r\n/**\r\n * Creates a new IOStreamReader instance.\r\n * @param source The source to read from (string, Iterable, AsyncIterable, ReadableStream).\r\n * @param definitions Optional initial definitions.\r\n * @param options Optional strictness and buffer settings.\r\n */\r\nexport function createStreamReader(\r\n  source: IOStreamSource,\r\n  definitions?: Definitions | null,\r\n  options?: StreamReaderOptions\r\n): IOStreamReader {\r\n  return new IOStreamReader(source, definitions, options);\r\n}\r\n"],"mappings":";;;;;;;;;;;;;;;;;;;;;;;;;;;;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AACA,oBAAmB;AACnB,oBAAkB;AAClB,kBAA4F;AAC5F,oBAAgC;AAQzB,MAAM,eAAoD;AAAA,EAC9C;AAAA,EACA;AAAA,EACA;AAAA,EAEjB,YAAY,QAAwB,aAAkC,SAA+B;AACnG,SAAK,aAAS,+BAAgB,MAAM;AACpC,SAAK,cAAc,eAAe;AAClC,SAAK,UAAU,WAAW,CAAC;AAAA,EAC7B;AAAA;AAAA;AAAA;AAAA;AAAA,EAMA,MAAM,UAAiC;AACrC,UAAM,QAAsB,CAAC;AAC7B,qBAAiB,QAAQ,MAAM;AAC7B,YAAM,KAAK,IAAI;AAAA,IACjB;AACA,WAAO;AAAA,EACT;AAAA;AAAA;AAAA;AAAA,EAKA,QAAQ,OAAO,aAAa,IAA+B;AACzD,UAAM,mBAAmB,KAAK,QAAQ,oBAAoB;AAE1D,QAAI,cAAc;AAClB,QAAI,cAAwB,CAAC;AAC7B,QAAI,aAAa;AAEjB,QAAI,OAA2B,KAAK;AACpC,QAAI,oBAAwC,KAAK,QAAQ;AAGzD,QAAI,2BAA0C;AAC9C,QAAI,eAAyB,CAAC;AAC9B,QAAI,cAAc;AAClB,QAAI,WAA0B;AAC9B,QAAI,YAAY;AAChB,UAAM,UAAU,IAAI,yBAAa;AAGjC,UAAM,eAAe,YAAmC;AACtD,UAAI,aAAa,WAAW,EAAG,QAAO,CAAC;AAEvC,YAAM,cAAc;AAAA;AAAA,QAElB,2BAA2B,GAAG,wBAAwB;AAAA,IAAO;AAAA,QAC7D,aAAa,KAAK,IAAI;AAAA,QACtB;AAAA,MACF,EAAE,KAAK,EAAE;AAET,qBAAe,CAAC;AAChB,oBAAc;AAEd,YAAM,SAAkB,CAAC;AACzB,UAAI;AACJ,UAAI;AAEF,cAAM,WAAO,cAAAA,SAAM,aAAa,MAAM,MAAM,QAAI,cAAAA,SAAM,aAAa,MAAM,MAAM;AAAA,MACjF,SAAS,KAAU;AAEjB,eAAO,CAAC;AAAA,UACN,MAAM;AAAA,UACN,YAAY,2BAA4B,yBAAyB,QAAQ,OAAO,EAAE,EAAE,KAAK,KAAK,YAAc,qBAAqB;AAAA,UACjI,OAAO;AAAA,UACP,OAAO;AAAA,QACT,CAAC;AAAA,MACH;AAEA,YAAM,MAAoB,CAAC;AAC3B,YAAM,WAAW,IAAI;AAGrB,WAAK,CAAC,YAAY,SAAS,WAAW,MAAM,OAAO,SAAS,GAAG;AAC5D,eAAO,CAAC;AAAA,UACP,MAAM;AAAA,UACN,YAAY,2BAA4B,yBAAyB,QAAQ,OAAO,EAAE,EAAE,KAAK,KAAK,YAAc,qBAAqB;AAAA,UACjI,OAAO;AAAA,UACP,OAAO,OAAO,CAAC;AAAA,QACjB,CAAC;AAAA,MACH;AAEA,UAAI,CAAC,YAAY,SAAS,WAAW,EAAG,QAAO;AAE/C,eAAS,KAAK,GAAG,KAAK,SAAS,QAAQ,MAAM;AAC3C,cAAM,UAAU,SAAS,IAAI,EAAE;AAC/B,YAAI,CAAC,QAAS;AAEd,cAAM,aAAc,QAAQ,aACvB,QAAQ,WAAW,WAAW,GAAG,IAAI,QAAQ,aAAa,IAAI,QAAQ,UAAU,KAChF,qBAAqB;AAE1B,cAAM,OAAO,QAAQ;AAGrB,YAAI,QAAQ,OAAQ,KAAa,OAAO,QAAQ,MAAM,cAAc,OAAO,KAAK,WAAW,YAAY;AAErG,cAAI,eAAe;AACnB,qBAAW,QAAQ,MAAa;AAC9B,gBAAI,QAAS,KAAa,SAAS;AACjC,kBAAI,KAAK;AAAA,gBACP,MAAM;AAAA,gBACN;AAAA,gBACA,OAAO;AAAA,gBACP,OAAO;AAAA,cACT,CAAC;AAAA,YACH,OAAO;AACL,kBAAI,KAAK;AAAA,gBACP,MAAM;AAAA,gBACN;AAAA,gBACA,OAAO;AAAA,gBACP,OAAO;AAAA,cACT,CAAC;AAAA,YACH;AACA;AAAA,UACF;AAAA,QACF,WAAW,MAAM;AACf,cAAK,KAAa,SAAS;AACzB,gBAAI,KAAK,EAAE,MAAM,MAAM,YAAY,OAAO,eAAe,OAAO,KAAK,CAAC;AAAA,UACxE,OAAO;AACL,gBAAI,KAAK,EAAE,MAAM,YAAY,OAAO,eAAe,OAAO,OAAU,CAAC;AAAA,UACvE;AAAA,QACF;AAAA,MACF;AAGA,UAAI,OAAO,SAAS,KAAK,IAAI,SAAS,GAAG;AACvC,YAAI,IAAI,SAAS,CAAC,EAAE,QAAQ,OAAO,CAAC;AAAA,MACtC;AAEA,aAAO;AAAA,IACT;AAGA,qBAAiB,SAAS,KAAK,QAAQ;AACrC,YAAM,WAAO,+BAAkB,QAAQ,OAAO,KAAK,CAAC;AACpD,mBAAa;AAEb,UAAI,UAAU,SAAS,kBAAkB;AACvC,cAAM,IAAI,MAAM,4CAA4C,gBAAgB,IAAI;AAAA,MAClF;AAEA,YAAM,EAAE,OAAO,WAAW,aAAa,QAAI,qCAAwB,SAAS;AAC5E,kBAAY;AAEZ,iBAAW,WAAW,OAAO;AAC3B,cAAM,OAAO;AACb,cAAM,UAAU,KAAK,KAAK;AAE1B,YAAI,CAAC,YAAY;AACb,cAAI,QAAQ,WAAW,KAAK,GAAG;AAC/B,yBAAa;AAEb,kBAAM,aAAa,YAAY,SAAS,GAAG,YAAY,KAAK,IAAI,CAAC;AAAA;AAAA,IAAY;AAAA;AAC7E,0BAAc,CAAC;AAEf,kBAAM,eAAwB,CAAC;AAC/B,kBAAM,YAAY,WAAO,cAAAA,SAAM,YAAY,MAAM,YAAY,QAAI,cAAAA,SAAM,YAAY,MAAM,YAAY;AAErG,mBAAO,UAAU,QAAQ,eAAe;AACxC,kBAAM,SAAS,UAAU,QAAQ;AACjC,gBAAI,CAAC,qBAAqB,kBAAkB,cAAAC,SAAQ;AAChD,kCAAoB;AAAA,YACxB;AAEA,uCAA2B,YAAY,QAAQ,OAAO;AACtD;AAAA,UACA;AAEA,sBAAY,KAAK,IAAI;AACrB;AAAA,QACJ;AAEA,YAAI,QAAQ,WAAW,KAAK,GAAG;AAC3B,gBAAMC,WAAU,MAAM,aAAa;AACnC,qBAAW,QAAQA,SAAS,OAAM;AAElC,qCAA2B,YAAY,QAAQ,OAAO;AACtD,qBAAW;AACX;AAAA,QACJ;AAEA,YAAI,QAAQ,WAAW,GAAG;AACtB,cAAI,aAAa,MAAM;AACvB,yBAAa,KAAK,IAAI;AAAA,UACtB;AACA;AAAA,QACJ;AAEA,YAAI,aAAa,SAAS,QAAQ,WAAW,GAAG,KAAK,QAAQ,WAAW,GAAG,IAAI;AAC3E,gBAAMA,WAAU,MAAM,aAAa;AACnC,qBAAW,QAAQA,SAAS,OAAM;AAAA,QACtC;AAEA,uBAAW,+BAAkB,MAAM,QAAQ;AAE3C,qBAAa,KAAK,IAAI;AACtB,uBAAe,KAAK;AACpB,YAAI,cAAc,kBAAkB;AAChC,gBAAM,IAAI,MAAM,4CAA4C,gBAAgB,qBAAqB;AAAA,QACrG;AAAA,MACF;AAAA,IACF;AAEA,QAAI,UAAU,KAAK,EAAE,SAAS,GAAG;AAC/B,UAAI,CAAC,YAAY;AACf,oBAAY,KAAK,SAAS;AAAA,MAC5B,OAAO;AACL,qBAAa,KAAK,SAAS;AAAA,MAC7B;AAAA,IACF;AAEA,QAAI,CAAC,cAAc,YAAY,SAAS,GAAG;AACzC,qBAAe;AACf,oBAAc,CAAC;AACf,mBAAa;AAAA,IACf;AAEA,UAAM,UAAU,MAAM,aAAa;AACnC,eAAW,QAAQ,QAAS,OAAM;AAAA,EACpC;AACF;AAQO,SAAS,mBACd,QACA,aACA,SACgB;AAChB,SAAO,IAAI,eAAe,QAAQ,aAAa,OAAO;AACxD;","names":["parse","Schema","flushed"]}