EventStream: push で入れ、async iterator で出す
utils/event-stream.ts は pi-ai のストリーミングプリミティブだ。一つのジェネリック EventStream<T, R> で、producer は push() でイベントを詰め、consumer は for await でイベントを引き、result() で最終結果を取る。下層は WebSocket でも Observable でもなく、手書きのキュー + waiter 配列で、push ベースの producer API を pull ベースの async iterator に翻訳する。すべての provider は最終的にイベントを AssistantMessageEventStream(このファイルの中の EventStream の具体サブクラス)に push し、呼び出し側は差を意識しない。
役割
- キュー + waiter:
push時、待っている consumer がいれば直接 resolve し、いなければqueueに詰む。packages/ai/src/utils/event-stream.ts:20-35参照。 - async iterator:
[Symbol.asyncIterator]がAsyncIterator<T>を実装する。キューを先に消費し、その後 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 だけでバックプレッシャーを持たない。Observable を async iterator に変換するには from で适配が必要だ。80 行のクラスが二つの役割を一つに圧縮し、さらに「終了状態イベントが result Promise を発火する」まで自前で持つ。complete 関数は stream(...).result() の一行で AssistantMessage を取れ、consumer が手動で蓄積する必要がない。
なぜ AssistantMessageEventStream はジェネリック引数ではなくサブクラスなのか? 「どのイベントが終了状態か、終了状態からどう結果を抽出するか」は provider によらないポリシーだからだ。すべての provider は AssistantMessageEvent を push し、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]: キュー優先・done で exit・それ以外は 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 双方のやり取り:
境界と失敗
- done 後の push:
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 は一律 AssistantMessageEvent を push し、呼び出し側は一律 for await + result() を使う。provider がどう使うかは Anthropic SSE 実装、トップレベルのディスパッチは stream/complete ファサード を参照。