fix(ai): optimize EventStream queue

closes #9055
This commit is contained in:
Mario Zechner
2026-09-07 22:26:47 +02:00
parent 9211da1723
commit b2602be77c
3 changed files with 127 additions and 7 deletions
+4
View File
@@ -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
+29 -7
View File
@@ -1,9 +1,31 @@
import type { AssistantMessage, AssistantMessageEvent } from "../types.ts";
class FifoQueue<T> {
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<T, R = T> implements AsyncIterable<T> {
private queue: T[] = [];
private waiting: ((value: IteratorResult<T>) => void)[] = [];
private queue = new FifoQueue<T>();
private waiting = new FifoQueue<(value: IteratorResult<T>) => void>();
private done = false;
private finalResultPromise: Promise<R>;
private resolveFinalResult!: (result: R) => void;
@@ -27,11 +49,11 @@ export class EventStream<T, R = T> implements AsyncIterable<T> {
}
// 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<T, R = T> implements AsyncIterable<T> {
}
// 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<T, R = T> implements AsyncIterable<T> {
async *[Symbol.asyncIterator](): AsyncIterator<T> {
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<IteratorResult<T>>((resolve) => this.waiting.push(resolve));
const result = await new Promise<IteratorResult<T>>((resolve) => this.waiting.enqueue(resolve));
if (result.done) return;
yield result.value;
}
+94
View File
@@ -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<number, number>(
(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<number, number>(
() => 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<number, number>(
() => 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<number, string>(
() => 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<number, number>(
() => 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 });
});
});