Skip to content

第20章 先开票,后办事——pi-ai ​

本章导读: 源码路线第 ② 站(核心层)。仍在 pi-ai 包,但深入 api/lazy.ts 和 api/transform-messages.ts。这一站要拆开那个让小a发愣的 AssistantMessageEventStream——厂商零散的增量事件,是怎么被收拢成循环能直接消费的统一形态的。

17 章走完,小a终于让模型“出发”了。可 Provider 交回来的东西,又一次让他愣住——不是一段完整回答,而是一个 AssistantMessageEventStream。他盯着这个流,问老z:“我到底该怎么接?等它全部结束,还是一个一个拿?”

“都不对。”老z说,“第19章 目录与户口——pi-ai已经完成‘模型对象到 Provider’的分派,但 Provider 交回的不是一个完成的字符串,而是 AssistantMessageEventStream。你现在要核对三件事:异步设置为什么不会阻塞入口;一个真实 adapter 如何把厂商流变成统一事件;取消和异常最终以什么状态向上报告。”

本章证据固定在 pi-mono 提交 583f153d502aa8e958eefdb9af0fbd3344e68f95、@earendil-works/pi-ai 0.83.0。选择 packages/ai/src/api/openai-responses.ts 作为一个真实 adapter,不是声称所有 Provider 或 API 都采用它的字段和失败策略。不覆盖所有厂商协议、模型定价规则、WebSocket adapter,也不把通用 retry 工具的存在写成所有调用都会重试。

返回流之前,先把设置留到稍后 ​

“ModelsImpl.stream 与 streamSimple 都调用 lazyStream。这意味着函数可以立刻给调用者一个可迭代的流,而认证解析、动态模块加载、Provider 调用在后台进行。”老z说,“它不是‘请求一定已经开始’,更不是一个普通 Promise<AssistantMessage>。”

ts
// packages/ai/src/api/lazy.ts,lazyStream(583f153d,节选)
export function lazyStream(model, setup): AssistantMessageEventStream {
	const outer = new AssistantMessageEventStream();
	setup()
		.then((inner) => forwardStream(outer, inner))
		.catch((error) => {
			const message = createSetupErrorMessage(model, error);
			outer.push({ type: "error", reason: "error", error: message });
			outer.end(message);
		});
	return outer;
}

小a:“为什么要这样绕一圈?直接等 setup 完成再返回不行吗?”

“逐段看。”老z指着代码说,“第一行构造 outer,因此调用点同步拿到同一类型的输出,不必区分‘尚在认证’和‘已在接收网络事件’。第二段不 await setup(),而以 promise continuation 接上 forwardStream——这正是延迟设置:ModelsImpl 中的 requireProvider、applyAuth 以及 provider.stream* 被放进传入的 async 回调。这里的‘延迟’是相对同步返回而言,并不承诺等到第一次 for await 才执行,因为 setup() 已立即调用。”

“第三段将 setup rejection 归一为一个 AssistantMessage:空 content、零 usage、stopReason: "error"、错误文字;随后 push error 事件并 end(message)。所以未知 Provider 或认证失败对消费者呈现为流内错误,不是调用 Models.stream 时的同步 throw。”

比喻:先开票,后办事

lazyStream 像先给你一张取餐号(outer stream),后厨(setup)慢慢做。做砸了,也是通过“取餐口”告诉你“这单做不了”,而不是在你下单那一刻就摔盘子。

forwardStream 逐个 push 内流事件,迭代结束后若内流带有 .result() 就用它结束 outer。这一规则保留了 adapter 给出的最终消息;不能从这里推断任何事件一定包含 text delta。

阶段产物正常下一步设置失败时
Models.stream* 返回outer stream调用者订阅/迭代尚无同步异常
setup()inner async iterableforwardStream.catch 创建 error message
inner eventsstart、delta、done/error 等原样推到 outeradapter 决定事件内容
inner 完结result() 或 undefinedouter.end由 adapter 的终止行为决定

lazy API 是第二层延迟,不是第二种事件协议 ​

小a:“那 Provider 内部的适配器模块呢?也是懒加载?”

“对,但那是第二层延迟。”老z打开 lazy.ts:

lazyApi(load) 返回一个 ProviderStreams 对象。其两个方法都再以 lazyStream 包住 await load() 后的模块方法:

ts
// packages/ai/src/api/lazy.ts,lazyApi(583f153d,节选)
return {
	stream: (model, context, options) =>
		lazyStream(model, async () => (await load()).stream(model, context, options)),
	streamSimple: (model, context, options) =>
		lazyStream(model, async () => (await load()).streamSimple(model, context, options)),
};

“外层 ModelsImpl 的 lazy stream 覆盖‘Provider/认证设置’;这里的 lazy stream 覆盖‘适配器模块加载’。二者都产出同一种 AssistantMessageEventStream,因此可能出现嵌套转发;模块 import cache 是否去重由宿主运行时负责,源码注释明确提到 host import cache,但没有为任意运行环境证明性能结果。”

“收益是未使用的 adapter 不必在根入口加载;代价是模块加载错误推迟到第一次请求,且排错时要区分 auth setup、import setup 和网络请求三层。”

真实 adapter 的一次正常事件流 ​

小a:“那挑一个真实 adapter,看它正常怎么走?”

“看 openai-responses.ts。”老z说,“它的 stream 在函数开始就创建 AssistantMessageEventStream,然后启动一个 async IIFE。它先建一个 stopReason: "pending" 的 output,成功后才推 start、让 processResponsesStream 更新它,再推 done。”

ts
// packages/ai/src/api/openai-responses.ts,stream(583f153d,节选)
const { data: openaiStream, response } = await retryProviderRequest(
	() => client.responses.create(params, requestOptions).withResponse(),
	{ maxRetries: options?.maxRetries, maxRetryDelayMs: options?.maxRetryDelayMs, signal: options?.signal },
);
await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model);
stream.push({ type: "start", partial: output });
await processResponsesStream(openaiStream, output, stream, model, { /* 节选 */ });
stream.push({ type: "done", reason: output.stopReason, message: output });
stream.end();

“注意这段的顺序。”老z提醒,“请求经过 retryProviderRequest,但 adapter 明确传入 maxRetries: 0 给 OpenAI SDK 自身,并把 retry 策略交给外层 helper。是否重试仍由 options.maxRetries、信号和 helper 的可重试判定决定,不能把代码存在简化为‘任何错误必重试’。”

“onResponse 发生在拿到 HTTP response 后、推 start 前;它提供 status 和转换成 record 的 headers。随后 processResponsesStream 是上游事件到统一 delta 的实际转换位置;本章不展开其每个 OpenAI event case,原因是那会把厂商协议细节误当成跨 Provider 契约。”

一个正常序列可读作:

text
buildParams / createClient
  → retryProviderRequest 发起 Responses 请求
  → 收到 HTTP response,调用 onResponse
  → push start(partial output)
  → processResponsesStream 连续修改 output 并 push text/thinking/tool-call 事件
  → output 得到非 pending、非 aborted/error 的 stopReason
  → push done(message),end()

“对 Agent 来说,最关键的不是上游 event 名,而是它会先看到一个可更新的 partial,再得到最终的 AssistantMessage。这也是下章 streamAssistantResponse 可以把 start 写入上下文、把 delta 转成 message_update 的原因。”

transform 与 compat:统一形状不等于抹平差异 ​

小a:“前面 applyAuth 里有 transformHeaders,这里又有 getCompat——它们是同一回事吗?”

“不是。它们解决的层次不同。”老z说,“ModelsImpl.applyAuth 中的 transformHeaders 是请求发送前的 Models 专用变换;openai-responses.ts 又会针对模型的 compat 建立默认兼容配置。”

ts
// packages/ai/src/api/openai-responses.ts,getCompat(583f153d,节选)
function getCompat(model: Model<"openai-responses">): Required<OpenAIResponsesCompat> {
	return {
		supportsDeveloperRole: model.compat?.supportsDeveloperRole ?? true,
		sessionAffinityFormat: model.compat?.sessionAffinityFormat ?? detectSessionAffinityFormat(model),
		supportsLongCacheRetention: model.compat?.supportsLongCacheRetention ?? true,
		// 其余能力开关已删节
	};
}

“transformHeaders 的输入是已经合并的 header record,适合宿主加追踪或路由信息;它不会改 adapter 的 message 结构。getCompat 从 model.compat 和 provider/baseUrl 推断能力开关,随后被 buildParams、grammar tool 输入等路径使用。它说明同一 API 名下仍可有 capability 差异,不能说明默认 true 的能力一定被远端接受。”

机制发生位置作用对象失败或限制
transformHeadersModelsImpl.applyAuth已合并 header抛错成为 lazy setup error
onPayloadadapter 组装 params 后provider payload返回值可替换 params;hook 失败进入 catch
getCompatadapter 参数构造前Model capability 标志默认值是本地选择,不是服务端证明
processResponsesStream请求建立后上游事件→统一事件具体字段只适用于该 adapter

小a指着表里 getCompat 一行问:“那个 sessionAffinityFormat 默认走 detectSessionAffinityFormat,它到底看什么?”

“两处判断。”老z说,“一处是 model.provider === "openrouter" 或 baseUrl 里含 openrouter.ai,任中其一就选 openrouter,否则默认 openai。这是本地启发式,不是服务端声明——之后 createClient 建会话头时用上它:有 sessionId 且是 openrouter,就发 x-session-id;是 openai 分支则发 session_id 并附 x-client-request-id。所以 compat 标志真正改变了发送的 header,而默认值只代表本地选型。”

“那这个分支算错误路径吗?”小a追问。

“更值得记的是 getClientApiKey 的位置差异。stream 里它在 async IIFE 的 try 内,缺 key 抛错会进 catch、变成 error event;而 streamSimple 在构造 stream 之前就同步调用它。直接调 adapter 的 streamSimple 时,缺 key 是同步 throw,不走流内错误——只有经由 Models 集合的 async setup 闭包,才被归一成 error event。这恰好印证了上节那句:绕过 ModelsImpl,契约要你自己满足。”

abort、错误与不完整流不能被吞掉 ​

小a:“那流中途出问题呢?abort、没结束就断、报错——都怎么处理?”

“真实 adapter 的尾部有三道检查:请求后若 signal.aborted 就抛错;上游流结束仍是 pending 就报‘没有 stop reason’;若 output 已是 aborted/error 也转为错误路径。这防止一个没有完整终止信息的流被误报为成功。”

ts
// packages/ai/src/api/openai-responses.ts,stream 的收尾与 catch(583f153d,节选)
if (options?.signal?.aborted) throw new Error("Request was aborted");
if (output.stopReason === "pending") throw new Error("OpenAI Responses stream ended without a stop reason");

// catch 中:
output.stopReason = options?.signal?.aborted ? "aborted" : "error";
output.errorMessage = formatOpenAIResponsesError(error);
stream.push({ type: "error", reason: output.stopReason, error: output });
stream.end();

“catch 之前还删除了 content block 上仅用于流式解析的 index、partialJson、customInput。这是一个可确认的清理行为,不能被夸张成‘所有中间数据绝不留存’,因为其他层的日志与 store 不在这个函数内。”

路径触发条件事件/最终 stop reason可恢复性边界
正常完成parser 填好有效 stop reasondone,原 stop reason调用者可继续处理工具调用
主动中止AbortSignal 已 abortederror,aborted是否重试由上层决定
请求/解析异常SDK、hook、参数或网络抛错error,errorerrorMessage 经 formatter 规范化
上游无终止原因结束仍 pendingerror,error防止把半截响应当成功
lazy setup 失败auth、Provider 或模块加载拒绝error,error没有进入本 adapter 的正常 start

小a盯着表问:“中止为什么也用 error event,而不是 done?”

“这里事件类型表达‘非正常完成通道’,而 reason 区分 aborted 与 error。”老z说,“消费者不能只看 event type;第21章 双层循环——pi-agent-core会看到 loop 依 stop reason 停止,而不是把 aborted 当工具调用为空。”

“那 formatOpenAIResponsesError 会把 HTTP 状态丢掉吗?我在 errorMessage 里能读到什么?”小a问。

“它是 formatProviderError(normalizeProviderError(error), "OpenAI API error") 的合成。normalize 先按 SDK 字段形状探测状态——statusCode、status、$metadata.httpStatusCode、$response.statusCode 按序取第一个数字;body 则只接受普通对象或字符串,类实例和流对象会被当作没有 body,并截断到 4000 字符。format 只在能取到状态和 body 时才拼成 前缀 (状态): body;否则保留 SDK 自己的 error.message,避免用噪音替换有用文本。”

“所以 errorMessage 是一个显示用的合成字符串,不是结构化错误对象。”老z总结,“它有前缀、尽力保留状态和 body,但消费者若要做分支,应优先看 stopReason 和事件本身,而不是解析这段文本——aborted 与 error 的差别已经由 reason 给出,文本只是给人读的。”

消费方应保存最终消息,而非自行拼接文本 ​

小a:“那我是不是把 text_delta 拼起来,就拿到完整回答了?”

“会丢东西。”老z说,“adapter 在流中发布 partial event,是为了让消费方及时显示或更新状态。**最终 AssistantMessage 还包含 usage、provider、model、stop reason,以及可能的 tool call content。**因此只把 text_delta 拼成字符串会丢掉 Agent 下一步所需的信息。”

消费策略能得到什么丢失或风险
只显示 delta低延迟文本体验tool call、usage、stop reason 不完整
只等待 result完整 message没有流式可见性
同时消费 event 与 resultpartial UI 加最终状态必须避免把 partial 当持久化终值
收到 error 立即忽略可快速返回抹去 aborted 与 provider error 的区别

“lazyStream 的 forwardStream 不重写内层 event。所以 event 序列的细粒度语义仍由 adapter 决定。但 setup failure 会被统一合成为 error event,这正是调用者能以同一个订阅面处理配置失败的原因。这不是‘所有失败都有相同 message 文本’的保证。createSetupErrorMessage 用 error.message 或 String(error) 填写 errorMessage;adapter 则可能使用自己的 formatter。”

“重试也要在消费语义之外看。retryProviderRequest 覆盖的是本 adapter 发出 Responses 请求时的那段 await——它不包围之后 Agent 工具执行,也不包围 session append。若请求已开始输出再中断,最终是否可以重试还取决于调用者是否接受重复生成。**对有外部副作用的 tool result,不能把模型请求重试规则直接搬过去。**这是一条通用工程判断,不是 OpenAI Responses adapter 的自动保护。”

收益、代价与未覆盖范围 ​

基于源码的推断: 将 API adapter 统一为事件流,Agent 可以在同一套循环中显示文本、思考和 tool-call 增量;同时 Provider 仍能通过 typed options、compat 和 payload hook 保留差异。代价是错误来源跨多层,且一个通用事件类型不能自动抹平计费、缓存或取消的远端语义。

本章没有验证真实服务是否响应、重试退避的全部算法、processResponsesStream 的每一条事件映射,也没有比较 Anthropic、Google 或 WebSocket adapter。这些需要分别读各实现和在受控凭据下运行,不能由 OpenAI Responses 的代码替代。

小结 ​

这一章回答的是:Provider 交回的那个 AssistantMessageEventStream,到底该怎么接。答案是它不是一个"等全部结束再读"的完整对象,而是一条有生命周期的增量事件流——lazyStream 让入口立即给出统一流句柄,把认证、模块加载的失败统一转换成 error message;adapter 的正常路径是请求、start、事件转换、done,错误路径用 reason 区分 abort 与 error,并拒绝没有 stop reason 的半截流。

  • 边界:header transform、payload hook 与 compat 分别在不同层解决不同问题——它们不互相替代。消费者不能只看 event type,还要看 reason:aborted 与 error 是两条不同的恢复路径。

记住:流是有生命周期的,不是一次返回。 增量事件、stop reason、半截流拒绝——这三样合起来定义了"流的终止形态"。循环要依据 stop reason 决定下一步,而不是把 aborted 当成"工具调用为空"。

源码走查 ​

  1. 阅读 packages/ai/src/api/lazy.ts 的 createSetupErrorMessage、forwardStream、lazyStream;确认 setup 错误用零 usage 的 assistant message 结束 outer stream。
  2. 在 packages/ai/src/models.ts 搜索 lazyStream(model;分别标出 stream 和 streamSimple 的 setup 闭包,确认它们都在闭包内调用 Provider。
  3. 打开 packages/ai/src/api/openai-responses.ts 的 stream;从 retryProviderRequest 跟到 stream.push({ type: "start",写下 start 出现在 HTTP response 的哪个阶段。
  4. 搜索 getCompat 与 buildParams;选择一个 compatibility flag,确认它被用于参数构造,而非由根 Models 自动决定。
  5. 在同一 stream 的 catch 中检查 partialJson 删除和 stopReason 赋值;分别触发信号已 aborted 与普通 Error 时,推演 reason 的差异。