/** * Async input queue — copy of the helper from `harnesses/claude.ts:37-68`. * * The Claude harness uses this exact shape as the input prompt to the Claude * Agent SDK's long-lived `query()`. The pi session loop uses the same pattern * so the non-blocking live-conversation behavior matches Claude byte-for-byte * at the queue level: pushMessage() never awaits the model, and each queued * message is consumed as its own turn. * * Deliberately NO bulk-drain helper: an earlier `drainPending()` let the * session fold mid-turn messages into the in-flight turn, which broke the * channel manager's one-routing-target-per-push FIFO (see * PI-PARITY-AUDIT-2026-06-11.md D1-1). One message in → one turn out. */ export interface AsyncQueue extends AsyncIterable { push(item: T): void; end(): void; } export function createAsyncQueue(): AsyncQueue { const pending: T[] = []; let resolve: ((value: IteratorResult) => void) | null = null; let done = false; return { push(item: T) { if (done) return; if (resolve) { resolve({ value: item, done: false }); resolve = null; } else { pending.push(item); } }, end() { done = true; if (resolve) resolve({ value: undefined as any, done: true }); }, [Symbol.asyncIterator]() { return { next(): Promise> { if (pending.length > 0) { return Promise.resolve({ value: pending.shift()!, done: false }); } if (done) return Promise.resolve({ value: undefined as any, done: true }); return new Promise((r) => { resolve = r; }); }, }; }, }; }