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