/** * Minimal AssistantMessageEventStream-compatible factory. * * pi-ai v0.77+ no longer re-exports the `AssistantMessageEventStream` class from * its package root (it is type-only there), and there is no public subpath to * the concrete class. The vendored Command Code backend only needs an object * with `push(event)` + `end()` that is also `AsyncIterable`, so this tiny * async-queue satisfies that contract and is iterated identically by pi's * streaming consumer. * * Used solely as the `createStream` dependency for `createStreamCommandCode`. */ type EventLike = Record interface Awaiter { resolve: (r: IteratorResult) => void } /** * A minimal async event stream: producers call push()/end(); consumers iterate * with for-await-of. Buffered events are delivered in order; once end() is * called and the buffer drains, iteration completes. */ export class MiniEventStream implements AsyncIterable { private buffer: EventLike[] = [] private awaiters: Awaiter[] = [] private finished = false push(event: EventLike): void { if (this.finished) return const waiter = this.awaiters.shift() if (waiter) { waiter.resolve({ value: event, done: false }) } else { this.buffer.push(event) } } end(): void { this.finished = true const pending = this.awaiters this.awaiters = [] for (const w of pending) w.resolve({ value: undefined, done: true }) } async next(): Promise> { if (this.buffer.length > 0) { return { value: this.buffer.shift() as EventLike, done: false } } if (this.finished) return { value: undefined, done: true } return new Promise>((resolve) => { this.awaiters.push({ resolve }) }) } [Symbol.asyncIterator](): AsyncIterator { return { next: () => this.next(), } } } /** Factory matching the `createStream: () => AssistantMessageEventStreamLike` dep. */ export function createMiniEventStream(): MiniEventStream { return new MiniEventStream() }