/** * A simple push-based async queue that implements AsyncIterable. * Values are pushed via `.send(item)` and the iteration ends when `.close()` is called. */ export class PushChannel implements AsyncIterable { private readonly queue: T[] = []; private readonly resolvers: Array<(result: IteratorResult) => void> = []; private closed = false; /** Push a value to waiting consumers (or enqueue if none are waiting). */ send(item: T): void { if (this.closed) throw new Error('PushChannel is closed'); if (this.resolvers.length > 0) { const resolve = this.resolvers.shift()!; resolve({ value: item, done: false }); } else { this.queue.push(item); } } /** Signal end-of-stream; iteration will finish after queued items. */ close(): void { this.closed = true; for (const resolve of this.resolvers.splice(0)) { resolve({ value: undefined as unknown as T, done: true }); } } [Symbol.asyncIterator](): AsyncIterator { return { next: (): Promise> => { if (this.queue.length > 0) { return Promise.resolve({ value: this.queue.shift()!, done: false }); } if (this.closed) { return Promise.resolve({ value: undefined as unknown as T, done: true }); } return new Promise>((resolve) => { this.resolvers.push(resolve); }); }, return: (): Promise> => { return Promise.resolve({ value: undefined as unknown as T, done: true }); }, }; } } /** * Create a push-based channel with `.send(item)` and `.close()` for async iteration. */ export function pushChannel(_opts?: { highWaterMark?: number }): PushChannel { return new PushChannel(); }