Skip to content

Anthropic SSE 實作:從位元組到 AssistantMessage

源码版本v0.73.1

providers/anthropic.ts 是 9 家內建 provider 裡最重的一家——自己實作 SSE 增量解碼、Anthropic 事件轉 AssistantMessageEvent、OAuth/Claude Code 身份注入、cache_control 策略、thinking 與 tool call 的 block 增量拼接。它不依賴 Anthropic SDK 的串流解析,而是直接讀 response.body 的位元組流,自己跑 SSE 協定。最終產物是 push 進 AssistantMessageEventStream 的事件序列,消費者 for await 拿到增量。

職責

  1. SSE 位元組流解碼:iterateSseMessagesReadableStream<Uint8Array> 轉成 ServerSentEvent async generator,支援 \r\n / \n / 跨 chunk 拼接。見 packages/ai/src/providers/anthropic.ts:321-369
  2. 單行解碼:decodeSseLine 按 SSE 規範處理 event: / data: / : 註解 / 空行 flush。見 packages/ai/src/providers/anthropic.ts:266-290
  3. 事件轉型別化物件:iterateAnthropicEvents 把 SSE data 用 parseJsonWithRepair 解析成 RawMessageStreamEvent,過濾非 Anthropic 事件,校驗 message_start / message_stop 配對。見 packages/ai/src/providers/anthropic.ts:380-419
  4. 串流主入口:streamAnthropic 建 client、buildParams、調 client.messages.create({stream:true})、把事件 push 到 AssistantMessageEventStream。見 packages/ai/src/providers/anthropic.ts:421-510
  5. Simple 版本:streamSimpleAnthropicSimpleStreamOptions 翻譯成完整 AnthropicOptions,區分 adaptive thinking(Opus 4.6 / Sonnet 4.6)和 budget-based thinking。見 packages/ai/src/providers/anthropic.ts:717-756
  6. 參數建構:buildParams 拼 messages / system / tools / thinking,OAuth token 強制注入 Claude Code 身份。見 packages/ai/src/providers/anthropic.ts:868-925
  7. 訊息轉換:convertMessages 把內部 Message[] 轉成 Anthropic MessageParam[],處理 text/image/tool 塊、cache_control、surrogate 清理。見 packages/ai/src/providers/anthropic.ts:978-1070

設計動機

為什麼不直接用 @anthropic-ai/sdkclient.messages.stream()?因為 pi-ai 要自己控制事件形狀——所有 provider 最終都 push 同一種 AssistantMessageEvent,消費側 for await 一套程式碼處理所有 provider。SDK 自帶的 stream helper 給的是 SDK 自己的事件型別,再轉一道反而繞。自己跑 SSE 還能在解析失敗時把 sse.raw 一起塞進錯誤訊息,排查 CDN 截斷、代理改寫回應這類問題好用。

為什麼把 convertMessagesbuildParams 拆開?因為 OAuth token 與普通 API key 走的 system prompt 建構不同,buildParams 決定要不要注入 Claude Code 身份,convertMessages 只管單條訊息格式。兩者各管一層,改 cache 策略不動訊息轉換,改訊息轉換不動 OAuth 邏輯。

關鍵檔案

SSE 解碼核心,空行觸發 flush 回傳完整事件:

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

Anthropic 事件迭代器在配對失敗時拋錯,避免靜默吞掉被截斷的流:

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

OAuth token 走特殊 system prompt,強制帶 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 } : {}),
			},
		];

資料流

位元組流到事件流的全鏈路:

邊界與失敗

小結

anthropic.ts 自己跑 SSE,把位元組流 → ServerSentEventRawMessageStreamEventAssistantMessageEvent 四道轉換串起來,OAuth 身份、cache_control、thinking 策略都在 buildParams / convertMessages 裡處理。事件流的資料結構看 EventStream 非同步迭代器,provider 怎麼註冊進表看 Provider 抽象與內建