Anthropic SSE 實作:從位元組到 AssistantMessage
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 拿到增量。
職責
- SSE 位元組流解碼:
iterateSseMessages把ReadableStream<Uint8Array>轉成ServerSentEventasync generator,支援\r\n/\n/ 跨 chunk 拼接。見packages/ai/src/providers/anthropic.ts:321-369。 - 單行解碼:
decodeSseLine按 SSE 規範處理event:/data:/:註解 / 空行 flush。見packages/ai/src/providers/anthropic.ts:266-290。 - 事件轉型別化物件:
iterateAnthropicEvents把 SSE data 用parseJsonWithRepair解析成RawMessageStreamEvent,過濾非 Anthropic 事件,校驗message_start/message_stop配對。見packages/ai/src/providers/anthropic.ts:380-419。 - 串流主入口:
streamAnthropic建 client、buildParams、調client.messages.create({stream:true})、把事件 push 到AssistantMessageEventStream。見packages/ai/src/providers/anthropic.ts:421-510。 - Simple 版本:
streamSimpleAnthropic把SimpleStreamOptions翻譯成完整AnthropicOptions,區分 adaptive thinking(Opus 4.6 / Sonnet 4.6)和 budget-based thinking。見packages/ai/src/providers/anthropic.ts:717-756。 - 參數建構:
buildParams拼 messages / system / tools / thinking,OAuth token 強制注入 Claude Code 身份。見packages/ai/src/providers/anthropic.ts:868-925。 - 訊息轉換:
convertMessages把內部Message[]轉成 AnthropicMessageParam[],處理 text/image/tool 塊、cache_control、surrogate 清理。見packages/ai/src/providers/anthropic.ts:978-1070。
設計動機
為什麼不直接用 @anthropic-ai/sdk 的 client.messages.stream()?因為 pi-ai 要自己控制事件形狀——所有 provider 最終都 push 同一種 AssistantMessageEvent,消費側 for await 一套程式碼處理所有 provider。SDK 自帶的 stream helper 給的是 SDK 自己的事件型別,再轉一道反而繞。自己跑 SSE 還能在解析失敗時把 sse.raw 一起塞進錯誤訊息,排查 CDN 截斷、代理改寫回應這類問題好用。
為什麼把 convertMessages 和 buildParams 拆開?因為 OAuth token 與普通 API key 走的 system prompt 建構不同,buildParams 決定要不要注入 Claude Code 身份,convertMessages 只管單條訊息格式。兩者各管一層,改 cache 策略不動訊息轉換,改訊息轉換不動 OAuth 邏輯。
關鍵檔案
packages/ai/src/providers/anthropic.ts:421-510—streamAnthropic主體,async IIFE 內建 client + 串流消費。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))主消費迴圈。packages/ai/src/providers/anthropic.ts:321-369—iterateSseMessages位元組流切行。packages/ai/src/providers/anthropic.ts:266-290—decodeSseLine單行解析。packages/ai/src/providers/anthropic.ts:380-419—iterateAnthropicEventsSSE → 型別化事件 + 配對校驗。packages/ai/src/providers/anthropic.ts:717-756—streamSimpleAnthropic,simple → full options 翻譯。packages/ai/src/providers/anthropic.ts:762-810—createClient,beta header 拼接 + OAuth 偵測。packages/ai/src/providers/anthropic.ts:868-925—buildParams,OAuth 注入 Claude Code 身份。packages/ai/src/providers/anthropic.ts:978-1070—convertMessages,訊息塊轉換。
SSE 解碼核心,空行觸發 flush 回傳完整事件:
// 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 事件迭代器在配對失敗時拋錯,避免靜默吞掉被截斷的流:
// 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 身份:
// 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 } : {}),
},
];資料流
位元組流到事件流的全鏈路:
邊界與失敗
- 流被截斷:
sawMessageStart && !sawMessageEnd拋Anthropic stream ended before message_stop,見packages/ai/src/providers/anthropic.ts:416-418。 - JSON 解析失敗:
parseJsonWithRepair先嘗試修復再拋,錯誤訊息帶sse.data和sse.raw,見packages/ai/src/providers/anthropic.ts:408-413。 - OAuth 注入:
isOAuthToken偵測sk-ant-oat前綴,強制 Claude Code 身份 system prompt,見packages/ai/src/providers/anthropic.ts:758-760與packages/ai/src/providers/anthropic.ts:883-897。 - cache_control 策略:
getCacheControl(model, cacheRetention)決定加在 system / messages 哪些塊上,見packages/ai/src/providers/anthropic.ts:54-110。 - Adaptive thinking 偵測:
supportsAdaptiveThinking(model.id)區分 Opus 4.6 / Sonnet 4.6(走 effort)與老模型(走 budget tokens),見packages/ai/src/providers/anthropic.ts:681-715。 - Temperature 與 thinking 互斥:
thinkingEnabled時跳過params.temperature,見packages/ai/src/providers/anthropic.ts:910-912。
小結
anthropic.ts 自己跑 SSE,把位元組流 → ServerSentEvent → RawMessageStreamEvent → AssistantMessageEvent 四道轉換串起來,OAuth 身份、cache_control、thinking 策略都在 buildParams / convertMessages 裡處理。事件流的資料結構看 EventStream 非同步迭代器,provider 怎麼註冊進表看 Provider 抽象與內建。