Implémentation Anthropic SSE : des octets à AssistantMessage
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
- Décodage du flux d'octets SSE :
iterateSseMessagestransforme unReadableStream<Uint8Array>en un async generator deServerSentEvent, gère\r\n/\n/ le collage inter-chunk. Voirpackages/ai/src/providers/anthropic.ts:321-369. - Décodage ligne par ligne :
decodeSseLinetraite selon la spec SSE lesevent:/data:/ commentaires:/ lignes vides qui flush. Voirpackages/ai/src/providers/anthropic.ts:266-290. - Événements → objets typés :
iterateAnthropicEventsparse la data SSE viaparseJsonWithRepairenRawMessageStreamEvent, filtre les événements non Anthropic, valide l'appariementmessage_start/message_stop. Voirpackages/ai/src/providers/anthropic.ts:380-419. - Entrée principale du streaming :
streamAnthropicconstruit le client, appelle buildParams, invoqueclient.messages.create({stream:true}), et pousse les événements dansAssistantMessageEventStream. Voirpackages/ai/src/providers/anthropic.ts:421-510. - Version Simple :
streamSimpleAnthropictraduitSimpleStreamOptionsenAnthropicOptionscomplet, distingue adaptive thinking (Opus 4.6 / Sonnet 4.6) et budget-based thinking. Voirpackages/ai/src/providers/anthropic.ts:717-756. - Construction des paramètres :
buildParamsassemble messages / system / tools / thinking; le token OAuth force l'injection de l'identité Claude Code. Voirpackages/ai/src/providers/anthropic.ts:868-925. - Conversion des messages :
convertMessagestransforme lesMessage[]internes enMessageParam[]Anthropic, gère les blocs text/image/tool, le cache_control, le nettoyage des surrogate. Voirpackages/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
packages/ai/src/providers/anthropic.ts:421-510— corps destreamAnthropic, IIFE async qui construit le client + consomme le flux.packages/ai/src/providers/anthropic.ts:487-489—client.messages.create({ ...params, stream: true })+stream.push({ type: "start" }).packages/ai/src/providers/anthropic.ts:494-494—for await (const event of iterateAnthropicEvents(response, signal))boucle principale de consommation.packages/ai/src/providers/anthropic.ts:321-369—iterateSseMessagesdécoupe le flux d'octets en lignes.packages/ai/src/providers/anthropic.ts:266-290—decodeSseLineparse une ligne.packages/ai/src/providers/anthropic.ts:380-419—iterateAnthropicEventsSSE → événements typés + validation d'appariement.packages/ai/src/providers/anthropic.ts:717-756—streamSimpleAnthropic, traduction simple → full options.packages/ai/src/providers/anthropic.ts:762-810—createClient, assemblage du beta header + détection OAuth.packages/ai/src/providers/anthropic.ts:868-925—buildParams, injection Claude Code par OAuth.packages/ai/src/providers/anthropic.ts:978-1070—convertMessages, conversion des blocs de message.
Le cœur du décodage SSE : la ligne vide déclenche le flush et renvoie l'événement complet :
// 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é :
// 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 :
// 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
- Flux tronqué :
sawMessageStart && !sawMessageEndjetteAnthropic stream ended before message_stop, voirpackages/ai/src/providers/anthropic.ts:416-418. - Échec de parsing JSON :
parseJsonWithRepairtente d'abord une réparation avant de jeter, le message d'erreur embarquesse.dataetsse.raw, voirpackages/ai/src/providers/anthropic.ts:408-413. - Injection OAuth :
isOAuthTokendétecte le préfixesk-ant-oatet force le system prompt d'identité Claude Code, voirpackages/ai/src/providers/anthropic.ts:758-760etpackages/ai/src/providers/anthropic.ts:883-897. - Stratégie cache_control :
getCacheControl(model, cacheRetention)décide sur quels blocs system / messages l'ajouter, voirpackages/ai/src/providers/anthropic.ts:54-110. - Détection adaptive thinking :
supportsAdaptiveThinking(model.id)distingue Opus 4.6 / Sonnet 4.6 (qui passent par effort) des modèles plus anciens (qui passent par budget tokens), voirpackages/ai/src/providers/anthropic.ts:681-715. - Mutuellement exclusifs, temperature et thinking : quand
thinkingEnabledest vrai,params.temperatureest sauté, voirpackages/ai/src/providers/anthropic.ts:910-912.
Récapitulatif
anthropic.ts fait tourner lui-même le SSE, enchaînant quatre conversions flux d'octets → ServerSentEvent → RawMessageStreamEvent → AssistantMessageEvent; 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.