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 ベースの producer API を pull ベースの async iterator に翻訳する。すべての provider は最終的にイベントを AssistantMessageEventStream(このファイルの中の EventStream の具体サブクラス)に push し、呼び出し側は差を意識しない。

役割

  1. キュー + waiter: push 時、待っている consumer がいれば直接 resolve し、いなければ queue に詰む。packages/ai/src/utils/event-stream.ts:20-35 参照。
  2. async iterator: [Symbol.asyncIterator]AsyncIterator<T> を実装する。キューを先に消費し、その後 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 だけでバックプレッシャーを持たない。Observable を async iterator に変換するには from で适配が必要だ。80 行のクラスが二つの役割を一つに圧縮し、さらに「終了状態イベントが result Promise を発火する」まで自前で持つ。complete 関数は stream(...).result() の一行で AssistantMessage を取れ、consumer が手動で蓄積する必要がない。

なぜ AssistantMessageEventStream はジェネリック引数ではなくサブクラスなのか? 「どのイベントが終了状態か、終了状態からどう結果を抽出するか」は provider によらないポリシーだからだ。すべての provider は AssistantMessageEvent を push し、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 双方のやり取り:

境界と失敗

  • 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 を終了状態として扱う: 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 は一律 AssistantMessageEvent を push し、呼び出し側は一律 for await + result() を使う。provider がどう使うかは Anthropic SSE 実装、トップレベルのディスパッチは stream/complete ファサード を参照。