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 門面