8.6 接入真实模型(可选)
本页分析版本earendil-works/pi@c13ffe12026-07-30本章是全书唯一需要 API Key 的内容,而且完全可选
step01 到 step07 一个字节都不会发到网上。本章对应的 step08 也只在环境变量 ANTHROPIC_API_KEY 存在时才真的发请求;没有 Key 时它自动回退到 FakeModel,行为与 step07 逐字相同。
本章的实践任务 npm run demo 同样是完全离线的:它不读 Key、不建连接,却把本章的三个知识点全部演示了一遍。没有 API Key 的读者可以照常读完、照常跑完,唯一少掉的是「模型真的回了你一句话」这一下——而那一下并不是知识点。
本章解决什么问题:step03 定下
StreamFn这个抽象时,我们说「以后换成真模型,上层不用改」。这句话到现在为止还是一张空头支票。本章把它兑现:写一个直连 Anthropic Messages API 的StreamFn实现,然后确认 Agent Loop(Agent 循环)、工具、会话(Session)、取消、扩展(Extension)一行都不用改。顺带把「流式响应在网络上到底长什么样」这件事看清楚。 前置知识:8.2 消息历史与流式事件(StreamFn与StreamEvent的契约)、8.5 取消与 Extension Hook、2.6 异步迭代器与 for await、3.3 流式输出(Streaming)。 学习目标:读完并跑完 demo 后你能—— ① 说出一个 Provider(模型服务提供方)适配器要做的三件事,并指出它们各自对应哪个函数; ② 读懂一段 Server-Sent Events(SSE,服务器发送事件)报文,并解释「网络分块」与「协议事件」为什么必须解耦; ③ 解释工具调用的参数为什么是分片到达的 JSON 字符串,以及为什么必须等到块结束才能JSON.parse; ④ 说出 API Key 的三条最低安全要求,并指出 Pi 把凭证存在哪、用什么权限存; ⑤ 用「换实现时上层要改多少行」这把尺子,判断一个抽象值不值。
建立直觉:一个抽象值不值,看换实现时上层要改多少行
抽象是有成本的。多一层 StreamFn,就多一个类型要维护、多一次跳转要读、多一层调试时要穿过的壳。这笔成本什么时候回本?只有一个时刻能回答:真的换实现的那一天。
现在就是那一天。回顾 step03 定下的那行类型(src/types.ts:102):
export type StreamFn = (context: Context, options?: StreamOptions) => AsyncIterable<StreamEvent>;它规定的东西少得可怜:进去一个 Context,出来一串 StreamEvent,外加一条不写在签名里的契约——不许 throw,错误编码成 error 事件(src/types.ts:99-100 的注释)。正因为它规定得少,FakeModel 和真模型才可能是同一个类型。
于是 step08 的接线点只有一处。src/main.ts:72:
const { choice, streamFn } = createModel(env.ANTHROPIC_API_KEY);step07 那一行原本是 const streamFn = createFakeModel(rules, { delayMs: 40 });。除此之外 main.ts 里的实质改动只剩两处:import 换了一个来源,启动横幅多打一行「模型:……」。其余的差异都是注释与系统提示词措辞——没有一行是被迫改的。
src/ 下的其他文件(agent-loop.ts、tools/、session.ts、extensions.ts、types.ts、render.ts、fake-model.ts、rules.ts)与 step07 逐字相同,这一点本章的实践任务会让你亲手 diff 一遍。新增的只有 src/providers/ 这个目录。
图 8.6-1 step08 只换了虚线以下的部分
实线是 step01–step07 已经写完的代码,虚线以下是本章新增的两个可替换实现。
这张图要看的是虚线的位置:它落在 StreamFn 这一层,而不是落在 runAgentLoop 里。如果当初 Agent Loop 直接调 fetch,虚线就得画在 agent-loop.ts 内部,换模型意味着改循环——而循环是同时被工具、取消、扩展依赖的那段代码,动它的代价远大于动一个叶子模块。图上 FAKE 对应 src/fake-model.ts,REAL 对应本章新增的 src/providers/anthropic.ts,两者都返回 StreamFn 类型的函数,LOOP 拿到哪一个都不知情。
那么这个新增的叶子模块要做什么?只有三件事。
Context(系统提示词、消息数组、工具清单)翻译成对方 HTTP 接口要的 JSON 形状。2. 流式解析:把 HTTP 响应体的字节流,切成一个个完整的协议事件。
3. 事件翻译:把对方的协议事件,翻译成我们统一的
StreamEvent。这三件事在 step08 里分别是
toRequestBody、readSseMessages、toStreamEvents;在 Pi 里它们住在 packages/ai/src/api/ 下同一个文件的三段里。换一家服务商,只有这三件事要重写。最小示例:三个离线场景
step08 的目录相比 step07 只多了一个 src/providers/:
labs/mini-agent-harness/step08-real-provider/
├── src/
│ ├── providers/
│ │ ├── sse.ts ← 只做 SSE 格式解析,不认识 Anthropic
│ │ ├── anthropic.ts ← 适配器的三件事
│ │ ├── recorded-stream.ts ← 手写的样例 SSE,供离线演示
│ │ └── select.ts ← 有 Key 用真模型,没 Key 用 FakeModel
│ ├── main.ts ← 相比 step07 只换掉了造 streamFn 的那一行
│ ├── demo.ts ← 本章的离线演示
│ └── …(其余文件与 step07 完全相同)
└── expected-output.txt运行:
cd labs/mini-agent-harness/step08-real-provider
npm install
npm run demo # 完全离线,不需要 API Keydemo 的完整输出与 expected-output.txt 逐字相同(写作者真实运行核验过)。下面按三个场景拆开讲。
场景 1:Context → 请求体
demo 拿一段写死的三条消息上下文(一次用户提问 + 一次带工具调用的助手回复 + 一条工具结果),只调用 toRequestBody,把结果打印出来——不发送。真实输出(节选自 expected-output.txt 第 5–66 行):
{
"model": "claude-opus-5",
"max_tokens": 16000,
"system": "你是一个会用工具的助手。",
"messages": [
{ "role": "user", "content": [ { "type": "text", "text": "算一下 12*8" } ] },
{ "role": "assistant", "content": [
{ "type": "text", "text": "这个我算一下。" },
{ "type": "tool_use", "id": "toolu_demo", "name": "calc",
"input": { "expression": "12*8" } } ] },
{ "role": "user", "content": [
{ "type": "tool_result", "tool_use_id": "toolu_demo", "content": "96" } ] }
],
"tools": [ { "name": "calc", "description": "…", "input_schema": { "…": "…" } } ],
"stream": true
}(tools 数组里的 description 与 input_schema 在真实输出中是完整展开的,见 expected-output.txt 第 47–64 行;这里为省版面折叠。)
对着我们自己的类型看,翻译关系是这样的:
我们的表示(src/types.ts) | Anthropic 请求体里的位置 |
|---|---|
Context.systemPrompt | 顶层 system 字段,不是 messages 里的一条 |
UserMessage | role: "user",内容是 text 块 |
AssistantMessage 的 TextContent | role: "assistant",内容是 text 块 |
ToolCallContent | role: "assistant" 的 tool_use 块;arguments 改名成 input |
ToolResultMessage(第三种 role) | role: "user" 的 tool_result 块;toolCallId 改名成 tool_use_id |
ToolSpec.parameters | tools[].input_schema |
最值得停一下的是倒数第二行。我们的 toolResult 是独立的第三种 role(step05 定的),而 Anthropic 把工具结果塞在用户消息里。同一个概念、两种摆法——这就是适配器存在的全部理由。翻译代码在 src/providers/anthropic.ts:81-92:
case "toolResult":
return {
role: "user",
content: [
{
type: "tool_result",
tool_use_id: message.toolCallId,
content: message.content.map((block) => block.text).join(""),
...(message.isError ? { is_error: true } : {}),
},
],
};翻译完还有一步收尾:相邻的同角色消息必须合并成一条(src/providers/anthropic.ts:107-116)。
const messages: ApiMessage[] = [];
for (const message of context.messages) {
const converted = toApiMessage(message);
const previous = messages[messages.length - 1];
if (previous && previous.role === converted.role) {
previous.content.push(...converted.content);
} else {
messages.push(converted);
}
}为什么非合并不可?模型一次回复里可以发起多个 Tool Call(step05 讲过),于是我们会产生多条 toolResult 消息。如果原样翻译成多条 role: "user" 的消息,等于告诉模型「这些结果是分好几轮陆续给你的」。模型下次就会倾向于一次只调一个工具——你在数据里教会了它一件你并不想教的事。
场景 2:样例 SSE → StreamEvent
这是本章的核心。先看真实输出(expected-output.txt 第 72–82 行):
样例响应共 1217 个字符,被切成 33 个网络块。
第 1 块的内容是:"event: message_start\ndata: {\"type\":\"m"
——切点落在事件中间,所以解析器必须自己缓冲。
start
text_delta "这个我"
text_delta "算一下。"
toolcall calc {"expression":"12*8"}
done stopReason=toolUse,2 个内容块
工具参数是分片到达的("{\"expres" + "sion\":\"12*8\"}"),拼完整才能 JSON.parse。先看输入长什么样。 SSE 是一种极简的文本协议:一个普通的 HTTP 请求,服务端不一次性写完响应体,而是沿着这条不关闭的连接持续写文本。每个「事件」由若干行组成,事件之间用一个空行分隔:
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"这个我"}}
event: content_block_stop
data: {"type":"content_block_stop","index":0}event: 行是事件名,data: 行是负载(Anthropic 的负载是一段 JSON)。规范允许一个事件有多行 data:,也允许以冒号开头的注释行(常被用作心跳包)——src/providers/sse.ts:74-92 的 parseBlock 把这两种情况都处理了。
再看难点在哪。 demo 那行「共 1217 个字符,被切成 33 个网络块」不是装饰。toChunks(src/providers/recorded-stream.ts:56-62)把这段文本按 37 个字符一刀切开,代码注释说明这个数是刻意挑的——它不整除任何事件的长度,于是每一刀都落在奇怪的地方。第 1 块正好是 37 个字符:event: message_start(20)+ 换行(1)+ data: {"type":"m(16)——它在 JSON 的中间断掉了。
这不是人为刁难,这就是网络的真实样子:网络给你的分块和协议的分隔毫无关系。一个事件可能被切成三块到达,三个事件也可能挤在一块里到达。所以解析器必须自己缓冲。
图 8.6-2 从网络字节到 StreamEvent 的四段流水线
每个箭头上标的是「跨过这条边界的东西是什么」。四段各自是一个异步生成器,首尾相接。
这张图要注意的是每段各解决一个「切点问题」,且顺序不能换:
- 第一段
decodeChunks(src/providers/sse.ts:31-39)解决的是字节层面的切点。一个 UTF-8 汉字占 3 个字节,网络完全可能把「我」这个字切在两块中间。new TextDecoder()配decode(chunk, { stream: true })会把不完整的字节留在内部缓冲区,等下一块到齐再吐出来;循环结束后还要decoder.decode()无参调用一次,把尾巴冲出来。少了{ stream: true },你会在日志里看到随机出现的乱码字符,而且极难复现——因为它取决于网络怎么切。 - 第二段
readSseMessages(src/providers/sse.ts:49-71)解决的是协议层面的切点:
let buffer = "";
for await (const chunk of chunks) {
// 统一换行符:有的服务端用 \r\n,先归一化能省掉一堆分支。
buffer += chunk.replace(/\r\n/g, "\n");
let separator = buffer.indexOf("\n\n");
while (separator !== -1) {
const block = buffer.slice(0, separator);
buffer = buffer.slice(separator + 2);
const message = parseBlock(block);
if (message) yield message;
separator = buffer.indexOf("\n\n");
}
}注意那个 while 而不是 if:一块里可能挤着好几个完整事件,必须循环取到取不出为止。循环外还有一次收尾(sse.ts:69-70),处理「最后一个事件没有以空行结尾」的情况。这个 buffer + 反复找分隔符 + 收尾 的三件套,是所有流式协议解析的共同套路,换成 WebSocket 的分帧、换成换行分隔的 JSONL,骨架都一样。
- 第三段
toStreamEvents才开始认识 Anthropic。
Anthropic 的事件顺序是固定的(src/providers/anthropic.ts:179-184 的注释):
message_start
→ content_block_start / content_block_delta × N / content_block_stop(每个内容块一组)
→ message_delta(带 stop_reason)
→ message_stop翻译成我们的 StreamEvent 时,并不是一对一:
| Anthropic 事件 | 我们发出的 StreamEvent |
|---|---|
message_start | start |
content_block_start(text 或 tool_use) | 不发事件,只在 pending 里登记这个块 |
content_block_delta 的 text_delta | text_delta |
content_block_delta 的 input_json_delta | 不发事件,只把 JSON 分片累积进字符串 |
content_block_stop(text 块) | 不发事件,把整段文本收进 content |
content_block_stop(tool_use 块) | toolcall(此时才 JSON.parse) |
message_delta 的 stop_reason | 不发事件,记下来留给最后 |
message_stop | done,带上攒了一路的 AssistantMessage |
SSE 的 event: error,或 JSON 解析失败 | error |
一半的上游事件不产生下游事件——它们只是在攒状态。攒状态的容器是 toStreamEvents 开头这四行(src/providers/anthropic.ts:192-195):
const pending = new Map<number, PendingBlock>();
const content: ContentBlock[] = [];
let stopReason: StopReason = "stop";
let started = false;pending 用 index 作键,是因为内容块可以交错:模型可以先开一个文本块、再开一个工具块,两个块的 delta 事件混着到。没有 index 就没法把 delta 归位。
最容易踩的坑是工具参数。 它以 input_json_delta 分片到达,样例里被切成 {"expres 和 sion":"12*8"} 两片——单看第一片根本不是合法 JSON。所以 delta 阶段只能累积不能解析(src/providers/anthropic.ts:242-245):
} else if (payload.delta?.type === "input_json_delta" && current.kind === "toolCall") {
// 只累积,不解析:现在拿到的多半是半截 JSON。
current.json += payload.delta.partial_json ?? "";
}一直等到 content_block_stop 才动手(anthropic.ts:257-267,随后在 268-276 行组装出 ToolCallContent 并发出 toolcall 事件):
let args: Record<string, unknown>;
try {
// 空字符串表示这个工具没有参数,对应 {}。
args =
current.json.trim().length > 0
? (JSON.parse(current.json) as Record<string, unknown>)
: {};
} catch {
yield { type: "error", error: `工具参数不是合法 JSON:${current.json}` };
return;
}这段脏活干完,上层才能拿到 arguments 是对象的 ToolCallContent。回头看 step03 定 ToolCallContent.arguments: Record<string, unknown> 而不是 string 的决定——那个决定的成本,就是此刻这十几行代码;收益是从 agent-loop.ts 到 tools/validate.ts 到 render.ts,没有任何一处需要知道参数曾经是分片的。抽象不消灭复杂度,只决定复杂度住在哪一间房。
src/providers/anthropic.ts:148-153 声明的 AnthropicEvent 里,type、index、content_block、delta 全部带 ?。这不是偷懒,是故意的:这份对象来自 JSON.parse,运行时没有任何东西保证它的形状——服务端可能加了新字段、代理可能删了字段、网络可能给你半截数据。声明成可选,TypeScript 就会在每次取值时逼你判空。和 step05 的 validate 是同一个道理:边界之外的数据不可信。 还有一处值得看的是停止原因的映射(src/providers/anthropic.ts:160-171):tool_use → toolUse,refusal → error,其余(包括 max_tokens)一律 → stop。我们的教学版只有四种 StopReason,所以「被 max_tokens 截断」和「正常说完」归成了一类——代码注释里写明了这个取舍:max_tokens 的内容本身是有效的,只是没说完。Pi 在这一点上分得更细,下一节会看到。
场景 3:换模型不影响上层
选模型的规则被单独抽成一个纯函数(src/providers/select.ts:20-25),只为了能被非交互地演示:
export function chooseModel(apiKey: string | undefined): ModelChoice {
if (apiKey && apiKey.trim().length > 0) {
return { kind: "anthropic", reason: "检测到 ANTHROPIC_API_KEY,使用真实 Provider。" };
}
return { kind: "fake", reason: "未检测到 ANTHROPIC_API_KEY,回退到 FakeModel(离线)。" };
}demo 分别用「没 Key」和「有 Key」两种输入调它一次,真实输出(expected-output.txt 第 86–89 行):
ANTHROPIC_API_KEY=(未设置) → 未检测到 ANTHROPIC_API_KEY,回退到 FakeModel(离线)。
ANTHROPIC_API_KEY=sk-ant-xxx → 检测到 ANTHROPIC_API_KEY,使用真实 Provider。
本 demo 固定使用 FakeModel,保证离线可复现。然后用 FakeModel 跑一次完整循环(expected-output.txt 第 91–101 行):
你> 算一下 12*8
[第 1 轮]
助手> 先算数。
[工具调用] calc {"expression":"12*8"}
[工具结果] calc → 96
[第 1 轮结束]
[第 2 轮]
助手> 12*8 = 96。
[第 2 轮结束]
上层代码与 step07 完全相同,只是 streamFn 换了一个实现。(共 4 条消息)这段输出与 step07 逐字相同——这正是重点。一个成功的抽象,它兑现的那一刻应该是「什么都没发生」。
真正把两个分支接起来的是 createModel(src/providers/select.ts:32-42),它的返回类型里 streamFn: StreamFn,两个分支返回同一个类型:
const streamFn =
choice.kind === "anthropic" && apiKey
? createAnthropicModel({ apiKey })
: createFakeModel(rules, { delayMs: 30 });把三段接起来:createAnthropicModel
前面三件事都是纯函数,最后一步才碰网络(src/providers/anthropic.ts:300-341)。核心只有一个 fetch:
response = await fetch(API_URL, {
method: "POST",
headers: {
"content-type": "application/json",
// 认证头是 x-api-key,不是常见的 Authorization: Bearer。
"x-api-key": options.apiKey,
"anthropic-version": API_VERSION,
},
body: JSON.stringify(body),
signal,
});四处细节:
- 认证头是
x-api-key,不是很多接口习惯的Authorization: Bearer。写错了会拿到 401,而不是「参数不对」这种好懂的报错。 anthropic-version是必填头,值是"2023-06-01"(anthropic.ts:32)。代码注释说明缺了它会直接 400。它是协议版本不是模型版本——服务端靠它决定用哪套字段语义回应你,所以它跟你选claude-opus-5还是别的模型没有关系。signal直接透传给fetch。step07 建立的取消链条到这里自然延续:用户按 Ctrl+C →AbortController.abort()→fetch断开连接。不需要为真模型再写一套取消逻辑。- 错误一律编码成事件,不 throw。三个失败出口都遵守这条契约,下面是前两个(第三个是包住
yield* toStreamEvents(...)的 try/catch,见anthropic.ts:335-339):
} catch (error) {
// 契约:不 throw。连不上、被取消,都变成事件。
yield aborted(signal) ?? { type: "error", error: `请求失败:${describe(error)}` };
return;
}
if (!response.ok || !response.body) {
// 出错时响应体是一段 JSON,不是 SSE,所以整段读出来当错误信息。
const detail = await response.text().catch(() => "");
yield { type: "error", error: `HTTP ${response.status}:${detail.slice(0, 300)}` };
return;
}第二段那句注释指向一个真实的坑:请求失败时响应体不是 SSE,而是一整段 JSON 错误对象。如果不判 response.ok 就直接喂给 SSE 解析器,你会得到一个「解析不出任何事件」的空流,然后在别处看到一个莫名其妙的「流被中断」——离真正的原因(比如 Key 写错了)隔着三层。
被取消时的行为和 FakeModel 保持一致(anthropic.ts:344-347):
function aborted(signal: AbortSignal | undefined): StreamEvent | undefined {
if (!signal?.aborted) return undefined;
return { type: "done", message: { role: "assistant", content: [], stopReason: "aborted" } };
}还有一处类型上的小动作值得说明(anthropic.ts:334):
const chunks = decodeChunks(response.body as unknown as AsyncIterable<Uint8Array>);Node.js 的 fetch 返回的 ReadableStream 是可以直接 for await 的,但 TypeScript 自带的 DOM 类型定义里没有这个签名,所以要断言一次。这个断言是在说「我知道我在 Node 里跑」——如果这段代码要搬到浏览器,这里就是第一个要改的地方。
图 8.6-3 一次真实调用的完整时序
从 Agent Loop 发出请求,到第一个 StreamEvent 回到 Agent Loop 的全过程。
这张图要注意两点。其一,L 这一列在 step07 和 step08 里完全相同——它发出的和收到的都只有 Context 与 StreamEvent,中间那五列是新写的,但对它不可见。其二,图上没有任何一条箭头把异常抛回 L:所有失败路径最后都变成一个 T-->>L: error,这是 StreamFn 契约的直接后果。真实源码里 S、B、P、T 分别对应 createAnthropicModel 返回的匿名生成器、toRequestBody、sse.ts 的两个函数、toStreamEvents。
API Key 安全:凭证只应活在进程环境里
真模型这一步引入了全书唯一的机密数据。三条最低要求,按「出事概率」从高到低排:
git add 的东西。Key 一旦进入 git 历史,删一次提交是不够的,必须去服务商控制台吊销并重新生成。2. 绝不打印:不写进日志、不写进错误信息、不写进会话文件。回显 Key 的调试语句最容易在事后被忘掉。
3. 只从环境变量读:让凭证只活在进程的环境里,随进程一起消失。
step08 是这么做的:
- 读:整个项目里出现 Key 的地方只有一处,
src/main.ts:72的env.ANTHROPIC_API_KEY。它顺着createModel传给createAnthropicModel,最后只出现在 HTTP 请求头里。 - 不打印:
chooseModel返回的reason是「检测到 ANTHROPIC_API_KEY,使用真实 Provider。」——它只说有没有,不说是什么(src/providers/select.ts:22)。启动横幅打印的是这句话(main.ts:112),不是 Key。 - 不落盘:step06 的会话文件写的是一行 header(版本号、会话 id、创建时间、工作目录)加一条条消息(
src/session.ts:22-43),凭证不在其中。这不是巧合,是「会话文件只存对话」这条边界带来的附赠品。
命令行上的用法也有讲究。ANTHROPIC_API_KEY=sk-... npm start 最省事,但它会进入 shell 的历史文件。更稳妥的做法是把 Key 放进一个只有自己可读的文件(chmod 600),由 shell 启动时加载,并确认该文件在版本控制的忽略清单里;或者交给系统自带的钥匙串工具保管。
Pi 是怎么做的(源码事实)。凭证的解析顺序写在 anthropicApiKeyAuth() 里:先看已存储的凭据,再看 ANTHROPIC_AUTH_TOKEN(它作为 Authorization: Bearer 头使用),最后依次看 ANTHROPIC_OAUTH_TOKEN 和 ANTHROPIC_API_KEY。
resolve: async ({ ctx, credential }) => {
if (credential?.key) {
return { auth: { apiKey: credential.key }, env: credential.env, source: "stored credential" };
}
const authToken = await ctx.env(ANTHROPIC_AUTH_TOKEN_ENV);
if (authToken) {
return {
auth: { headers: { Authorization: `Bearer ${authToken}` } },
source: ANTHROPIC_AUTH_TOKEN_ENV,
};
}
for (const envVar of [ANTHROPIC_OAUTH_TOKEN_ENV, ANTHROPIC_API_KEY_ENV]) {
const apiKey = await ctx.env(envVar);
if (apiKey) return { auth: { apiKey }, source: envVar };
}
return undefined;
},anthropicApiKeyAuthresolve 回调(第 16–34 行)。解析顺序是:已存储凭据优先,然后依次尝试三个环境变量。注意每个分支都带 source,用来在界面上说明「这个 Key 是从哪来的」——而不是把 Key 本身显示出来。三个环境变量名定义在 packages/ai/src/env-api-keys.ts:29-31。 「已存储凭据」存在 auth.json 里,路径是 agent 目录下的 auth.json(默认 ~/.pi/agent/auth.json):
AuthStorageAuthStorage.create() 默认落在 getAgentDir() 下的 auth.json。同文件 FileAuthStorageBackend 建目录时用 mode: 0o700(第 38 行)、写文件时用 mode: 0o600 并额外 chmodSync(…, 0o600)(第 21、45 行)——即「只有文件属主可读写」。 从源码结构看,Pi 把「凭证从哪来」(resolve)与「凭证存在哪」(AuthStorage)拆成了两件事:前者住在 pi-ai 里、对每家服务商各写一份,后者住在 coding-agent 里、只有一份并负责文件权限。我们的教学版只实现了「从环境变量来」这一条路径,但边界画在同一个地方——凭证不进业务代码,只在最后一刻进请求头。
回到 Pi 源码:同一个骨架,多十倍的细节
我们这个 351 行的 anthropic.ts(加上 92 行的 sse.ts),和 Pi 的 packages/ai/src/api/anthropic-messages.ts 骨架完全一样:手写 SSE 解析 → 事件过滤 → 映射成统一事件。差别在细节的密度上。逐项对照最能说明问题。
第一处:StreamFn 的契约,Pi 写得更明确。
// Generic StreamFunction with typed options.
//
// Contract:
// - Must return an AssistantMessageEventStream.
// - Once invoked, request/model/runtime failures should be encoded in the
// returned stream, not thrown.
// - Error termination must produce an AssistantMessage with stopReason
// "error" or "aborted" and errorMessage, emitted via the stream protocol.
export type StreamFunction<TApi extends Api = Api, TOptions extends StreamOptions = StreamOptions> = (
model: Model<TApi>,
context: Context,
options?: TOptions,
) => AssistantMessageEventStream;StreamFunctionStreamFunction 与契约注释。第 316–317 行那句「失败应当编码进返回的流里,而不是抛出」就是我们从 step03 起遵守的那条约定的原文。 两处签名差异值得注意:Pi 多一个 model 参数(因为一个 StreamFunction 要服务同一协议下的多个模型),返回的是 AssistantMessageEventStream 这个具体类而不是裸的异步迭代器(因为它还要支持「只要最终结果」的 result() 调用)。我们的教学版把这两件事都省了。
第二处:SSE 解析,Pi 也是手写的。
iterateSseMessagesTextDecoder 加 buffer 加循环消费,区别是它按行切(consumeLine)并用一个 SseDecoderState 累积多行 data:,而我们按空行切成块再解析。第 398–400 行每读一块就检查一次 signal?.aborted——这是我们没做的一层保险。 iterateAnthropicEventsevent: error 直接抛、用 parseJsonWithRepair(而不是裸 JSON.parse)解析负载,并在第 482–484 行校验「流必须以 message_stop 结束」。我们的 toStreamEvents 用 started 标志做了同一件事的简化版(anthropic.ts:292-295)。 parseJsonWithRepair 是我们没有的东西:真实世界里会遇到被代理截断、被中间层改写的畸形 JSON,Pi 选择尽力修复而不是直接失败。
第三处:分片工具参数,Pi 边收边解析。
} else if (event.delta.type === "input_json_delta") {
const index = blocks.findIndex((b) => b.index === event.index);
const block = blocks[index];
if (block && block.type === "toolCall") {
block.partialJson += event.delta.partial_json;
block.arguments = parseStreamingJson(block.partialJson);
stream.push({
type: "toolcall_delta",
contentIndex: index,
delta: event.delta.partial_json,
partial: output,
});
}
}input_json_deltaparseStreamingJson,把「半截 JSON」尽力解析成一个对象,并发出 toolcall_delta 事件。块结束时(第 694–705 行)再解析一次定稿,并删掉 partialJson 这个脚手架字段。 这是我们和 Pi 最大的一处行为差异,值得说清楚代价与收益。我们只在块结束时解析一次:实现简单,toolCall.arguments 永远是完整可信的对象,代价是界面上没法在参数打字的过程中就显示「它正准备读哪个文件」。Pi 每片都解析一次:界面可以实时显示参数轮廓(Pi 的 TUI 就是这么做的),代价是 arguments 在中途是一个可能残缺的对象,所有消费方都得知道这件事——toolcall_delta 事件带的 partial 字段就是在提醒这一点。一种看法是:教学版该选前者,产品该选后者。
第四处:停止原因,Pi 分得更细。
mapStopReasonend_turn 与 pause_turn 与 stop_sequence → stop,max_tokens → length,tool_use → toolUse,refusal 与 sensitive → error 并附带 errorMessage,未知取值直接抛异常。 对照我们的 mapStopReason(src/providers/anthropic.ts:160-171),差别有三:Pi 有独立的 length 表示「被 token 上限截断」,我们把它并进了 stop;Pi 对拒答带上可读的 errorMessage,我们只给一个 error;Pi 对未知的停止原因抛异常,宁可失败也不猜——这是一个刻意的选择:静默地把新出现的取值当成 stop,会让「模型其实没说完」这件事永远不被发现。
第五处:我们没有而 Pi 有的东西。 重试与限流退避(retryProviderRequest)、超时、缓存控制、用量与费用统计、thinking 内容块、图片输入、按需加载厂商 SDK。还有一个数量级上的差异:Pi 支持多种 API 协议、几十个内置 Provider,而我们只处理 text 与 tool_use 两种内容块。
Response 后立刻 .asResponse() 取出原始响应体,SSE 解析和事件映射全部自己写(packages/ai/src/api/anthropic-messages.ts:559-573)。这是一个务实的分工:认证、重试、连接管理这类和协议演进强相关的事交给 SDK,而「一次生成过程如何被拆成事件」这件事关乎 Pi 自己的统一事件模型,必须自己掌握。本步骤的约束是「除 tsx 与 typescript 外零运行依赖」,所以连发请求也自己写了一遍——教学上正好把 SSE 这层看清楚,真实项目里没必要重复这份工作。 实践任务
labs/mini-agent-harness/step08-real-provider目标:不用 API Key,验证适配器的三件事全部工作;并能说出「换实现时上层改了几行」。
步骤 1 · 跑离线 demo:
cd labs/mini-agent-harness/step08-real-provider
npm install
npm run demo预期现象:输出与 expected-output.txt 逐字相同,共三段场景。全程不联网,不需要任何环境变量。
步骤 2 · 确认「上层没改」:把 step07 和 step08 的源码目录做一次逐文件比较:
diff -rq ../step07-abort-extensions/src src \
| grep -v providers预期现象:只有两个文件不同——main.ts(接线点)与 demo.ts(本步骤新写的离线演示);agent-loop.ts、tools/、session.ts、extensions.ts、types.ts、render.ts、fake-model.ts、rules.ts 全部一致。这条 diff 就是本章那句「一行都不用改」的证据。
再看 main.ts 到底改了什么:diff ../step07-abort-extensions/src/main.ts src/main.ts。你会看到实质改动只有三处:换掉一个 import、把造 streamFn 的那一行换成 createModel(env.ANTHROPIC_API_KEY)、启动横幅多打一行模型来源。其余差异是文件头注释与系统提示词的措辞,与「换模型」这件事无关。
步骤 3 · 不设 Key 跑交互模式:
npm start预期现象:启动横幅第二行显示「模型:未检测到 ANTHROPIC_API_KEY,回退到 FakeModel(离线)。」,之后的对话行为与 step07 完全一致(可以试 算一下 12*8、说个长的 再按 Ctrl+C、/exit)。
步骤 4 · 证明解析器与分块无关:把 src/providers/recorded-stream.ts:56 的 size = 37 分别改成 1(一个字符一块,最极端的碎片)和 5000(整段挤在一块),各跑一次 npm run demo。
预期现象(写作者实测):第一行的「被切成 N 个网络块」会变成 1217 和 1,但下面那五行事件——start / 两条 text_delta / toolcall / done——一字不变。这就是「协议解析不依赖网络分块」的可观察定义。改回 37 再往下走。
步骤 5 · 改坏解析器,看它怎么崩(这一步最有教学价值):把 src/providers/sse.ts:49-71 里 readSseMessages 的整个函数体换成「假设一块就是一个事件」的写法——删掉 buffer,直接对每个 chunk 调 parseBlock:
for await (const chunk of chunks) {
const message = parseBlock(chunk.replace(/\r\n/g, "\n"));
if (message) yield message;
}预期现象(写作者实测,size 仍为 37):场景 2 的事件列表整个消失,只剩一行——
error 无法解析的事件数据:{"type":"m第 1 块在 JSON 中间断掉,翻译器立刻放弃。接着把 size 改成 5000 让整段挤进一块:错误不会消失,只会换一副面孔——变成一条把十一段 JSON 拼在一起的 error,因为丢掉缓冲的同时也丢掉了事件边界。两种分块都错,只是错得不一样,这正是这类 bug 难查的原因:它不随响应变小而消失,只随响应变化而变形。看完把这两处照原样改回去。
步骤 6(可选,需要 API Key,会产生费用):
ANTHROPIC_API_KEY=你的Key npm start预期现象:横幅显示「检测到 ANTHROPIC_API_KEY,使用真实 Provider。」,之后是真实模型的回复。它不一定会像 FakeModel 那样乖乖调用 calc——真模型只是参考工具的 name / description / parameters 自己判断。如果它总不调用工具,先去改工具描述和系统提示词,那才是你能控制的部分。
如何判断成功:① 步骤 1 的输出与 expected-output.txt 一致;② 步骤 2 的 diff 只剩 main.ts 与 demo.ts 两个文件;③ 步骤 4 里三种分块大小给出完全相同的事件列表;④ 你能不看书说出「为什么工具参数必须等 content_block_stop 才能 JSON.parse」;⑤ 你能指出如果去掉 TextDecoder 的 { stream: true },会在什么条件下出现乱码。
常见错误:
- 用
node src/demo.ts而不是npm run demo:本项目靠tsx直接跑 TypeScript,直接用 node 会因为类型语法报错。 - 把 Key 写进
package.json的 script 里图省事:那是入库的最快路径,正是本章说的第一条底线。 - 步骤 6 收到
HTTP 401:先确认 Key 没有多余的空格或换行(从网页复制常带上),再确认用的是x-api-key而不是Authorization。 - 步骤 6 收到
HTTP 404且信息里提到 model:src/providers/anthropic.ts:301的默认模型名不一定在你的账号下可用,换成你有权限的模型 id。
对应源码位置:src/providers/sse.ts:31-92(解码与分帧)、src/providers/anthropic.ts:103-137(请求体)、189-296(事件翻译)、300-341(fetch)、src/providers/select.ts:20-42(选模型)、src/main.ts:72(唯一的接线点)。
本章小结
- 一个 Provider 适配器只做三件事:请求体翻译(
toRequestBody)、流式解析(decodeChunks+readSseMessages)、事件翻译(toStreamEvents)。换一家服务商,只有这三件事要重写。 - 请求体翻译不是字段改名:三种 role 降成两种要合并,系统提示词从数组成员变成顶层字段;相邻同角色消息不合并,会在数据里教会模型「不要并行调用工具」。
- SSE 解析的唯一难点是网络分块与协议分隔毫无关系。对策是三件套:
TextDecoder的stream模式处理半个汉字,buffer加循环找分隔符处理半个事件,流末尾再收一次尾。 - 工具参数以
input_json_delta分片到达,单片不是合法 JSON,必须等content_block_stop拼完整再解析。上层之所以能拿到「arguments是对象」的干净结果,是因为这份脏活在适配器里干完了。 - 错误全部编码成
error事件而不是抛异常。这条从 step03 起就定下的契约,到真实网络这里才显出全部价值:失败可能发生在流的任何位置,统一成事件后调用方只有一套处理路径。 - API Key 的三条底线:绝不入库、绝不打印、只从环境变量读。step08 里 Key 只出现在
main.ts:72和请求头两处;Pi 把已存储的凭证放在auth.json,目录0700、文件0600。 - 判断一个抽象值不值,标准是「换实现时上层要改多少行」。step08 的答案是:
src/下只有main.ts需要动,动的是一个import和一行赋值;其余文件与 step07 逐字相同。 - 关键术语:Provider(模型服务提供方)、适配器(Adapter)、SSE(Server-Sent Events,服务器发送事件)、分块(chunk)与分帧(framing)、分片 JSON(partial JSON)、
StreamFn/StreamFunction契约、停止原因(Stop Reason)、API Key。 - 关键源码索引:
- 实验代码:
labs/mini-agent-harness/step08-real-provider/src/providers/sse.ts:31-92、anthropic.ts:65-137(请求体)、148-296(事件翻译)、300-347(fetch 与取消)、select.ts:20-42 - Pi:
packages/ai/src/types.ts:312-324(StreamFunction契约)、packages/ai/src/api/anthropic-messages.ts:387-444(SSE 解析)、446-485(事件过滤与 JSON 修复)、654-666(分片工具参数)、1325-1351(停止原因映射) - 凭证:
packages/ai/src/providers/anthropic.ts:16-34、packages/ai/src/env-api-keys.ts:29-31、packages/coding-agent/src/core/auth-storage.ts:180-182
- 实验代码:
- 自测问题:
- 一段 SSE 响应共 1217 个字符,被切成 33 块到达。如果解析器假设「一块 = 一个事件」,会发生什么?把块大小改成 5000 让整段挤进一块之后,错误为什么没有消失,而是换了个样子?反过来,正确的解析器为什么在块大小为 1 和 5000 时输出一模一样?
- 我们的
ToolResultMessage是第三种 role,而请求体里它变成了role: "user"。如果有三个并行工具结果,翻译成三条独立的 user 消息会有什么后果? toStreamEvents收到content_block_delta且delta.type === "input_json_delta"时,为什么一个StreamEvent都不发?Pi 在同样的位置发了toolcall_delta,它多付出了什么代价?- 真模型返回 HTTP 500 时,
runAgentLoop会收到什么?如果适配器改成直接throw,render.ts和main.ts分别要多写什么? - 假设你要给这个 harness 再接一家用「换行分隔的 JSON」而不是 SSE 的模型服务,
src/下哪些文件需要改,哪些一定不用改?
全书到这里
八个部分走完了。回头看,这条路是这样铺的:
第一、二部分把 TypeScript 和异步迭代器讲成工具而不是语法清单;第三部分把 Agent 的几个核心概念——消息、上下文、流式、工具调用、循环、会话——各自拆开;第四到第七部分进 Pi 的源码,从仓库结构一路读到 Agent Loop、pi-ai、工具系统、Session 存储、扩展与技能;第八部分反过来,自己动手把这些概念重新搭一遍,八步、每步一个新概念、每步都能跑。
而这最后一章证明的事情,其实是前面所有设计决策的验收:当 StreamFn 的两个实现——一个按剧本吐字、一个跨越互联网——可以在一行代码里互换,而 Agent Loop、工具执行、会话持久化、取消、扩展全部无动于衷时,前面那些看起来啰嗦的边界才算真的成立了。 这也是读 Pi 源码时最值得带走的一样东西:Pi 的复杂度不小,但它的每一层边界都在回答同一个问题——什么变了的时候,什么可以不变。
从这里往哪走,取决于你想要什么:
- 想让这个 mini harness 更像 Pi:
src/providers/anthropic.ts里的toStreamEvents加上对thinking内容块的处理;给createAnthropicModel加上重试与超时;在session.ts里补上parentId真正的分支能力。每一项都能在 Pi 里找到对照实现。 - 想给 Pi 加东西而不是重写一个:去 7.1 Extension 系统、7.3 自定义工具与斜杠命令、7.4 自定义 Provider——本章手写的这套适配器,正是理解 7.4 的最短路径。
- 想继续读源码:6.1 pi-ai:统一的模型接口 是本章的完整版,6.2 Provider 与模型注册 讲几十个 provider 是怎么组装起来的,5.4 流式事件如何传播到界面 接着讲事件离开 pi-ai 之后的旅程。
- 想查东西:术语表、源码索引、附录 · 延伸阅读。
本书尚未展开的内容(留给你自己去读源码的部分):OAuth 登录流程、多模态输入(图片)、提示词缓存与用量计费、packages/evals 的评测框架、TUI 的渲染细节,以及 Pi 的远程过程调用(RPC,Remote Procedure Call)服务端模式。这些都不是本书骨架的一部分,但每一个都在源码里有完整实现——现在你已经知道该从哪个文件开始读了。