Skip to content

EventStream : push en entrée, itérateur async en sortie

源码版本v0.73.1

utils/event-stream.ts est la primitive de streaming de pi-ai—une classe générique EventStream<T, R> : le producer y pousse des événements via push(), le consommateur les tire via for await et récupère le résultat final via result(). En dessous, pas de WebSocket, pas d'Observable : juste une queue écrite à la main + un tableau de waiter, qui traduit une API producer push-based en itérateur async pull-based. Tous les provider finissent par pousser leurs événements dans un AssistantMessageEventStream (la sous-classe concrète de EventStream dans ce fichier), l'appelant les traite tous pareil.

Responsabilités

  1. Queue + waiter : à push, si un consommateur attend, on le resolve directement, sinon on pousse dans queue, voir packages/ai/src/utils/event-stream.ts:20-35.
  2. Itérateur async : [Symbol.asyncIterator] implémente AsyncIterator<T>, consomme d'abord la queue puis attend les waiter, voir packages/ai/src/utils/event-stream.ts:49-61.
  3. Signal de complétion : isComplete(event) décide quel événement compte comme « terminal » et déclenche resolveFinalResult(extractResult(event)), voir packages/ai/src/utils/event-stream.ts:23-26.
  4. Promise de résultat final : result(): Promise<R> renvoie le résultat extrait de l'événement terminal, découplé du for await, voir packages/ai/src/utils/event-stream.ts:63-65.
  5. Filet de sécurité end : end(result?) force la fin, vide tous les waiter, voir packages/ai/src/utils/event-stream.ts:37-47.
  6. Sous-classe AssistantMessage : AssistantMessageEventStream considère done / error comme terminaux, done.message / error.error comme résultat, voir packages/ai/src/utils/event-stream.ts:68-82.

Motivation de design

Pourquoi ne pas utiliser directement le EventEmitter de Node ou un Observable RxJS ? Parce que le consommateur veut la syntaxe for await—le rythme de tirage est contrôlé par le consommateur, le backpressure vient naturellement. EventEmitter pousse sans backpressure; un Observable doit passer par from pour devenir un itérateur async. Une classe de 80 lignes compresse les deux, et embarque en plus « l'événement terminal déclenche la result Promise » : la fonction complete fait stream(...).result() en une ligne pour récupérer un AssistantMessage, sans que le consommateur ait à accumuler à la main.

Pourquoi AssistantMessageEventStream est une sous-classe plutôt qu'un paramètre générique ? Parce que « quel événement est terminal, comment extraire le résultat du terminal » est une stratégie indépendante du provider—tous les provider poussent du AssistantMessageEvent, done.message est le résultat final. La sous-classe fige la stratégie, les provider n'ont pas à passer chacun leur isComplete / extractResult.

Fichiers clés

push est la seule entrée du producer, deux chemins : attendant un waiter ou enfiler :

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);
	}
}

L'itérateur async consomme d'abord les événements déjà enfilés, puis attend si vide :

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 fige la stratégie, le provider n'a pas à se soucier de la détection terminal :

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");
			},
		);
	}
}

Flux de données

L'interaction producer / consommateur :

Frontières et échecs

  • push sur un stream déjà done : if (this.done) return ignore silencieusement; si le provider repousse après une error, rien ne casse, voir packages/ai/src/utils/event-stream.ts:21-21.
  • end avant un événement terminal : si end(result?) est appelé sans résultat, finalResultPromise peut ne jamais résoudre; l'appelant doit s'assurer qu'un événement terminal a été poussé au préalable ou passer un result, voir packages/ai/src/utils/event-stream.ts:37-41.
  • error comme terminal : l'extractResult de AssistantMessageEventStream renvoie event.error sur error, donc result() récupère aussi l'objet erreur sur un flux d'erreur, sans hang, voir packages/ai/src/utils/event-stream.ts:75-77.
  • Plusieurs consommateurs : plusieurs for await en parallèle se disputent la file des waiter, les événements se répartissent sur différents consommateurs—mais en pratique il n'y a qu'un seul consommateur; le code ne restreint pas le multi-consommateur.

Récapitulatif

EventStream est un adaptateur push-to-async-iterator écrit à la main en 80 lignes; la sous-classe AssistantMessageEventStream fige la stratégie terminal. Les provider poussent tous du AssistantMessageEvent; les appelants font tous for await + result(). L'usage côté provider se lit dans Implémentation Anthropic SSE, et le dispatch de haut niveau dans Façade stream/complete.