EventStream:push 进、async iterator 出
utils/event-stream.ts 是 pi-ai 的流式原语——一个泛型 EventStream<T, R>,producer 用 push() 往里塞事件,consumer 用 for await 拉事件,用 result() 拿最终结果。底层不是 WebSocket、不是 Observable,就是一个手写的队列 + waiter 数组,把 push-based 的 producer API 翻译成 pull-based 的 async iterator。所有 provider 最终都把事件 push 进 AssistantMessageEventStream(本文件里 EventStream 的具体子类),调用方零差别对待。
职责
- 队列 + waiter:
push时如果有等待的 consumer 就直接 resolve,否则塞进queue,见packages/ai/src/utils/event-stream.ts:20-35。 - async iterator:
[Symbol.asyncIterator]实现AsyncIterator<T>,先消费 queue 再等 waiter,见packages/ai/src/utils/event-stream.ts:49-61。 - 完成信号:
isComplete(event)判定哪个事件算「终态」,触发resolveFinalResult(extractResult(event)),见packages/ai/src/utils/event-stream.ts:23-26。 - 最终结果 Promise:
result(): Promise<R>返回终态事件抽取的结果,与for await解耦,见packages/ai/src/utils/event-stream.ts:63-65。 - end 兜底:
end(result?)强制结束,清空所有 waiter,见packages/ai/src/utils/event-stream.ts:37-47。 - AssistantMessage 子类:
AssistantMessageEventStream把done/error当终态,done.message/error.error当结果,见packages/ai/src/utils/event-stream.ts:68-82。
设计动机
为什么不直接用 Node EventEmitter 或 RxJS Observable?因为 consumer 要 for await 语法——拉取节奏由 consumer 控制,backpressure 自然有。EventEmitter 是 push 不带 backpressure,Observable 转成 async iterator 要 from 适配。一个 80 行的类把这两件事压到一起,还自带「终态事件触发 result Promise」,complete 函数一行 stream(...).result() 就能拿到 AssistantMessage,不需要 consumer 手动累积。
为什么 AssistantMessageEventStream 是子类而不是泛型参数?因为「什么事件算终态、怎么从终态抽结果」是 provider 无关的策略——所有 provider 都 push AssistantMessageEvent,done.message 就是最终结果。子类把策略写死,provider 不用各自传 isComplete / extractResult。
关键文件
packages/ai/src/utils/event-stream.ts:4-18—EventStream<T, R>类声明 + 构造函数收isComplete/extractResult两个策略。packages/ai/src/utils/event-stream.ts:20-35—push:终态判定 + 等 waiter 或入队。packages/ai/src/utils/event-stream.ts:37-47—end:强制 done + 清空 waiter。packages/ai/src/utils/event-stream.ts:49-61—[Symbol.asyncIterator]:queue 优先、done 退出、否则等 Promise。packages/ai/src/utils/event-stream.ts:63-65—result:返回终态 Promise。packages/ai/src/utils/event-stream.ts:68-82—AssistantMessageEventStream子类,定终态策略。packages/ai/src/utils/event-stream.ts:84-87—createAssistantMessageEventStream工厂,扩展用。
push 是 producer 唯一入口,等 waiter 或入队两种路径:
// packages/ai/src/utils/event-stream.ts:20-35
push(event: T): void {
if (this.done) return;
if (this.isComplete(event)) {
this.done = true;
this.resolveFinalResult(this.extractResult(event));
}
// Deliver to waiting consumer or queue it
const waiter = this.waiting.shift();
if (waiter) {
waiter({ value: event, done: false });
} else {
this.queue.push(event);
}
}async iterator 优先消费已入队事件,空了再等:
// packages/ai/src/utils/event-stream.ts:49-61
async *[Symbol.asyncIterator](): AsyncIterator<T> {
while (true) {
if (this.queue.length > 0) {
yield this.queue.shift()!;
} else if (this.done) {
return;
} else {
const result = await new Promise<IteratorResult<T>>((resolve) => this.waiting.push(resolve));
if (result.done) return;
yield result.value;
}
}
}AssistantMessageEventStream 把策略写死,provider 不用关心终态判定:
// packages/ai/src/utils/event-stream.ts:68-82
export class AssistantMessageEventStream extends EventStream<AssistantMessageEvent, AssistantMessage> {
constructor() {
super(
(event) => event.type === "done" || event.type === "error",
(event) => {
if (event.type === "done") {
return event.message;
} else if (event.type === "error") {
return event.error;
}
throw new Error("Unexpected event type for final result");
},
);
}
}数据流
producer 与 consumer 双方的交互:
边界与失败
- push 到已 done:
if (this.done) return静默忽略,provider 在error后再 push 不会出问题,见packages/ai/src/utils/event-stream.ts:21-21。 - 未触发终态就 end:
end(result?)不传 result 时finalResultPromise可能永远不 resolve,调用方应保证终态事件先 push 过或传 result,见packages/ai/src/utils/event-stream.ts:37-41。 - error 当终态:
AssistantMessageEventStream的extractResult在error时返回event.error,所以result()在错误流上也能拿到错误对象,而不是 hang,见packages/ai/src/utils/event-stream.ts:75-77。 - 多 consumer:多个
for await同时跑会争抢 waiter 队列,事件会被分到不同 consumer——但实际场景都是单 consumer,代码没限制多 consumer 行为。
小结
EventStream 是 80 行的手写 push-to-async-iterator 适配器,AssistantMessageEventStream 子类定死终态策略。provider 一律 push AssistantMessageEvent,调用方一律 for await + result()。provider 怎么用它的看 Anthropic SSE 实现,顶层派发看 stream/complete 门面。