Skip to content

6.1 pi-ai:统一的模型接口 ​

本页分析版本earendil-works/pi@16787ad2026-09-21

本章解决什么问题:Pi 支持几十家模型服务商,但 Agent Loop 里只有一行 models.stream(...)。这中间的「统一层」到底统一了什么、怎么实现的。 前置知识:2.6 异步迭代器与 for await、3.3 流式输出、3.4 工具调用、4.3 monorepo 与 package 地图。 学习目标:① 说出 pi-ai 的四层结构与各层职责;② 看懂统一消息类型与 12 种流式事件的协议约定;③ 讲清 complete() 与 stream() 的关系以及 EventStream 的缓冲机制;④ 跟着一条完整调用链读完 Anthropic 协议实现的 SSE 转换;⑤ 不用 API Key 就能用 faux provider 跑出一次完整的假对话。

建立直觉:几十家 provider,一个形状 ​

Anthropic 的流式接口返回 content_block_delta,OpenAI 的 Responses 接口返回 response.output_text.delta,Google 返回的又是另一套字段名。它们描述的是同一件事——「模型刚刚多吐了几个字」——但字段名、嵌套层次、停止原因的取值全都不同。

没有统一层会怎样?Agent Loop 里就必须写 if (provider === "anthropic") { ... } else if (provider === "openai") { ... },而且每加一家服务商,Agent Loop、工具执行、界面渲染、Session 存储都要跟着改一遍。上下文里存下来的历史消息也会带上厂商烙印,换个模型继续对话就变成了数据迁移工程。

pi-ai 的做法是在中间插一层「形状统一」:进去的是统一的 Context(在交给 provider 之前会被规范化成 TranscriptContext),出来的是统一的 AssistantMessageEvent 事件流,中间由谁去调哪家 HTTP 接口,调用方完全不需要知道。这一层可以拆成四块来理解:

📘 概念pi-ai 的四层(four layers of pi-ai)
数据层(types.ts):消息、内容块、上下文长什么样。
协议层(types.ts 的事件定义 + utils/event-stream.ts):一次生成过程如何被拆成事件、如何被消费。
路由层(models.ts):`Models` 集合持有若干 `Provider`,按 `model.provider` 把请求派发给对应的 provider,顺带解析凭据。
实现层(api/*.ts):每个文件是一种线上协议(wire protocol)的适配器,负责把厂商响应翻译成统一事件。

包本身是独立发布的:npm 包名 @earendil-works/pi-ai,锁定版本 0.87.0,自述为 "Unified LLM API with automatic model discovery and provider configuration"(官方说明,见 packages/ai/package.json:2-4)。在仓库内部它只依赖一个很小的 @earendil-works/pi-telemetry(遥测契约与类型工具,packages/ai/package.json:69),不依赖 agent、coding-agent 等上层包——这是 4.3 那张依赖图的最底层。

根入口刻意保持「无副作用」:index.ts 开头的注释写明,这里不带生成的模型目录、不带 provider 工厂、不带 OAuth 实现;provider 工厂要从子路径 @earendil-works/pi-ai/providers/* 引入,协议实现在 @earendil-works/pi-ai/api/*,工具函数在 @earendil-works/pi-ai/utils/*(源码事实,packages/ai/src/index.ts:4-8 与 package.json 的 exports)。这样只用到一两家 provider 的程序,打包时不会把几十家的实现全拖进来。

门面:Models 集合 ​

pi-ai 的对外入口不是一堆全局函数,而是一个可以创建多个实例的集合对象。createModels() 造一个空集合,builtinModels() 造一个已注册全部内置 provider 的集合(源码事实):

ts
// packages/ai/src/providers/all.ts:137-143
export function builtinModels(options?: CreateModelsOptions): MutableModels {
	const models = createModels(options);
	for (const provider of builtinProviders()) {
		models.setProvider(provider);
	}
	return models;
}
earendil-works/pi@16787ad第 137–143 行在 GitHub 查看 ↗
「注册」在这里是显式的函数调用,不是 import 副作用。同文件 90-134 行的 builtinProviders 返回 41 个 provider 工厂的调用结果(逐个数 92-132 行的数组元素;这 41 个 id 恰好与 types.ts 里 KnownProvider 联合的 41 个成员一一对应)。

集合上的方法分三类:查目录(getProviders / getModels / getModel / refresh,以及只列出已配好凭据的 getAvailable)、管凭据(getAuth / login / logout / checkAuth)、发请求(stream / complete / streamSimple / completeSimple,以及处理「延迟响应」的 streamDeferred / fetchDeferred / cancelDeferred,下文协议层会讲)。

earendil-works/pi@16787ad第 163–235 行在 GitHub 查看 ↗
Models interface:provider 集合 + auth 解析 + 请求门面。注意 stream 返回 AssistantMessageEventStream,complete 返回 Promise<AssistantMessage>。

四个发请求的方法其实只有两个维度的差别。第一个维度是「流还是结果」——这一点被实现得非常直白:

ts
// packages/ai/src/models.ts:695-701
	async complete<TApi extends Api>(
		model: Model<TApi>,
		context: Context,
		options?: ModelsApiStreamOptions<TApi>,
	): Promise<AssistantMessage> {
		return this.stream(model, context, options).result();
	}
earendil-works/pi@16787ad第 695–701 行在 GitHub 查看 ↗
complete 不是另一条实现路径,而是在 stream 出来的流上取最终结果。completeSimple 与 streamSimple 的关系同理,见同文件 712-718 行。

这句话值得停一下:pi-ai 里没有「非流式」的代码路径。所有 provider 实现都只写流式一种,非流式只是「订阅这条流的最终结果」。一种看法是:这样每个 provider 少写一半实现,两条路径也不会各自长出不同的 bug;代价是即使只想要一句完整回答,也要付出建流、入队、事件分发的开销。

第二个维度是 stream 与 streamSimple 的区别:前者接受该 provider 协议专属的选项(例如 Anthropic 的 thinkingEnabled / thinkingBudgetTokens),后者接受一个统一的推理档位 ThinkingLevel = "minimal" | "low" | "medium" | "high" | "xhigh" | "max"(types.ts:84),由各实现自己翻译成厂商参数。Anthropic 的翻译逻辑在 api/anthropic-messages.ts:858-904:没请求推理就直接 thinkingEnabled: false(869-874 行),支持自适应推理的模型映射到 effort(878-885 行),老模型则换算成 token 预算(887-903 行)。翻译完照样回落到同文件的 stream。

数据层:一条消息长什么样 ​

统一层的第一块地基是消息类型。Pi 有四种消息,联合成 Message(源码事实):

ts
// packages/ai/src/types.ts:491-553(节选,UserMessage 与 ToolResultMessage 压成了较少的行)
export interface SystemMessage {
	role: "system";
	content: string | TextContent[];
	sections?: Record<string, string | null>;
	toolsAdded?: Tool[];
	toolsRemoved?: ToolReference[];
	timestamp: number;
}
export interface UserMessage { role: "user"; content: string | (TextContent | ImageContent)[]; timestamp: number }
export interface AssistantMessage {
	role: "assistant";
	content: (TextContent | ThinkingContent | ToolCall)[];
	api: Api;
	provider: ProviderId;
	model: string;
	// …(省略:responseModel / responseId / diagnostics / deferred / errorMessage / endTurn 等)
	usage: Usage;
	stopReason: StopReason;
	timestamp: number;
}
export type ToolResultMessage<TDetails = JsonValue> = IsJsonCompatible<TDetails> extends true
	? { role: "toolResult"; toolCallId: string; toolName: string; content: (TextContent | ImageContent)[];
	    details?: JsonRepresentation<TDetails>; isError: boolean; timestamp: number /* …(省略:usage) */ }
	: never;
export type Message = SystemMessage | UserMessage | AssistantMessage | ToolResultMessage;
earendil-works/pi@16787ad第 364–553 行在 GitHub 查看 ↗
内容块 TextContent / ThinkingContent / ImageContent / ToolCall 在 364-394 行,Usage 与 StopReason 在 396-419 行,JsonValue 与 JSON 兼容性检查在 421-467 行,四种消息与 Message 联合在 481-553 行。

这是一个典型的可辨识联合(2.3 讲过):role 是判别字段。几个值得注意的设计(分析解释):

  • 助手消息的内容块只允许三种:文本、思考(thinking)、工具调用(types.ts:517)。图片只能出现在用户消息和工具结果里——在这套对话类型里模型不输出图片(图像生成是另一套 AssistantImages 类型,types.ts:564-574),这条约束直接写进了类型而不是靠约定。
  • 工具调用是内容块,不是独立消息。ToolCall.arguments 是已经解析好的对象(types.ts:386-394),调用方拿到的不是 JSON 字符串;而且它的类型是 JsonObject,只允许 JSON 能表示的值(types.ts:421-422)。
  • 工具结果是独立消息(role: "toolResult"),用 toolCallId 回指。这解释了 Agent Loop 为什么能把一轮工具调用拆成「助手消息 + 若干工具结果消息」追加进历史。ToolResultMessage 现在是一个条件类型(types.ts:539-551):details 的类型参数必须是 JSON 兼容的,否则整个类型变成 never,编译期就报错。从源码结构看,这是为了保证任何一条消息都能原样写进会话文件或跨进程传输,不会混进函数、undefined 之类序列化后会走样的值。
  • 每条消息都带 timestamp,助手消息还带 api / provider / model / usage / stopReason——历史记录里保存的是「这句话是谁在什么配置下生成的」,而不只是文本。这是 Session 能跨模型接力的前提。助手消息还可能带 endTurn(provider 是否明确表示本轮结束),注释写明它只为调试保留,目前不影响 agent 的控制流(types.ts:531-535)。

系统提示词也是一条消息 ​

第四种消息 SystemMessage 值得单独说。请求的输入侧仍然可以写成熟悉的 Context = { systemPrompt?, messages, tools? }(types.ts:617-621),工具用 TypeBox schema 描述参数(types.ts:600-605)。但这只是给调用方的简写:Models 的公开入口会先调一次 normalizeContext(),把 systemPrompt 和 tools 折叠成 transcript 开头的一条 system 消息(utils/transcript.ts:30-34),得到只有 messages 字段的 TranscriptContext(types.ts:631-634)。TranscriptContext 在类型上带一个「品牌」标记(一个只存在于类型系统里的 unique symbol,运行时并没有这个字段),所以只有 normalizeContext() 能产出它——源码注释的说法是「a raw Context cannot reach provider code by accident」。

好处是对话中途可以改系统提示词和工具:往 transcript 里再追加一条 system 消息,content 追加指令、sections 按名字替换或删除提示词片段、toolsAdded / toolsRemoved 增减工具。按顺序「回放」全部 system 消息,就得到当前的提示词与工具集,对应的函数是 getCurrentSystemPrompt() 与 getCurrentTools()(utils/transcript.ts:99-102、58-66)。provider 实现都从 context.messages 里这样读,context.systemPrompt / context.tools 在那一层根本不存在(官方说明,packages/ai/README.md:1412)。支持中途 system 消息的模型原样收到每一条;不支持的,由 resolveTranscript() 折叠成开头一条(utils/transcript.ts:108-120)。Anthropic 实现一进门就做这两件事(api/anthropic-messages.ts:517-518)。

协议层:12 种事件 ​

生成过程被拆成 12 种事件,全部定义在一个联合类型里:

packages/ai/src/types.ts · AssistantMessageEvent
earendil-works/pi@16787ad第 652–668 行在 GitHub 查看 ↗
统一流事件协议。637-651 行的注释写明约定:先 start,再增量,最后以 done 或 error 二选一终止;请求还没开始生成就失败时,可以只有一个 error。
事件何时发出关键字段
start流开始,此时 partial.content 还是空数组partial
text_start一个文本块开始contentIndex
text_delta文本块多了一段delta、contentIndex
text_end文本块结束content(完整文本)
thinking_start一个思考块开始contentIndex
thinking_delta思考块多了一段delta
thinking_end思考块结束content
toolcall_start一次工具调用开始contentIndex
toolcall_delta工具参数的 JSON 又来了一截delta(JSON 片段)
toolcall_end工具调用完整toolCall(含 id/name/arguments)
done正常终止reason:stop / length / toolUse / deferred;message
error异常终止reason:aborted / error;error(同样是一条 AssistantMessage)

三条约定必须记住:

  1. 每个事件都带 partial,即「截至此刻已经拼好的那条 AssistantMessage」。界面层不需要自己攒字符串,直接渲染 partial 即可(5.4 流式事件如何传播到界面会看到 Pi 的 TUI 就是这么做的)。但要注意它是一个共享的活对象,不是事件发生那一刻的快照:provider 会持续原地修改同一条消息,哪怕旧事件还在队列里没被取走(types.ts:645-650 的注释,README.md:663 也这么说)。所以处理事件时读它可以,把它存起来当历史记录就不行。
  2. 块之间可能交错。官方文档明确说明:不同内容块的事件不保证连续,可能出现 text_start → text_delta → toolcall_start → text_delta 这种顺序,消费者必须用 contentIndex 把 delta 关联回它所属的块(packages/ai/README.md:682)。
  3. 终止事件二选一,且 error 也携带一条完整消息。也就是说「失败」不是抛异常,而是一条 stopReason 为 error 或 aborted 的助手消息。types.ts:335-345 把这一点写成了 StreamFunction 的契约:流一旦返回,请求、模型、运行时的失败都必须编码进流里。契约里只留了一个例外——直接调用某个协议实现的 streamSimple() 而凭据缺失时,会同步抛错(例如 Anthropic 的 assertRequestAuth,api/anthropic-messages.ts:863);走 Models 门面则不会,因为那里整个请求准备都包在 lazyStream 里,异常会变成流里的 error 事件。这是整个 pi-ai 错误模型的基石——通过 Models 调用时,调用方永远只需要处理流,不需要同时处理 try/catch 和流。

表里 done 的第四种 reason——deferred——对应一个新能力:延迟响应。调用 streamSimple 时传 deferred: true(types.ts:329-330),支持它的 provider 会先返回一条 stopReason: "deferred" 的消息,里面带一个可持久化的句柄 DeferredHandle(types.ts:469-478),真正的结果稍后再用 models.fetchDeferred() / streamDeferred() 取回,或用 cancelDeferred() 取消(models.ts:720-755)。于是 StopReason 一共 7 种取值:pending(只出现在生成中的 partial 里)、stop、length、toolUse、error、aborted、deferred(types.ts:419)。官方 README 的「Stop Reasons」一节(README.md:921-930)和事件表(README.md:678)目前都还没列出 deferred,以源码为准。

EventStream:推与拉之间的那个缓冲 ​

事件是 provider 实现「推」出来的,调用方是用 for await 一条条「拉」的。两边速度不一样:网络一次可能到来 5 个事件,而消费者还在渲染第一个。中间必须有个缓冲。

最小示例:30 行手写一个 ​

先脱离 Pi 看这个模式的本质。把下面的内容存成 mini-stream.mjs,用 node mini-stream.mjs 运行:

js
class MiniStream {
	queue = [];
	waiting = [];
	done = false;
	constructor() { this.result = new Promise((r) => (this.resolveResult = r)); }
	push(event) {
		if (this.done) return;
		if (event.type === "done") { this.done = true; this.resolveResult(event.text); }
		const waiter = this.waiting.shift();
		if (waiter) waiter({ value: event, done: false }); // 有人在等 → 直接交付
		else this.queue.push(event);                        // 没人等 → 先排队
	}
	async *[Symbol.asyncIterator]() {
		while (true) {
			if (this.queue.length > 0) yield this.queue.shift();
			else if (this.done) return;
			else {
				const r = await new Promise((resolve) => this.waiting.push(resolve));
				if (r.done) return;
				yield r.value;
			}
		}
	}
}

const s = new MiniStream();
s.push({ type: "delta", text: "你" });   // 生产者一口气推完,不等消费者
s.push({ type: "delta", text: "好" });
s.push({ type: "done", text: "你好" });
for await (const e of s) console.log(e.type, e.text); // 消费者迟到,但一条不丢
console.log("result =", await s.result);

实际运行输出(真实运行,Node 26):

text
delta 你
delta 好
done 你好
result = 你好

三个要点:队列接住消费者来不及取的事件;waiting 数组存放消费者的「等待凭证」,生产者一到就直接兑现;resolveResult 让「只要最终结果」的调用方不必遍历整条流。

回到 Pi 源码 ​

Pi 的实现就是这个结构,只是把「什么算终止」和「终止时取什么」做成了构造参数:

earendil-works/pi@16787ad第 26–89 行在 GitHub 查看 ↗
通用事件流:push 入队或直接唤醒消费者见 43-58 行;Symbol.asyncIterator 支持 for await 见 72-84 行;result 返回终止时兑现的 Promise 见 86-88 行。

和最小示例相比只有一处工程化的差别:队列不是普通数组加 shift()(每次取头都要把后面的元素整体前移),而是同文件 3-23 行的 FifoQueue——用两个栈拼成的队列,入队压进 incoming,出队时 outgoing 空了才把 incoming 整个倒过去,均摊下来每次都是常数时间。一次长回答可能有成千上万个事件,这个差别就有意义了。

AssistantMessageEventStream(同文件 91-105 行)只是给它填上具体规则:done 或 error 算终止,终止时取出 event.message 或 event.error 作为最终的 AssistantMessage。前面那句 complete() = stream().result() 到这里就闭环了——result() 返回的正是这个 finalResultPromise。

如果要把流式进度持久化(比如写进会话存储,崩溃后还原出生成到一半的消息),直接存事件并不合适:每个事件都带着越来越大的 partial。pi-ai 为此另提供了 AssistantMessageFrameEncoder 与 reduceAssistantMessageFrames(),把事件流压成紧凑的「帧」(frame)序列、再按需还原(utils/assistant-message-frame.ts,官方说明见 README.md 的「Compact Assistant Message Frames」一节)。6.3 讲的 AgentHarness 就用它保存生成中的助手消息;本章不展开。

⚠️ 常见误解以为 error 事件之后还要 catch
error 事件本身就是终止事件,result() 会正常 resolve 出一条 stopReason: "error" 的消息,而不是 reject。所以 await models.complete(...) 不会因为模型报错而抛异常——你必须检查返回消息的 stopReason。这一点在 6.3 pi-agent-core 的 Agent Loop 里非常关键。

图解:请求怎么走完这四层 ​

图 6.1-1 pi-ai 的四层与一次请求的去程回程
从上往下的实线是去程,带「AssistantMessageEvent 事件流」标签的那条是回程:事件直接回到调用方,不被中间层再包一层。虚线表示路由层与实现层都以 types.ts 为共同契约。

读图要关注两处。第一,auth 在路由层解析,不在实现层:applyAuth(models.ts:648-677)解析出 apiKey、合并 headers、必要时替换 baseUrl,再把结果作为普通选项传给 provider;调用方若传了 transformHeaders,它在所有 headers 合并完之后最后运行一次(667 行)。协议实现只管「有 key 就发请求」,不关心 key 从环境变量、凭据文件还是 OAuth 刷新来。第二,实现层是懒加载的:anthropicMessagesApi() 全文只有一行 lazyApi(() => import("./anthropic-messages.ts"))(api/anthropic-messages.lazy.ts:4),厂商 SDK 直到第一次真正发起请求才被加载。lazyStream(api/lazy.ts:46-61)让 stream() 仍然同步返回一条流,异步的加载与鉴权在背后进行,失败时转成流内的 error 事件——又一次践行「错误进流不抛出」的契约。

完整调用链(源码事实,可用 grep 逐环复核):

models.stream()(models.ts:679-692)→ normalizeContext()(utils/transcript.ts:30-34)→ lazyStream()(api/lazy.ts:46-61)→ applyAuth()(models.ts:648-677)→ provider.stream(models.ts:851)→ dispatch()(models.ts:803-814)→ lazyApi().stream(api/lazy.ts:75-76)→ stream(api/anthropic-messages.ts:511)→ resolveTranscript()(同文件 517)→ iterateAnthropicEvents()(同文件 470-509)→ iterateSseMessages()(同文件 411-468)→ 事件循环 push 统一事件(同文件 601-789)→ done(815 行)或 error(825 行)→ EventStream.result()(utils/event-stream.ts:86-88)。

解剖一个实现:anthropic-messages.ts ​

api/ 下的每个模块恰好导出 stream 与 streamSimple 两个函数,因此模块本身就满足 ProviderStreams 接口(types.ts:270-292 的注释明确写了这一点;支持延迟响应的模块还可以额外导出 fetchDeferred / cancelDeferred)。挑 Anthropic 这个看,因为它的 SSE 处理是手写的,读起来最完整。

第一步:把字节流切成 SSE 事件 ​

服务端发回的是 SSE(Server-Sent Events,服务器推送事件)文本流,形如 event: content_block_delta 换行 data: {...} 换行空行。pi-ai 没有用 SDK 自带的流封装,而是自己解码:

ts
// packages/ai/src/api/anthropic-messages.ts:411-468(节选,签名已折行压缩)
async function* iterateSseMessages(body: ReadableStream<Uint8Array>, signal?: AbortSignal) {
	const reader = body.getReader();
	const decoder = new TextDecoder();
	const state: SseDecoderState = { event: null, data: [], raw: [] };
	let buffer = "";
	try {
		while (true) {
			if (signal?.aborted) throw new Error("Request was aborted");
			const { value, done } = await reader.read();
			if (done) break;
			buffer += decoder.decode(value, { stream: true });
			let consumed = consumeLine(buffer);
			while (consumed) {
				buffer = consumed.rest;
				const event = decodeSseLine(consumed.line, state);
				if (event) yield event; // 只有空行触发 flush 时才有事件产出
				consumed = consumeLine(buffer);
			}
		}
		// …(省略:443-463 行,流读完后 flush 残留 buffer 与最后一个未闭合事件)
	} finally {
		reader.releaseLock();
	}
}
earendil-works/pi@16787ad第 411–468 行在 GitHub 查看 ↗
手写 SSE 解码:按行切分,356-380 行的 decodeSseLine 负责识别 event: 与 data: 字段,空行触发一次事件产出。

外面还套了一层过滤:iterateAnthropicEvents()(470-509 行)只放行 6 种 Anthropic 消息事件(ANTHROPIC_MESSAGE_EVENTS,331-338 行),用 parseJsonWithRepair 解析(491 行),并在流结束时校验「见过 message_start 就必须见到 message_stop」(506-508 行),否则报「流提前结束」。

第二步:六类上游事件 → 统一事件 ​

主循环在 601-789 行,是一棵大 if/else 树。对照表如下(源码事实):

Anthropic 事件做了什么push 出的统一事件
message_start记录 responseId、实际响应的模型、初始 usage无(start 已在 596 行发出)
content_block_start按 text / thinking / redacted_thinking / tool_use 建块text_start / thinking_start / toolcall_start
content_block_delta按 delta 子类型追加内容text_delta / thinking_delta / toolcall_delta;signature_delta 只累积签名,不发事件
content_block_stop收尾、清理临时字段text_end / thinking_end / toolcall_end
message_delta保存原始停止原因 rawStopReason、mapStopReason 映射停止原因、累计 usage 并算成本无
message_stop由 iterateAnthropicEvents 用于完整性校验无

停止原因的映射在 mapStopReason(1494-1520 行):end_turn → stop,max_tokens → length,tool_use → toolUse,refusal 与 sensitive → error 并带上说明文案,未知取值直接抛异常(宁可显式失败也不装作正常)。

图 6.1-2 一次 Anthropic 流式请求的事件转换
从上到下是时间顺序。注意 start 事件在 HTTP 响应头拿到之后、任何数据到达之前就发出,让界面能立刻显示「正在生成」。

这张图对应四段真实源码:A->>H 是 587-594 行经 retryProviderRequest 包裹的请求;A-->>U: start 是 596 行;中间三段 delta 来自 601-789 行的主循环;最后的 done 在 815 行、error 在 825 行。要留意 start 的位置——它在拿到响应之后立刻发出,晚于建立连接但早于第一个 token。

第三步:工具参数的增量 JSON ​

最需要单独说的是工具调用参数。Anthropic 把参数 JSON 拆成 input_json_delta 一片片发,中途的字符串必然是残缺的(例如 {"path":"READ)。pi-ai 每收到一片就尝试解析一次:

ts
// packages/ai/src/api/anthropic-messages.ts:699-711
					} 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,
							});
						}
earendil-works/pi@16787ad第 699–711 行在 GitHub 查看 ↗
每一片 JSON 增量都会更新 partialJson 缓冲,并用 parseStreamingJson 重新解析出一个「尽力而为」的参数对象。

parseStreamingJson 有四层兜底(源码事实):先按标准 JSON 解析;失败则修复非法转义与裸控制字符后再解析;再失败交给 partial-json 这个第三方库做残缺解析(依赖声明见 packages/ai/package.json:75);仍失败就返回空对象,绝不抛异常。

earendil-works/pi@16787ad第 104–124 行在 GitHub 查看 ↗
流式 JSON 的四级兜底解析,保证 partial 里的 arguments 永远是一个可用对象。

为什么要费这个劲?因为界面希望在参数流完之前就显示「正在读取文件 README.md」。有了增量解析,partial.content[i].arguments 在中途就已经是 { path: "READ" } 这样的半成品对象,界面可以直接读字段而不必自己拼 JSON。到 content_block_stop 时再整体重解析一次并删掉 partialJson 缓冲(740-743 行),保证存进 Session 的只有干净的参数对象。

要注意这属于 Anthropic 实现的行为:下面实践任务用的 faux provider 只按块发出 JSON 字符串片段,不做增量解析,参数对象是在 toolcall_end 之前整体赋值的(faux.ts:408 与 faux.ts:420)。同一个事件协议,不同实现对「中途的 partial 有多准」可以有不同投入。

相关测试:packages/ai/test/anthropic-sse-parsing.test.ts:329 的用例 "repairs malformed SSE JSON and malformed streamed tool JSON" 直接喂进畸形 SSE 与畸形工具 JSON,验证这条修复链路。

faux provider:把「假模型」做成正式设施 ​

3.1 讲过用假模型来学习和测试的想法。pi-ai 把它工程化成了一个正式的 provider:fauxProvider()。

ts
// packages/ai/src/providers/faux.ts:687-699(节选,api 对象压成了一行)
export function fauxProvider(options: RegisterFauxProviderOptions = {}): FauxProviderHandle {
	const core = createFauxCore(options);
	const provider = createProvider({
		id: core.provider,
		auth: { apiKey: { name: "Faux", resolve: async () => ({ auth: {} }) } },
		models: core.models,
		api: { stream: core.stream, streamSimple: core.streamSimple, fetchDeferred: core.fetchDeferred, cancelDeferred: core.cancelDeferred },
	});
	// …(省略:700-709 行返回句柄,含 setResponses / appendResponses / state 等)
earendil-works/pi@16787ad第 687–710 行在 GitHub 查看 ↗
faux 走的是和真 provider 完全相同的 createProvider 装配路径,只是 auth 永远成功、api 实现来自内存队列。

关键在于它不是 mock,而是一个真 provider:注册进 Models 之后,models.stream() 的每一环——鉴权、dispatch、事件流——都真实执行,只有最后一步「发 HTTP 请求」被替换成「从队列里取一条预设消息」。因此它能验证的东西远多于一个假函数:

  • 脚本化响应:setResponses([...]) 排队,stream 每次调用 shift 一条(faux.ts:507);队列空了会推出一条 errorMessage 为 "No more faux responses queued" 的 error 事件(faux.ts:513-524)。队列元素还可以是函数,按上下文动态决定回什么(faux.ts:109-114,它收到的也是规范化后的 TranscriptContext)。
  • 真的流式:streamWithDeltas()(faux.ts:339-435)把消息按 3–5 个 token(约 12–20 字符)随机切块(faux.ts:270-280),逐块发出全套 *_start / *_delta / *_end 事件,还响应 AbortSignal——中途取消会得到 stopReason: "aborted"。
  • 估算用量:withUsageEstimate()(faux.ts:230-268)按约 4 字符 = 1 token 估算,并且在给了 sessionId 时用「前缀公共长度」模拟 prompt cache 的命中与写入。
  • 连延迟响应也能模拟:请求带了 deferred 选项时,faux 先流出一条带句柄的 deferred 消息(faux.ts:526-552),之后再用 fetchDeferred 取回真正的结果(569 行起)。

注意 faux 不在 builtinProviders() 的 41 个内置 provider 列表里(providers/all.ts:90-134),必须自己 models.setProvider(faux.provider)。它从根入口导出(packages/ai/src/index.ts:36),所以 import 时不需要子路径。

Pi 自己就靠它测试:packages/ai/test/faux-provider.test.ts 有 23 个用例,其中 382 行的 "streams an exact event order for fixed-size chunks" 用固定块大小断言了完整事件序列;coding-agent 那边也有测试直接用它(例如 packages/coding-agent/test/fixtures/faux-session-worker.ts)。要留意的是,coding-agent 最常用的测试脚手架 test/test-harness.ts 并没有用 fauxProvider,而是自己定义了一个同样叫 faux 的模型和流函数(test-harness.ts:47-62)——名字相同,是两套东西。

实践任务 ​

🛠 实践任务用 faux provider 造一次流式假对话

目标:不需要任何 API Key,亲眼看到 models.stream() 吐出的统一事件序列,并验证「complete 就是 stream 的结果」这条设计。

准备:确认 Pi 仓库已装好依赖(仓库根目录存在 node_modules/.bin/tsx)。

步骤 1:在 Pi 仓库之外新建文件 faux-demo.mts(后缀必须是 .mts,原因见下方常见错误),把第一行的路径换成你本机 Pi 仓库的绝对路径:

text
import { createModels, fauxAssistantMessage, fauxProvider, fauxText, fauxToolCall } from "/绝对路径/pi/packages/ai/src/index.ts";

const faux = fauxProvider();
const models = createModels();
models.setProvider(faux.provider);
faux.setResponses([fauxAssistantMessage([fauxText("先看看 README"), fauxToolCall("read_file", { path: "README.md" }, { id: "call_1" })], { stopReason: "toolUse" })]);

const stream = models.stream(faux.getModel() as any, { messages: [{ role: "user", content: "介绍一下这个仓库", timestamp: Date.now() }] });
for await (const event of stream) console.log(event.type, "delta" in event ? JSON.stringify(event.delta) : "");
const final = await stream.result();
console.log("stopReason =", final.stopReason, "| totalTokens =", final.usage.totalTokens);

步骤 2:在 Pi 仓库根目录执行(把路径换成你的文件位置):

./node_modules/.bin/tsx ~/faux-demo.mts

预期现象(本书作者真实运行输出,Node 26 + 锁定版本 0.87.0 源码):

start 
text_start 
text_delta "先看看 README"
text_end 
toolcall_start 
toolcall_delta "{\"path\":\"README."
toolcall_delta "md\"}"
toolcall_end 
done 
stopReason = toolUse | totalTokens = 15

如何判断成功:① 事件序列符合 start → 各块的 *_start/*_delta/*_end → done;② 最后一行 stopReason = toolUse,与你在 fauxAssistantMessage 里写的一致;③ 你能解释 totalTokens 为什么不是 0(提示:faux.ts:230-268 的 4 字符/token 估算)。注意 toolcall_delta 的条数与切分位置每次运行都可能不同——分块大小是随机的 3–5 个 token、即 12–20 个字符(faux.ts:270-280),作者多次运行分别得到过 "{\"path\":\"REA" + "DME.md\"}"、"{\"path\":\"README." + "md\"}",以及一次就发完的 "{\"path\":\"README.md\"}";事件类型的先后顺序则是固定的。

常见错误:

  • 文件后缀写成 .ts 且放在没有 "type": "module" 的目录下,会报 Top-level await is currently not supported with the "cjs" output format。改成 .mts 即可(这是作者真实踩到的报错)。
  • 忘了调用 setResponses:只会看到一个 error 事件,errorMessage 为 No more faux responses queued(faux.ts:513-524)——这本身也是一次有用的观察。
  • 改用 models.complete(...) 却期待看到事件:complete 只返回最终消息,中间事件不会打印。想两者都要,就照上面这样先遍历流、再 await stream.result()。

对应源码位置:packages/ai/src/providers/faux.ts:687-710(装配)、505-564(从队列取一条并交给流)、339-435(切块并推送事件);packages/ai/src/models.ts:679-692(Models.stream)、695-701(complete)。

本章小结 ​

  • pi-ai 分四层:数据(types.ts 的四种消息与内容块)、协议(12 种 AssistantMessageEvent + EventStream)、路由(models.ts 的 Models 与 Provider)、实现(api/*.ts 每个文件一种线上协议)。
  • 系统提示词和工具声明住在 transcript 的 system 消息里:Context 只是输入简写,normalizeContext() 把它变成 provider 唯一接受的 TranscriptContext;中途追加 system 消息就能改提示词与工具。
  • 对外只有一个集合对象:createModels() / builtinModels() 造出来,方法分为查目录、管凭据、发请求三类。complete() 就是 stream().result(),包里没有第二条非流式代码路径。
  • 事件协议的三条约定:每个事件带 partial(共享的活对象,不是快照);块可能交错、必须用 contentIndex 关联;终止只有 done / error 两种,且失败编码进流而不是抛异常。done 的 reason 可以是 deferred,表示结果稍后用句柄取回。
  • EventStream 是一个「推-拉桥接」:队列接住早到的事件,waiting 数组兑现晚到的消费者,finalResultPromise 服务只要结果的调用方。
  • Anthropic 实现的完整链路:手写 SSE 解码 → 事件过滤与 JSON 修复 → 六类上游事件映射成统一事件 → mapStopReason 归一化停止原因;工具参数用 parseStreamingJson 做四级兜底的增量解析。
  • faux provider 是一个真 provider,只把「发 HTTP」换成「读队列」,因此能在没有 API Key 的情况下跑通除网络之外的全部路径。
  • 关键术语:统一消息类型(Unified Message)、system 消息与规范化上下文(TranscriptContext)、流事件协议(AssistantMessageEvent)、事件流缓冲(EventStream)、延迟响应(deferred)、线上协议(wire protocol)、懒加载 API(lazyApi)、增量 JSON 解析(partial JSON)、faux provider。
  • 关键源码索引:
    • packages/ai/src/types.ts:364-553(内容块与消息)、617-634(Context 与 TranscriptContext)、652-668(事件协议)、335-345(StreamFunction 契约);packages/ai/src/utils/transcript.ts(normalizeContext 与回放函数)
    • packages/ai/src/models.ts:163-235(Models)、648-755(applyAuth / stream / complete / deferred)、784-884(createProvider 与 dispatch)
    • packages/ai/src/utils/event-stream.ts:3-105、packages/ai/src/utils/json-parse.ts:104-124
    • packages/ai/src/api/anthropic-messages.ts:411-468(SSE)、601-789(事件映射)、1494-1520(停止原因)、858-904(streamSimple)
    • packages/ai/src/providers/faux.ts:339-710、packages/ai/src/providers/all.ts:90-143
    • 测试:packages/ai/test/anthropic-sse-parsing.test.ts、packages/ai/test/faux-provider.test.ts
  • 自测问题:① 为什么说 pi-ai 里没有非流式代码路径?② 模型拒答(refusal)时,await models.complete(...) 会抛异常吗?返回的消息里哪个字段能看出来?③ 工具参数流到一半时,partial 里的 arguments 是什么?由哪个函数保证它一定是个对象?④ faux provider 和「把 stream 函数换成一个假函数」相比,多验证了哪些环节?
  • 下一章:6.2 Provider 与模型注册——41 个内置 provider 是怎么组装的、models.generated.ts 这份模型目录从哪来、自定义 provider 需要提供什么。本章尚未展开的内容还有:transformMessages 的跨模型接力(历史消息如何在换模型时被改写)、models-store.ts 的动态模型目录持久化、以及 OAuth 登录流程——前两者在 6.2 讲,OAuth 属于支线,本书只在 7.4 自定义 Provider 顺带提及。
✅ 自测问题参考答案先自己回答,再点开对照
  1. 因为 complete() 就是 stream().result()——它只是把同一条流跑完再把终局消息交出来,包里没有第二条「一次性请求、一次性响应」的代码路径。好处是所有 Provider 只需实现流式一种协议,取消、用量统计、错误编码这些逻辑也只写一遍;非流式想要的东西(一条完整消息)本来就是流式的终局事件里现成的。反过来「非流式实现流式」做不到。
  2. 不会抛异常。error 事件本身就是终止事件,result() 正常 resolve 出一条消息而不是 reject——这是 StreamFunction 契约规定的:失败必须编码进流,不许 throw。看 stopReason 字段:mapStopReason 把上游的 refusal 与 sensitive 都映射成 "error",并在 errorMessage 里带上说明文案。所以调用方必须自己检查 stopReason,不能因为「没抛异常」就当成功了。
  3. 是一个**「尽力而为」的半成品参数对象**——参数 JSON 还没收完时,{"path":"READ 会被解析成 { path: "READ" } 这样能直接读字段的东西,界面因此可以在参数流完之前就显示「正在读取 …」。保证它一定是对象的是 parseStreamingJson(utils/json-parse.ts:104-124):标准解析 → 修复非法转义与裸控制字符后再解析 → 交给 partial-json 做残缺解析 → 仍失败就返回空对象,绝不抛异常。
  4. faux 不是 mock,而是一个真 provider:它走完全相同的 createProvider 装配路径,注册进 Models 之后,鉴权、dispatch 路由、EventStream 推拉桥接、事件序列生成、用量估算(含 prompt cache 命中模拟)、AbortSignal 取消——每一环都真实执行,只有最后一步「发 HTTP 请求」被换成「从队列里取一条预设消息」。而「把 stream 函数换成假函数」是从最外层截断,上面这些环节一个都不会被验证到。

本书分析的 Pi 版本:earendil-works/pi@16787ad(2026-09-21)