Skip to content

EventStream:push 进、async iterator 出

源码版本v0.73.1

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 的具体子类),调用方零差别对待。

职责

  1. 队列 + waiter:push 时如果有等待的 consumer 就直接 resolve,否则塞进 queue,见 packages/ai/src/utils/event-stream.ts:20-35
  2. async iterator:[Symbol.asyncIterator] 实现 AsyncIterator<T>,先消费 queue 再等 waiter,见 packages/ai/src/utils/event-stream.ts:49-61
  3. 完成信号:isComplete(event) 判定哪个事件算「终态」,触发 resolveFinalResult(extractResult(event)),见 packages/ai/src/utils/event-stream.ts:23-26
  4. 最终结果 Promise:result(): Promise<R> 返回终态事件抽取的结果,与 for await 解耦,见 packages/ai/src/utils/event-stream.ts:63-65
  5. end 兜底:end(result?) 强制结束,清空所有 waiter,见 packages/ai/src/utils/event-stream.ts:37-47
  6. AssistantMessage 子类:AssistantMessageEventStreamdone / 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

关键文件

push 是 producer 唯一入口,等 waiter 或入队两种路径:

typescript
// 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 优先消费已入队事件,空了再等:

typescript
// 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 不用关心终态判定:

typescript
// 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 当终态:AssistantMessageEventStreamextractResulterror 时返回 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 门面