{"version":3,"sources":["../../src/streaming/writer.ts"],"sourcesContent":["import Definitions from '../core/definitions';\r\nimport Document from '../core/document';\r\nimport Header from '../core/header';\r\nimport SectionCollection from '../core/section-collection';\r\nimport { loadObject } from '../facade/load';\r\nimport { stringify } from '../facade/stringify';\r\nimport { stringifyDocument } from '../facade/stringify-document';\r\nimport { IOStreamTransport, StreamWriterOptions } from './types';\r\n\r\nexport class IOStreamWriter {\r\n  private readonly transport: IOStreamTransport;\r\n  private readonly defs: Definitions | null;\r\n  private readonly options: StreamWriterOptions;\r\n\r\n  private headerDefs: Definitions | null = null;\r\n  private headerText: string | null = null;\r\n  private currentSchemaName: string | null = null;\r\n\r\n  constructor(transport: IOStreamTransport, defs?: Definitions | null, options?: StreamWriterOptions) {\r\n    this.transport = transport;\r\n    this.defs = defs ?? null;\r\n    this.options = {\r\n      includeSchemas: options?.includeSchemas ?? true,\r\n      defsId: options?.defsId,\r\n      onError: options?.onError,\r\n    };\r\n  }\r\n\r\n  /** Sets header metadata (non-schema definitions). Must be called before getHeader(). */\r\n  setHeader(header: Definitions | null): void {\r\n    this.headerDefs = header;\r\n    this.headerText = null; // invalidate cache\r\n  }\r\n\r\n  /** Returns the full header chunk including the terminating `---`. */\r\n  getHeader(): string {\r\n    if (this.headerText) return this.headerText;\r\n\r\n    const header = new Header();\r\n\r\n    // 1) metadata header first\r\n    if (this.headerDefs) {\r\n      header.definitions.merge(this.headerDefs, true);\r\n    }\r\n\r\n    // 2) optional defs identifier (design hook)\r\n    if (this.options.defsId) {\r\n      header.definitions.set('defsId', this.options.defsId);\r\n    }\r\n\r\n    // 3) schema definitions (if configured)\r\n    if (this.options.includeSchemas && this.defs) {\r\n      header.definitions.merge(this.defs, false);\r\n    }\r\n\r\n    const doc = new Document(header, new SectionCollection());\r\n\r\n    // stringifyDocument includes header definitions but does NOT add the required terminator.\r\n    const headerBody = stringifyDocument(doc, { includeHeader: true });\r\n    const normalized = headerBody ? headerBody.trimEnd() : '';\r\n\r\n    this.headerText = normalized.length > 0 ? `${normalized}\\n---\\n` : `---\\n`;\r\n\r\n    // initialize current schema from defs ($schema) if present\r\n    this.currentSchemaName = '$schema';\r\n\r\n    return this.headerText;\r\n  }\r\n\r\n  /** Sends header immediately via transport. */\r\n  async sendHeader(): Promise<void> {\r\n    const h = this.getHeader();\r\n    await this.transport.send(h);\r\n  }\r\n\r\n  /** Emits a schema switch marker. Use `$schema` or omit to switch back to default. */\r\n  section(schemaName?: string): string {\r\n    const s = schemaName ? `--- ${schemaName}\\n` : `---\\n`;\r\n    this.currentSchemaName = schemaName ?? '$schema';\r\n    return s;\r\n  }\r\n\r\n  /**\r\n   * Serializes one item.\r\n   * If schemaName changes, prepends a section switch marker automatically.\r\n   */\r\n  write(data: any, schemaName?: string): string {\r\n    const effectiveSchema = schemaName ?? '$schema';\r\n    const parts: string[] = [];\r\n\r\n    try {\r\n      if (this.currentSchemaName !== effectiveSchema) {\r\n        // only include schema marker for non-default schema\r\n        parts.push(this.section(schemaName));\r\n      }\r\n\r\n      // Validate+wrap into InternetObject if schema available.\r\n      // If defs is not provided, loadObject will be schemaless.\r\n      const ioObj = this.defs\r\n        ? loadObject(data, this.defs, { schemaName: effectiveSchema })\r\n        : loadObject(data as any);\r\n\r\n      // stringify() outputs record without '~'\r\n      const row = this.defs\r\n        ? stringify(ioObj as any, this.defs, { schemaName: effectiveSchema })\r\n        : stringify(ioObj as any);\r\n\r\n      parts.push(`~ ${row}\\n`);\r\n      return parts.join('');\r\n    } catch (err: any) {\r\n      const action = this.options.onError ?? 'throw';\r\n\r\n      if (action === 'throw') {\r\n        throw err;\r\n      }\r\n\r\n      if (action === 'ignore') {\r\n        return '';\r\n      }\r\n\r\n      if (action === 'emit') {\r\n        // Switch to $error section to avoid schema validation issues on client\r\n        // We don't update this.currentSchemaName to '$error' permanently because\r\n        // we want the next write() to switch back to its intended schema if needed.\r\n        // However, write() logic checks this.currentSchemaName.\r\n        // So we MUST update this.currentSchemaName so the next write knows we are in $error.\r\n        parts.push(this.section('$error'));\r\n\r\n        // Simple error object\r\n        const errorObj = {\r\n          code: err.errorCode || 'error',\r\n          message: err.message || String(err)\r\n        };\r\n\r\n        // We use JSON.stringify for safety, assuming simple error structure\r\n        // IO format is compatible with JSON for simple objects\r\n        parts.push(`~ ${JSON.stringify(errorObj)}\\n`);\r\n\r\n        return parts.join('');\r\n      }\r\n\r\n      return '';\r\n    }\r\n  }\r\n\r\n  /**\r\n   * Writes a batch of items.\r\n   * Efficiently handles schema switching (only switches once if all items share schema).\r\n   */\r\n  writeBatch(items: object[], schemaName?: string): string {\r\n    let out = '';\r\n    for (const item of items) {\r\n      out += this.write(item, schemaName);\r\n    }\r\n    return out;\r\n  }\r\n\r\n  /**\r\n   * Serializes and sends one item via the transport.\r\n   */\r\n  async send(data: object, schemaName?: string): Promise<void> {\r\n    const chunk = this.write(data, schemaName);\r\n    if (chunk) {\r\n      await this.transport.send(chunk);\r\n    }\r\n  }\r\n\r\n  /**\r\n   * Serializes and sends a batch of items via the transport.\r\n   */\r\n  async sendBatch(items: object[], schemaName?: string): Promise<void> {\r\n    const chunk = this.writeBatch(items, schemaName);\r\n    if (chunk) {\r\n      await this.transport.send(chunk);\r\n    }\r\n  }\r\n}\r\n\r\nexport function createStreamWriter(\r\n  transport: IOStreamTransport | { write: (chunk: any) => boolean | void },\r\n  defs?: Definitions | null,\r\n  options?: StreamWriterOptions\r\n): IOStreamWriter {\r\n  // Duck-type check for Node.js Writable stream\r\n  if (transport && typeof (transport as any).write === 'function' && typeof (transport as any).send !== 'function') {\r\n    const writable = transport as any;\r\n    transport = {\r\n      send: (chunk: string | Uint8Array) => {\r\n        const ok = writable.write(chunk);\r\n        // Handle backpressure if needed? For now we just write.\r\n        if (!ok) {\r\n           // simple drain handler could be async but send() result is void|Promise<void>\r\n           // Implementing proper backpressure requires async send().\r\n           // For now, we fire and forget or trust the stream buffers.\r\n        }\r\n      }\r\n    };\r\n  }\r\n\r\n  return new IOStreamWriter(transport as IOStreamTransport, defs ?? null, options);\r\n}\r\n"],"mappings":";;;;;;;;;;;;;;;;;;;;;;;;;;;;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AACA,sBAAqB;AACrB,oBAAmB;AACnB,gCAA8B;AAC9B,kBAA2B;AAC3B,uBAA0B;AAC1B,gCAAkC;AAG3B,MAAM,eAAe;AAAA,EACT;AAAA,EACA;AAAA,EACA;AAAA,EAET,aAAiC;AAAA,EACjC,aAA4B;AAAA,EAC5B,oBAAmC;AAAA,EAE3C,YAAY,WAA8B,MAA2B,SAA+B;AAClG,SAAK,YAAY;AACjB,SAAK,OAAO,QAAQ;AACpB,SAAK,UAAU;AAAA,MACb,gBAAgB,SAAS,kBAAkB;AAAA,MAC3C,QAAQ,SAAS;AAAA,MACjB,SAAS,SAAS;AAAA,IACpB;AAAA,EACF;AAAA;AAAA,EAGA,UAAU,QAAkC;AAC1C,SAAK,aAAa;AAClB,SAAK,aAAa;AAAA,EACpB;AAAA;AAAA,EAGA,YAAoB;AAClB,QAAI,KAAK,WAAY,QAAO,KAAK;AAEjC,UAAM,SAAS,IAAI,cAAAA,QAAO;AAG1B,QAAI,KAAK,YAAY;AACnB,aAAO,YAAY,MAAM,KAAK,YAAY,IAAI;AAAA,IAChD;AAGA,QAAI,KAAK,QAAQ,QAAQ;AACvB,aAAO,YAAY,IAAI,UAAU,KAAK,QAAQ,MAAM;AAAA,IACtD;AAGA,QAAI,KAAK,QAAQ,kBAAkB,KAAK,MAAM;AAC5C,aAAO,YAAY,MAAM,KAAK,MAAM,KAAK;AAAA,IAC3C;AAEA,UAAM,MAAM,IAAI,gBAAAC,QAAS,QAAQ,IAAI,0BAAAC,QAAkB,CAAC;AAGxD,UAAM,iBAAa,6CAAkB,KAAK,EAAE,eAAe,KAAK,CAAC;AACjE,UAAM,aAAa,aAAa,WAAW,QAAQ,IAAI;AAEvD,SAAK,aAAa,WAAW,SAAS,IAAI,GAAG,UAAU;AAAA;AAAA,IAAY;AAAA;AAGnE,SAAK,oBAAoB;AAEzB,WAAO,KAAK;AAAA,EACd;AAAA;AAAA,EAGA,MAAM,aAA4B;AAChC,UAAM,IAAI,KAAK,UAAU;AACzB,UAAM,KAAK,UAAU,KAAK,CAAC;AAAA,EAC7B;AAAA;AAAA,EAGA,QAAQ,YAA6B;AACnC,UAAM,IAAI,aAAa,OAAO,UAAU;AAAA,IAAO;AAAA;AAC/C,SAAK,oBAAoB,cAAc;AACvC,WAAO;AAAA,EACT;AAAA;AAAA;AAAA;AAAA;AAAA,EAMA,MAAM,MAAW,YAA6B;AAC5C,UAAM,kBAAkB,cAAc;AACtC,UAAM,QAAkB,CAAC;AAEzB,QAAI;AACF,UAAI,KAAK,sBAAsB,iBAAiB;AAE9C,cAAM,KAAK,KAAK,QAAQ,UAAU,CAAC;AAAA,MACrC;AAIA,YAAM,QAAQ,KAAK,WACf,wBAAW,MAAM,KAAK,MAAM,EAAE,YAAY,gBAAgB,CAAC,QAC3D,wBAAW,IAAW;AAG1B,YAAM,MAAM,KAAK,WACb,4BAAU,OAAc,KAAK,MAAM,EAAE,YAAY,gBAAgB,CAAC,QAClE,4BAAU,KAAY;AAE1B,YAAM,KAAK,KAAK,GAAG;AAAA,CAAI;AACvB,aAAO,MAAM,KAAK,EAAE;AAAA,IACtB,SAAS,KAAU;AACjB,YAAM,SAAS,KAAK,QAAQ,WAAW;AAEvC,UAAI,WAAW,SAAS;AACtB,cAAM;AAAA,MACR;AAEA,UAAI,WAAW,UAAU;AACvB,eAAO;AAAA,MACT;AAEA,UAAI,WAAW,QAAQ;AAMrB,cAAM,KAAK,KAAK,QAAQ,QAAQ,CAAC;AAGjC,cAAM,WAAW;AAAA,UACf,MAAM,IAAI,aAAa;AAAA,UACvB,SAAS,IAAI,WAAW,OAAO,GAAG;AAAA,QACpC;AAIA,cAAM,KAAK,KAAK,KAAK,UAAU,QAAQ,CAAC;AAAA,CAAI;AAE5C,eAAO,MAAM,KAAK,EAAE;AAAA,MACtB;AAEA,aAAO;AAAA,IACT;AAAA,EACF;AAAA;AAAA;AAAA;AAAA;AAAA,EAMA,WAAW,OAAiB,YAA6B;AACvD,QAAI,MAAM;AACV,eAAW,QAAQ,OAAO;AACxB,aAAO,KAAK,MAAM,MAAM,UAAU;AAAA,IACpC;AACA,WAAO;AAAA,EACT;AAAA;AAAA;AAAA;AAAA,EAKA,MAAM,KAAK,MAAc,YAAoC;AAC3D,UAAM,QAAQ,KAAK,MAAM,MAAM,UAAU;AACzC,QAAI,OAAO;AACT,YAAM,KAAK,UAAU,KAAK,KAAK;AAAA,IACjC;AAAA,EACF;AAAA;AAAA;AAAA;AAAA,EAKA,MAAM,UAAU,OAAiB,YAAoC;AACnE,UAAM,QAAQ,KAAK,WAAW,OAAO,UAAU;AAC/C,QAAI,OAAO;AACT,YAAM,KAAK,UAAU,KAAK,KAAK;AAAA,IACjC;AAAA,EACF;AACF;AAEO,SAAS,mBACd,WACA,MACA,SACgB;AAEhB,MAAI,aAAa,OAAQ,UAAkB,UAAU,cAAc,OAAQ,UAAkB,SAAS,YAAY;AAChH,UAAM,WAAW;AACjB,gBAAY;AAAA,MACV,MAAM,CAAC,UAA+B;AACpC,cAAM,KAAK,SAAS,MAAM,KAAK;AAE/B,YAAI,CAAC,IAAI;AAAA,QAIT;AAAA,MACF;AAAA,IACF;AAAA,EACF;AAEA,SAAO,IAAI,eAAe,WAAgC,QAAQ,MAAM,OAAO;AACjF;","names":["Header","Document","SectionCollection"]}