Appearance
第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 iterable | forwardStream | .catch 创建 error message |
| inner events | start、delta、done/error 等 | 原样推到 outer | adapter 决定事件内容 |
| inner 完结 | result() 或 undefined | outer.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 的能力一定被远端接受。”
| 机制 | 发生位置 | 作用对象 | 失败或限制 |
|---|---|---|---|
transformHeaders | ModelsImpl.applyAuth | 已合并 header | 抛错成为 lazy setup error |
onPayload | adapter 组装 params 后 | provider payload | 返回值可替换 params;hook 失败进入 catch |
getCompat | adapter 参数构造前 | 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 reason | done,原 stop reason | 调用者可继续处理工具调用 |
| 主动中止 | AbortSignal 已 aborted | error,aborted | 是否重试由上层决定 |
| 请求/解析异常 | SDK、hook、参数或网络抛错 | error,error | errorMessage 经 formatter 规范化 |
| 上游无终止原因 | 结束仍 pending | error,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 与 result | partial 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 当成"工具调用为空"。
源码走查
- 阅读
packages/ai/src/api/lazy.ts的createSetupErrorMessage、forwardStream、lazyStream;确认 setup 错误用零 usage 的 assistant message 结束 outer stream。 - 在
packages/ai/src/models.ts搜索lazyStream(model;分别标出stream和streamSimple的 setup 闭包,确认它们都在闭包内调用 Provider。 - 打开
packages/ai/src/api/openai-responses.ts的stream;从retryProviderRequest跟到stream.push({ type: "start",写下 start 出现在 HTTP response 的哪个阶段。 - 搜索
getCompat与buildParams;选择一个 compatibility flag,确认它被用于参数构造,而非由根Models自动决定。 - 在同一
stream的 catch 中检查partialJson删除和stopReason赋值;分别触发信号已 aborted 与普通 Error 时,推演reason的差异。