{"version":3,"sources":["../../src/streaming/adapters.ts"],"sourcesContent":["import { IOStreamSource, IOStreamTransport, StreamChunk } from './types';\r\n\r\n/**\r\n * A helper to bridge event-based or callback-based data sources to the AsyncIterable\r\n * required by `openStream`.\r\n *\r\n * @returns An object containing:\r\n * - `source`: The AsyncIterable to pass to `openStream`.\r\n * - `push`: A function to push a new chunk of data.\r\n * - `close`: A function to signal the end of the stream (or an error).\r\n *\r\n * @example\r\n * ```ts\r\n * const { source, push, close } = createPushSource();\r\n * openStream(source);\r\n *\r\n * xhr.onprogress = () => push(xhr.responseText.substring(seen));\r\n * xhr.onload = () => close();\r\n * xhr.onerror = () => close(new Error('Network error'));\r\n * ```\r\n */\r\nexport function createPushSource(): {\r\n  source: AsyncIterable<StreamChunk>;\r\n  push: (chunk: StreamChunk) => void;\r\n  close: (error?: Error) => void;\r\n} {\r\n  const queue: StreamChunk[] = [];\r\n  let resolvers: ((value: IteratorResult<StreamChunk> | PromiseLike<IteratorResult<StreamChunk>>) => void)[] = [];\r\n  let done = false;\r\n  let error: Error | undefined;\r\n\r\n  const push = (chunk: StreamChunk) => {\r\n    if (done) return;\r\n    if (resolvers.length > 0) {\r\n      const resolve = resolvers.shift()!;\r\n      resolve({ value: chunk, done: false });\r\n    } else {\r\n      queue.push(chunk);\r\n    }\r\n  };\r\n\r\n  const close = (err?: Error) => {\r\n    if (done) return;\r\n    done = true;\r\n    error = err;\r\n    while (resolvers.length > 0) {\r\n      const resolve = resolvers.shift()!;\r\n      if (err) {\r\n        // We can't easily reject the iterator promise in a way that plays nice with for-await\r\n        // without handling it carefully. Standard async iterators usually throw on next().\r\n        // Here we resolve with a rejected promise logic if we could, but for simplicity\r\n        // we will let the next call handle the error state.\r\n        // Actually, resolving with a rejected promise is valid.\r\n        Promise.reject(err).catch(() => {}); // Prevent unhandled rejection\r\n        // But we need to pass the error to the consumer.\r\n        // The cleanest way in a manual iterator is to throw when next() is called.\r\n        // Since we are resolving a pending next() call:\r\n        // We can't \"reject\" the resolve callback directly.\r\n        // We need to store the error and let the promise reject.\r\n        resolve(Promise.reject(err));\r\n      } else {\r\n        resolve({ value: undefined, done: true });\r\n      }\r\n    }\r\n  };\r\n\r\n  const source: AsyncIterable<StreamChunk> = {\r\n    [Symbol.asyncIterator]() {\r\n      return {\r\n        next(): Promise<IteratorResult<StreamChunk>> {\r\n          if (queue.length > 0) {\r\n            return Promise.resolve({ value: queue.shift()!, done: false });\r\n          }\r\n          if (done) {\r\n            if (error) return Promise.reject(error);\r\n            return Promise.resolve({ value: undefined, done: true });\r\n          }\r\n          return new Promise<IteratorResult<StreamChunk>>((resolve, reject) => {\r\n            // We wrap resolve to handle the error rejection logic from close()\r\n            resolvers.push((res: any) => {\r\n               // If res is a promise (from close with error), it will bubble up\r\n               resolve(res);\r\n            });\r\n          });\r\n        }\r\n      };\r\n    }\r\n  };\r\n\r\n  return { source, push, close };\r\n}\r\n\r\n/**\r\n * A transport that accumulates all written data into a single string buffer.\r\n * Useful for environments where streaming upload is not possible, and you need\r\n * to generate the full payload before sending.\r\n */\r\nexport class BufferTransport implements IOStreamTransport {\r\n  private chunks: string[] = [];\r\n\r\n  send(chunk: string | Uint8Array): void {\r\n    if (typeof chunk === 'string') {\r\n      this.chunks.push(chunk);\r\n    } else {\r\n      // Best effort text decoding\r\n      this.chunks.push(new TextDecoder().decode(chunk));\r\n    }\r\n  }\r\n\r\n  getOutput(): string {\r\n    return this.chunks.join('');\r\n  }\r\n}\r\n"],"mappings":";;;;;;;;;;;;;;;;;;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAqBO,SAAS,mBAId;AACA,QAAM,QAAuB,CAAC;AAC9B,MAAI,YAAyG,CAAC;AAC9G,MAAI,OAAO;AACX,MAAI;AAEJ,QAAM,OAAO,CAAC,UAAuB;AACnC,QAAI,KAAM;AACV,QAAI,UAAU,SAAS,GAAG;AACxB,YAAM,UAAU,UAAU,MAAM;AAChC,cAAQ,EAAE,OAAO,OAAO,MAAM,MAAM,CAAC;AAAA,IACvC,OAAO;AACL,YAAM,KAAK,KAAK;AAAA,IAClB;AAAA,EACF;AAEA,QAAM,QAAQ,CAAC,QAAgB;AAC7B,QAAI,KAAM;AACV,WAAO;AACP,YAAQ;AACR,WAAO,UAAU,SAAS,GAAG;AAC3B,YAAM,UAAU,UAAU,MAAM;AAChC,UAAI,KAAK;AAMP,gBAAQ,OAAO,GAAG,EAAE,MAAM,MAAM;AAAA,QAAC,CAAC;AAMlC,gBAAQ,QAAQ,OAAO,GAAG,CAAC;AAAA,MAC7B,OAAO;AACL,gBAAQ,EAAE,OAAO,QAAW,MAAM,KAAK,CAAC;AAAA,MAC1C;AAAA,IACF;AAAA,EACF;AAEA,QAAM,SAAqC;AAAA,IACzC,CAAC,OAAO,aAAa,IAAI;AACvB,aAAO;AAAA,QACL,OAA6C;AAC3C,cAAI,MAAM,SAAS,GAAG;AACpB,mBAAO,QAAQ,QAAQ,EAAE,OAAO,MAAM,MAAM,GAAI,MAAM,MAAM,CAAC;AAAA,UAC/D;AACA,cAAI,MAAM;AACR,gBAAI,MAAO,QAAO,QAAQ,OAAO,KAAK;AACtC,mBAAO,QAAQ,QAAQ,EAAE,OAAO,QAAW,MAAM,KAAK,CAAC;AAAA,UACzD;AACA,iBAAO,IAAI,QAAqC,CAAC,SAAS,WAAW;AAEnE,sBAAU,KAAK,CAAC,QAAa;AAE1B,sBAAQ,GAAG;AAAA,YACd,CAAC;AAAA,UACH,CAAC;AAAA,QACH;AAAA,MACF;AAAA,IACF;AAAA,EACF;AAEA,SAAO,EAAE,QAAQ,MAAM,MAAM;AAC/B;AAOO,MAAM,gBAA6C;AAAA,EAChD,SAAmB,CAAC;AAAA,EAE5B,KAAK,OAAkC;AACrC,QAAI,OAAO,UAAU,UAAU;AAC7B,WAAK,OAAO,KAAK,KAAK;AAAA,IACxB,OAAO;AAEL,WAAK,OAAO,KAAK,IAAI,YAAY,EAAE,OAAO,KAAK,CAAC;AAAA,IAClD;AAAA,EACF;AAAA,EAEA,YAAoB;AAClB,WAAO,KAAK,OAAO,KAAK,EAAE;AAAA,EAC5B;AACF;","names":[]}