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 抽象与内置。