diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index d299278ee..a19015641 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed quadratic CPU usage when draining buffered `EventStream` events ([#9055](https://github.com/earendil-works/pi/issues/9055)). + ## [0.85.1] - 2026-09-05 ### Added diff --git a/packages/ai/src/utils/event-stream.ts b/packages/ai/src/utils/event-stream.ts index e5eef921d..c4f731d59 100644 --- a/packages/ai/src/utils/event-stream.ts +++ b/packages/ai/src/utils/event-stream.ts @@ -1,9 +1,31 @@ import type { AssistantMessage, AssistantMessageEvent } from "../types.ts"; +class FifoQueue { + private incoming: T[] = []; + private outgoing: T[] = []; + + get length(): number { + return this.incoming.length + this.outgoing.length; + } + + enqueue(value: T): void { + this.incoming.push(value); + } + + dequeue(): T | undefined { + if (this.outgoing.length === 0) { + while (this.incoming.length > 0) { + this.outgoing.push(this.incoming.pop()!); + } + } + return this.outgoing.pop(); + } +} + // Generic event stream class for async iteration export class EventStream implements AsyncIterable { - private queue: T[] = []; - private waiting: ((value: IteratorResult) => void)[] = []; + private queue = new FifoQueue(); + private waiting = new FifoQueue<(value: IteratorResult) => void>(); private done = false; private finalResultPromise: Promise; private resolveFinalResult!: (result: R) => void; @@ -27,11 +49,11 @@ export class EventStream implements AsyncIterable { } // Deliver to waiting consumer or queue it - const waiter = this.waiting.shift(); + const waiter = this.waiting.dequeue(); if (waiter) { waiter({ value: event, done: false }); } else { - this.queue.push(event); + this.queue.enqueue(event); } } @@ -42,7 +64,7 @@ export class EventStream implements AsyncIterable { } // Notify all waiting consumers that we're done while (this.waiting.length > 0) { - const waiter = this.waiting.shift()!; + const waiter = this.waiting.dequeue()!; waiter({ value: undefined as any, done: true }); } } @@ -50,11 +72,11 @@ export class EventStream implements AsyncIterable { async *[Symbol.asyncIterator](): AsyncIterator { while (true) { if (this.queue.length > 0) { - yield this.queue.shift()!; + yield this.queue.dequeue()!; } else if (this.done) { return; } else { - const result = await new Promise>((resolve) => this.waiting.push(resolve)); + const result = await new Promise>((resolve) => this.waiting.enqueue(resolve)); if (result.done) return; yield result.value; } diff --git a/packages/ai/test/event-stream.test.ts b/packages/ai/test/event-stream.test.ts new file mode 100644 index 000000000..897ce5913 --- /dev/null +++ b/packages/ai/test/event-stream.test.ts @@ -0,0 +1,94 @@ +import { describe, expect, it } from "vitest"; +import { EventStream } from "../src/utils/event-stream.ts"; + +// Regression tests for https://github.com/earendil-works/pi/issues/9055 +describe("EventStream", () => { + it("drains buffered events in order and ignores events pushed after completion", async () => { + const stream = new EventStream( + (event) => event === 3, + (event) => event, + ); + stream.push(1); + stream.push(2); + stream.push(3); + stream.push(4); + + expect(await stream.result()).toBe(3); + + const events: number[] = []; + for await (const event of stream) { + events.push(event); + } + expect(events).toEqual([1, 2, 3]); + }); + + it("preserves order when events arrive after buffered draining starts", async () => { + const stream = new EventStream( + () => false, + (event) => event, + ); + stream.push(1); + stream.push(2); + + const iterator = stream[Symbol.asyncIterator](); + expect(await iterator.next()).toEqual({ value: 1, done: false }); + + stream.push(3); + expect(await iterator.next()).toEqual({ value: 2, done: false }); + expect(await iterator.next()).toEqual({ value: 3, done: false }); + + stream.end(3); + expect(await iterator.next()).toEqual({ value: undefined, done: true }); + }); + + it("delivers events to waiting consumers in registration order", async () => { + const stream = new EventStream( + () => false, + (event) => event, + ); + const firstIterator = stream[Symbol.asyncIterator](); + const secondIterator = stream[Symbol.asyncIterator](); + const firstEvent = firstIterator.next(); + const secondEvent = secondIterator.next(); + + stream.push(1); + stream.push(2); + + expect(await firstEvent).toEqual({ value: 1, done: false }); + expect(await secondEvent).toEqual({ value: 2, done: false }); + }); + + it("drains buffered events after end and resolves the explicit result", async () => { + const stream = new EventStream( + () => false, + (event) => String(event), + ); + stream.push(1); + stream.push(2); + stream.end("complete"); + + expect(await stream.result()).toBe("complete"); + + const events: number[] = []; + for await (const event of stream) { + events.push(event); + } + expect(events).toEqual([1, 2]); + }); + + it("wakes all waiting consumers when ended without a result", async () => { + const stream = new EventStream( + () => false, + (event) => event, + ); + const firstIterator = stream[Symbol.asyncIterator](); + const secondIterator = stream[Symbol.asyncIterator](); + const firstEvent = firstIterator.next(); + const secondEvent = secondIterator.next(); + + stream.end(); + + expect(await firstEvent).toEqual({ value: undefined, done: true }); + expect(await secondEvent).toEqual({ value: undefined, done: true }); + }); +});