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