Appearance
1.4 StreamFunction 契约
StreamFunction 是 pi 架构中最重要的一条边界:它把"@pi/ai 底层怎么调模型"和 "@pi/agent-core 上层怎么用 Agent" 彻底解耦。这一节讲清楚这个契约。
问题:Agent 不依赖具体厂商
@pi/agent-core 在构造 Agent 时,完全不知道背后是 OpenAI 还是 Anthropic。它只依赖一个注入进来的函数 —— StreamFn。这就是依赖倒置:上层定义接口,下层提供实现。
StreamFn 的定义
packages/agent/src/types.ts:28:
ts
export type StreamFn = (
model: Model<Api>,
context: Context,
options?: SimpleStreamOptions,
) => AssistantMessageEventStream | Promise<AssistantMessageEventStream>;它和 @pi/ai 里的 StreamFunction(packages/ai/src/types.ts:316)是同一个形状,只是应用层给它起了个更短的名字 StreamFn。
契约(必须遵守的约定)
StreamFn 不是普通的函数,它有严格的错误处理契约:
- 不得 throw,不得返回 rejected promise。所有失败都要编码进返回的事件流里。
- 必须返回
AssistantMessageEventStream。 - 失败通过流协议表达:发出
error事件,并以一条stopReason === "error" || "aborted"且带errorMessage的AssistantMessage结束。
为什么?因为 Agent 循环要保证事件序列永远完整(start → … → done/error)。如果 StreamFn 直接 throw,循环就崩溃了,无法产生正常的 agent_end 关闭事件。
ts
// 错误示范:直接 throw
const badStreamFn: StreamFn = () => { throw new Error("boom"); };
// 正确示范:把错误编码进流
const goodStreamFn: StreamFn = (model, context, options) => {
const stream = createAssistantMessageEventStream();
queueMicrotask(() => {
stream.push({ type: "error", reason: "error",
error: { role: "assistant", content: [], stopReason: "error",
errorMessage: "boom", /* ... */ } });
stream.end(/* the error message */);
});
return stream;
};事件流如何一步步被消费
Agent 循环消费 StreamFn 返回的流时,遵循固定模式(见 packages/agent/src/agent-loop.ts:317-361):
ts
for await (const event of response) {
switch (event.type) {
case "start":
// 一次解析到 partial,压入上下文,发 message_start
break;
case "text_delta": case "thinking_delta": case "toolcall_delta":
case "text_start": case "text_end": /* ... */
// 更新上下文里最后一条 partial,发 message_update
break;
case "done": case "error":
const finalMessage = await response.result(); // 完整消息
// 回填上下文,发 message_end
return finalMessage;
}
}为什么 partial 每次都要带完整消息
sse 服务端返回的是增量 delta,但 AssistantMessageEvent 的每个事件里都带 partial(当前累积出的完整消息)。这让消费端"拿到任何事件都能渲染出当前的完整状态",无需自己拼 delta。Faux Provider 的 streamWithDeltas 正是这样构造每次的 partial(见 faux.ts:338)。
setDefaultStreamFn:默认实现的注入
@pi/agent-core 本身不 import 任何厂商实现,但允许宿主注入一个默认 StreamFn,这样上层代码可以省略显式传参(packages/agent/src/stream-fn.ts):
ts
export function setDefaultStreamFn(streamFn) { defaultStreamFn = streamFn; }
export function getDefaultStreamFn() { /* 未设置则 throw */ }@pi/coding-agent 在入口处注册了默认实现:
ts
// packages/coding-agent/src/core/sdk.ts:36
setDefaultStreamFn(streamSimple); // 来自 @earendil-works/pi-ai/compat小结
StreamFn是 Agent 与模型之间的唯一边界。- 契约的核心:错误编码进流,绝不 throw。
- 事件流按
start → delta → done/error消费,partial携带完整累积状态。 - 通过
setDefaultStreamFn注入默认实现,保持运行时与厂商解耦。
真实源码位置
StreamFn:packages/agent/src/types.ts:28- 消费流程:
packages/agent/src/agent-loop.ts:317-361 setDefaultStreamFn:packages/agent/src/stream-fn.ts- 默认注入:
packages/coding-agent/src/core/sdk.ts:36
面试角度:为什么这样设计 StreamFunction
Q1:为什么 Agent 不直接依赖具体厂商,而是注入一个 StreamFn? 这是依赖倒置:上层(Agent)定义接口,下层(厂商)提供实现。如果 Agent 直接 import OpenAI 的 SDK,那换 Anthropic 就得改 Agent 代码。注入 StreamFn 后,Agent 完全不关心背后是谁,只要对方满足这个函数形状即可 —— 这正是分层架构里"Agent 运行时与模型解耦"的关键。
Q2:契约要求"不 throw、错误编码进流",违反会怎样? 如果 StreamFn 直接 throw,Agent 循环所在的 for await 会中断,无法发出完整的 agent_end 事件序列,导致状态机卡在"流式中"、UI 无法收尾。把错误编码成 error 事件 + stopReason=error 的消息,循环才能正常走完关闭流程。这是"保证事件序列永远完整"的底线。
Q3:为什么 partial 每次都要带完整消息而不只是增量 delta? 消费端(如 agent-loop)收到事件后要更新上下文里"正在生成的那条消息",如果每次只给 delta,消费端得自己累积拼接。给 partial,消费端直接用它覆盖即可,逻辑更简单、也更难出错。
Q4:setDefaultStreamFn 解决了什么问题? 让 @pi/agent-core 这个"内核"完全不 import 任何厂商实现,同时允许宿主(如 coding-agent)注入一个默认 StreamFn。这样上层代码可以省略显式传参,而内核依旧保持厂商无关。
下一步看 Demo:最小 Faux Provider。