EventStream : push en entrée, itérateur async en sortie
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
- Queue + waiter : à
push, si un consommateur attend, on le resolve directement, sinon on pousse dansqueue, voirpackages/ai/src/utils/event-stream.ts:20-35. - Itérateur async :
[Symbol.asyncIterator]implémenteAsyncIterator<T>, consomme d'abord la queue puis attend les waiter, voirpackages/ai/src/utils/event-stream.ts:49-61. - Signal de complétion :
isComplete(event)décide quel événement compte comme « terminal » et déclencheresolveFinalResult(extractResult(event)), voirpackages/ai/src/utils/event-stream.ts:23-26. - Promise de résultat final :
result(): Promise<R>renvoie le résultat extrait de l'événement terminal, découplé dufor await, voirpackages/ai/src/utils/event-stream.ts:63-65. - Filet de sécurité end :
end(result?)force la fin, vide tous les waiter, voirpackages/ai/src/utils/event-stream.ts:37-47. - Sous-classe AssistantMessage :
AssistantMessageEventStreamconsidèredone/errorcomme terminaux,done.message/error.errorcomme résultat, voirpackages/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
packages/ai/src/utils/event-stream.ts:4-18— déclaration de classeEventStream<T, R>+ constructeur qui reçoit les deux stratégiesisComplete/extractResult.packages/ai/src/utils/event-stream.ts:20-35—push: détection du terminal + attendant le waiter ou enfile.packages/ai/src/utils/event-stream.ts:37-47—end: force done + vide les waiter.packages/ai/src/utils/event-stream.ts:49-61—[Symbol.asyncIterator]: queue d'abord, sortie sur done, sinon attend la Promise.packages/ai/src/utils/event-stream.ts:63-65—result: renvoie la Promise terminal.packages/ai/src/utils/event-stream.ts:68-82— sous-classeAssistantMessageEventStream, fige la stratégie terminal.packages/ai/src/utils/event-stream.ts:84-87— factorycreateAssistantMessageEventStream, pour les extensions.
push est la seule entrée du producer, deux chemins : attendant un waiter ou enfiler :
// 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 :
// 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 :
// 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) returnignore silencieusement; si le provider repousse après uneerror, rien ne casse, voirpackages/ai/src/utils/event-stream.ts:21-21. - end avant un événement terminal : si
end(result?)est appelé sans résultat,finalResultPromisepeut ne jamais résoudre; l'appelant doit s'assurer qu'un événement terminal a été poussé au préalable ou passer un result, voirpackages/ai/src/utils/event-stream.ts:37-41. - error comme terminal : l'
extractResultdeAssistantMessageEventStreamrenvoieevent.errorsurerror, doncresult()récupère aussi l'objet erreur sur un flux d'erreur, sans hang, voirpackages/ai/src/utils/event-stream.ts:75-77. - Plusieurs consommateurs : plusieurs
for awaiten 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.