Skip to content

1.4 StreamFunction 契约

StreamFunctionpi 架构中最重要的一条边界:它把"@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 里的 StreamFunctionpackages/ai/src/types.ts:316)是同一个形状,只是应用层给它起了个更短的名字 StreamFn

契约(必须遵守的约定)

StreamFn 不是普通的函数,它有严格的错误处理契约

  1. 不得 throw,不得返回 rejected promise。所有失败都要编码进返回的事件流里。
  2. 必须返回 AssistantMessageEventStream
  3. 失败通过流协议表达:发出 error 事件,并以一条 stopReason === "error" || "aborted" 且带 errorMessageAssistantMessage 结束。

为什么?因为 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 注入默认实现,保持运行时与厂商解耦。
真实源码位置
  • StreamFnpackages/agent/src/types.ts:28
  • 消费流程:packages/agent/src/agent-loop.ts:317-361
  • setDefaultStreamFnpackages/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