Anthropic SSE-Implementierung: von Bytes zu AssistantMessage
providers/anthropic.ts ist der schwerste der 9 eingebauten Provider — er implementiert selbst inkrementelle SSE-Dekodierung, übersetzt Anthropic-Events in AssistantMessageEvent, injiziert OAuth/Claude-Code-Identität, regelt cache_control, fügt Thinking- und Tool-Call-Blöcke inkrementell zusammen. Er nutzt nicht den Streaming-Parser des Anthropic-SDK, sondern liest den Byte-Stream von response.body direkt und spricht selbst das SSE-Protokoll. Endprodukt ist eine Folge von Events, die in den AssistantMessageEventStream gepusht werden; der Konsument erhält über for await die Inkremente.
Verantwortung
- SSE-Byte-Stream dekodieren:
iterateSseMessageswandeltReadableStream<Uint8Array>in einenServerSentEvent-Async-Generator um, unterstützt\r\n/\n/ chunk-übergreifendes Zusammenfügen. Siehepackages/ai/src/providers/anthropic.ts:321-369. - Einzelzeilen-Dekodierung:
decodeSseLinebehandeltevent:/data:/:-Kommentare / Leerzeile-flush nach SSE-Spec. Siehepackages/ai/src/providers/anthropic.ts:266-290. - Events in typisierte Objekte:
iterateAnthropicEventsparst SSE-data mitparseJsonWithRepairzuRawMessageStreamEvent, filtert nicht-Anthropic-Events, prüftmessage_start/message_stop-Paarung. Siehepackages/ai/src/providers/anthropic.ts:380-419. - Streaming-Haupteingang:
streamAnthropicbaut den Client, ruft buildParams auf, ruftclient.messages.create({stream:true})und pusht die Events in denAssistantMessageEventStream. Siehepackages/ai/src/providers/anthropic.ts:421-510. - Simple-Variante:
streamSimpleAnthropicübersetztSimpleStreamOptionsin vollständigeAnthropicOptionsund unterscheidet adaptive Thinking (Opus 4.6 / Sonnet 4.6) von budget-basiertem Thinking. Siehepackages/ai/src/providers/anthropic.ts:717-756. - Parameter-Aufbau:
buildParamssetzt messages / system / tools / thinking zusammen und erzwingt bei OAuth-Token die Claude-Code-Identität. Siehepackages/ai/src/providers/anthropic.ts:868-925. - Nachricht-Konvertierung:
convertMessageswandelt interneMessage[]in Anthropic-MessageParam[]um und behandelt Text-/Bild-/Werkzeug-Blöcke, cache_control sowie Surrogate-Bereinigung. Siehepackages/ai/src/providers/anthropic.ts:978-1070.
Entwurfsmotivation
Warum nicht einfach client.messages.stream() aus @anthropic-ai/sdk? Weil pi-ai die Event-Form selbst kontrollieren will — alle Provider pushen am Ende dasselbe AssistantMessageEvent, sodass der Konsument mit einem Code-Pfad über for await alle Provider behandelt. Die Stream-Helfer des SDK liefern die SDK-eigenen Event-Typen; ein erneutes Umwandeln wäre ein Umweg. SSE selbst zu fahren erlaubt außerdem, bei Parse-Fehlern sse.raw in die Fehlermeldung zu packen — nützlich beim Diagnostizieren von CDN-Abbrüchen oder Proxys, die die Antwort verändern.
Warum convertMessages und buildParams trennen? Weil OAuth-Token und normale API-Key unterschiedliche System-Prompt-Konstruktion haben; buildParams entscheidet, ob die Claude-Code-Identität injiziert wird, convertMessages kümmert sich nur um das Format einzelner Nachrichten. Beide haben je eine Schicht: Wer die Cache-Strategie ändert, berührt nicht die Nachrichtenkonvertierung und umgekehrt die OAuth-Logik nicht.
Wichtige Dateien
packages/ai/src/providers/anthropic.ts:421-510—streamAnthropic-Body, in einer async-IIFE Client gebaut + Streaming konsumiert.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))als Hauptschleife.packages/ai/src/providers/anthropic.ts:321-369—iterateSseMessages, Byte-Stream in Zeilen zerlegen.packages/ai/src/providers/anthropic.ts:266-290—decodeSseLine, Parsing einzelner Zeilen.packages/ai/src/providers/anthropic.ts:380-419—iterateAnthropicEvents, SSE → typisierte Events + Paarungsprüfung.packages/ai/src/providers/anthropic.ts:717-756—streamSimpleAnthropic, Simple → Full-Options-Übersetzung.packages/ai/src/providers/anthropic.ts:762-810—createClient, Beta-Header + OAuth-Erkennung.packages/ai/src/providers/anthropic.ts:868-925—buildParams, OAuth injiziert Claude-Code-Identität.packages/ai/src/providers/anthropic.ts:978-1070—convertMessages, Konvertierung der Nachrichtenblöcke.
Kern der SSE-Dekodierung: die Leerzeile triggert den Flush und liefert das ganze Event:
// 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;
}Der Anthropic-Event-Iterator wirft bei fehlerhafter Paarung, anstatt einen abgeschnittenen Stream stumm zu schlucken:
// packages/ai/src/providers/anthropic.ts:416-418
if (sawMessageStart && !sawMessageEnd) {
throw new Error("Anthropic stream ended before message_stop");
}OAuth-Token bekommt einen besonderen System-Prompt und erzwingt die Claude-Code-Identität:
// 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 } : {}),
},
];Datenfluss
Die ganze Kette vom Byte-Stream zum Event-Stream:
Grenzen und Fehler
- Stream abgeschnitten:
sawMessageStart && !sawMessageEndwirftAnthropic stream ended before message_stop; siehepackages/ai/src/providers/anthropic.ts:416-418. - JSON-Parse schlägt fehl:
parseJsonWithRepairversucht erst Reparatur, dann wird geworfen; die Fehlermeldung enthältsse.dataundsse.raw; siehepackages/ai/src/providers/anthropic.ts:408-413. - OAuth-Injektion:
isOAuthTokenerkennt dassk-ant-oat-Präfix und erzwingt den Claude-Code-System-Prompt; siehepackages/ai/src/providers/anthropic.ts:758-760undpackages/ai/src/providers/anthropic.ts:883-897. - cache_control-Strategie:
getCacheControl(model, cacheRetention)entscheidet, an welche Blöcke in system / messages es angehängt wird; siehepackages/ai/src/providers/anthropic.ts:54-110. - Adaptive-Thinking-Erkennung:
supportsAdaptiveThinking(model.id)unterscheidet Opus 4.6 / Sonnet 4.6 (gehen über effort) von älteren Modellen (gehen über Budget-Tokens); siehepackages/ai/src/providers/anthropic.ts:681-715. - Temperatur und Thinking schließen sich aus: bei
thinkingEnabledwirdparams.temperatureübersprungen; siehepackages/ai/src/providers/anthropic.ts:910-912.
Zusammenfassung
anthropic.ts fährt SSE selbst und verkettet vier Umwandlungen: Byte-Stream → ServerSentEvent → RawMessageStreamEvent → AssistantMessageEvent. OAuth-Identität, cache_control und Thinking-Strategie werden in buildParams / convertMessages behandelt. Die Datenstruktur des Event-Streams steht in EventStream Async-Iterator; wie ein Provider in die Tabelle kommt, in Provider-Abstraktion und Built-ins.