Skip to content

Implémentation Anthropic SSE : des octets à AssistantMessage

源码版本v0.73.1

providers/anthropic.ts est le plus lourd des 9 provider intégrés—il implémente lui-même le décodage incrémental SSE, la conversion des événements Anthropic en AssistantMessageEvent, l'injection d'identité OAuth/Claude Code, la stratégie cache_control, et l'assemblage incrémental des blocs thinking et tool call. Il ne s'appuie pas sur l'analyse streaming du SDK Anthropic; il lit directement le flux d'octets de response.body et fait tourner lui-même le protocole SSE. Le produit final est une séquence d'événements poussés dans AssistantMessageEventStream, que le consommateur récupère en incrémental via for await.

Responsabilités

  1. Décodage du flux d'octets SSE : iterateSseMessages transforme un ReadableStream<Uint8Array> en un async generator de ServerSentEvent, gère \r\n / \n / le collage inter-chunk. Voir packages/ai/src/providers/anthropic.ts:321-369.
  2. Décodage ligne par ligne : decodeSseLine traite selon la spec SSE les event: / data: / commentaires : / lignes vides qui flush. Voir packages/ai/src/providers/anthropic.ts:266-290.
  3. Événements → objets typés : iterateAnthropicEvents parse la data SSE via parseJsonWithRepair en RawMessageStreamEvent, filtre les événements non Anthropic, valide l'appariement message_start / message_stop. Voir packages/ai/src/providers/anthropic.ts:380-419.
  4. Entrée principale du streaming : streamAnthropic construit le client, appelle buildParams, invoque client.messages.create({stream:true}), et pousse les événements dans AssistantMessageEventStream. Voir packages/ai/src/providers/anthropic.ts:421-510.
  5. Version Simple : streamSimpleAnthropic traduit SimpleStreamOptions en AnthropicOptions complet, distingue adaptive thinking (Opus 4.6 / Sonnet 4.6) et budget-based thinking. Voir packages/ai/src/providers/anthropic.ts:717-756.
  6. Construction des paramètres : buildParams assemble messages / system / tools / thinking; le token OAuth force l'injection de l'identité Claude Code. Voir packages/ai/src/providers/anthropic.ts:868-925.
  7. Conversion des messages : convertMessages transforme les Message[] internes en MessageParam[] Anthropic, gère les blocs text/image/tool, le cache_control, le nettoyage des surrogate. Voir packages/ai/src/providers/anthropic.ts:978-1070.

Motivation de design

Pourquoi ne pas utiliser directement client.messages.stream() du @anthropic-ai/sdk ? Parce que pi-ai veut contrôler lui-même la forme des événements—tous les provider finissent par pousser le même type AssistantMessageEvent, le consommateur gère tous les provider avec un seul for await. Le helper de stream du SDK produit ses propres types d'événements, qu'il faudrait re-convertir, ce qui rajoute un détour. Faire tourner SSE soi-même permet aussi, en cas d'échec de parsing, d'inclure sse.raw dans le message d'erreur—pratique pour diagnostiquer les tronquatures CDN, les réponses réécrites par un proxy.

Pourquoi séparer convertMessages et buildParams ? Parce que le token OAuth et la clé API classique ne construisent pas le system prompt de la même façon; buildParams décide d'injecter ou non l'identité Claude Code, convertMessages ne s'occupe que du format d'un message. Chacun gère sa couche : modifier la stratégie de cache ne touche pas la conversion des messages, modifier la conversion des messages ne touche pas la logique OAuth.

Fichiers clés

Le cœur du décodage SSE : la ligne vide déclenche le flush et renvoie l'événement complet :

typescript
// packages/ai/src/providers/anthropic.ts:266-290
function decodeSseLine(line: string, state: SseDecoderState): ServerSentEvent | null {
	if (line === "") {
		return flushSseEvent(state);
	}
	state.raw.push(line);
	if (line.startsWith(":")) {
		return null;
	}
	const delimiterIndex = line.indexOf(":");
	const fieldName = delimiterIndex === -1 ? line : line.slice(0, delimiterIndex);
	let value = delimiterIndex === -1 ? "" : line.slice(delimiterIndex + 1);
	if (value.startsWith(" ")) {
		value = value.slice(1);
	}
	if (fieldName === "event") {
		state.event = value;
	} else if (fieldName === "data") {
		state.data.push(value);
	}
	return null;
}

L'itérateur d'événements Anthropic jette en cas d'appariement rompu, pour ne pas avaler silencieusement un flux tronqué :

typescript
// packages/ai/src/providers/anthropic.ts:416-418
	if (sawMessageStart && !sawMessageEnd) {
		throw new Error("Anthropic stream ended before message_stop");
	}

Le token OAuth passe par un system prompt spécial qui force l'identité Claude Code :

typescript
// packages/ai/src/providers/anthropic.ts:883-890
	if (isOAuthToken) {
		params.system = [
			{
				type: "text",
				text: "You are Claude Code, Anthropic's official CLI for Claude.",
				...(cacheControl ? { cache_control: cacheControl } : {}),
			},
		];

Flux de données

Toute la chaîne du flux d'octets au flux d'événements :

Frontières et échecs

Récapitulatif

anthropic.ts fait tourner lui-même le SSE, enchaînant quatre conversions flux d'octets → ServerSentEventRawMessageStreamEventAssistantMessageEvent; l'identité OAuth, le cache_control et la stratégie thinking se gèrent dans buildParams / convertMessages. La structure de données du flux d'événements se lit dans Itérateur async EventStream, et la façon dont un provider s'enregistre dans la table dans Abstraction provider et intégrations.