Skip to content

Anthropic SSE-Implementierung: von Bytes zu AssistantMessage

源码版本v0.73.1

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

  1. SSE-Byte-Stream dekodieren: iterateSseMessages wandelt ReadableStream<Uint8Array> in einen ServerSentEvent-Async-Generator um, unterstützt \r\n / \n / chunk-übergreifendes Zusammenfügen. Siehe packages/ai/src/providers/anthropic.ts:321-369.
  2. Einzelzeilen-Dekodierung: decodeSseLine behandelt event: / data: / :-Kommentare / Leerzeile-flush nach SSE-Spec. Siehe packages/ai/src/providers/anthropic.ts:266-290.
  3. Events in typisierte Objekte: iterateAnthropicEvents parst SSE-data mit parseJsonWithRepair zu RawMessageStreamEvent, filtert nicht-Anthropic-Events, prüft message_start / message_stop-Paarung. Siehe packages/ai/src/providers/anthropic.ts:380-419.
  4. Streaming-Haupteingang: streamAnthropic baut den Client, ruft buildParams auf, ruft client.messages.create({stream:true}) und pusht die Events in den AssistantMessageEventStream. Siehe packages/ai/src/providers/anthropic.ts:421-510.
  5. Simple-Variante: streamSimpleAnthropic übersetzt SimpleStreamOptions in vollständige AnthropicOptions und unterscheidet adaptive Thinking (Opus 4.6 / Sonnet 4.6) von budget-basiertem Thinking. Siehe packages/ai/src/providers/anthropic.ts:717-756.
  6. Parameter-Aufbau: buildParams setzt messages / system / tools / thinking zusammen und erzwingt bei OAuth-Token die Claude-Code-Identität. Siehe packages/ai/src/providers/anthropic.ts:868-925.
  7. Nachricht-Konvertierung: convertMessages wandelt interne Message[] in Anthropic-MessageParam[] um und behandelt Text-/Bild-/Werkzeug-Blöcke, cache_control sowie Surrogate-Bereinigung. Siehe packages/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

Kern der SSE-Dekodierung: die Leerzeile triggert den Flush und liefert das ganze Event:

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

Der Anthropic-Event-Iterator wirft bei fehlerhafter Paarung, anstatt einen abgeschnittenen Stream stumm zu schlucken:

typescript
// 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:

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 } : {}),
			},
		];

Datenfluss

Die ganze Kette vom Byte-Stream zum Event-Stream:

Grenzen und Fehler

Zusammenfassung

anthropic.ts fährt SSE selbst und verkettet vier Umwandlungen: Byte-Stream → ServerSentEventRawMessageStreamEventAssistantMessageEvent. 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.