本文是「Pi 源码拆解」系列第 2 篇。系列目录:
- 2026
- 06-19 Pi 源码拆解(一):极简 Coding Agent Harness 的分层设计
- 06-19 Pi 源码拆解(二): Agent 运行时机制(本篇)
上一篇介绍了 pi 的包结构和「刻意不做」清单。本文聚焦 packages/agent 中的 agent 运行时:主循环在 packages/agent/src/agent-loop.ts,全文件不到 800 行,不依赖 Node API。
本文前半介绍循环、事件、状态和挂点,后半分析运行时如何处理失败的。pi 将失败编码为 stopReason 为 error 的 assistant 消息,并沿普通消息路径写入会话历史。后续模型请求可以读取失败原因。
主循环:双层 while 与 steering 的生效边界
steering是 agent 已在运行时,用户追加的一条补充指令。例如,agent 正在重构src/auth.ts,用户再输入「保留现有 export」;这条输入不会终止当前模型流,也不会撤销已经产生的 tool call,而是先进入 steering 队列。当前 turn 完成并写入 assistant 消息、toolResult 后,运行时才在下一次模型请求前把这条消息加入 context。模型因此能同时读取上一轮执行记录和补充约束,再决定后续动作。
入口有两个:agentLoop() 接收新 prompt 并启动运行,agentLoopContinue() 不追加新消息、从当前 context 继续运行(agent-loop.ts:31、:64)。两者都汇入 runLoop()(agent-loop.ts:155)。
runLoop 采用双层 while。内层循环的进入条件为「存在待执行的 tool call,或队列中存在待注入消息」;每次迭代依次注入排队消息、流式获取 assistant 响应、执行 tool call 批次,并在收集结果后发出 turn_end。仅当内层循环结束、agent 即将结束运行且 followUp 队列返回新消息时(agent-loop.ts:263),外层循环才继续:运行时将消息加入 pending,并重新进入内层循环。
steering用于 agent 尚在运行时追加约束,消息在当前 turn 结束后的下一轮生效;followUp用于 agent 已经没有 tool call 或 pending 消息、原本将结束时再次追加输入。两者最终都写入 pending 并驱动新的 turn,区别在于运行时读取各自队列的时机。
在交互式 TUI 中,agent 正在流式生成时,普通 Enter 提交的输入会进入 steering 队列;Alt + Enter 提交的输入会进入 followUp 队列。前者适合补充会影响后续操作的约束,例如「保留现有 export」;后者适合让当前任务完整结束后再处理的下一项输入,例如「完成后再为改动补测试」。TUI 会分别显示 Steering: ... 和 Follow-up: ...。agent 空闲时,Alt + Enter 与普通 Enter 相同,都会直接启动一次新的 prompt,不会进入队列。
agentLoopContinue 要求 context 的最后一条消息不能是 assistant;不满足该条件时会直接抛错(agent-loop.ts:74-76)。最后一条消息必须能够经 convertToLlm 转换为 user 或 toolResult 消息,否则 provider 会拒绝该请求。该校验明确了继续运行与重试的前提:调用方需要先将 context 调整为 provider 可接受的消息序列,运行时不会在 assistant 消息后自动发起请求。
下面的伪代码省略事件细节,仅保留两层循环与队列读取边界。steering 不会中断当前流式调用,因为运行时仅在 turn 开始前和结束后读取该队列。
async function runLoop(context: Context) {
let pending = await getSteeringMessages();
do {
while (pending.length > 0 || hasToolCalls(context)) {
emit({ type: "turn_start" });
context.messages.push(...pending.drainOneOrAll());
const assistant = await streamAssistant(context); // 失败也返回一条消息
context.messages.push(assistant);
if (assistant.stopReason === "length") {
await failAllToolCalls(assistant.toolCalls);
} else {
await executeToolBatch(assistant.toolCalls);
}
emit({ type: "turn_end" });
pending.push(...await getSteeringMessages());
}
pending.push(...await getFollowUpMessages());
} while (pending.length > 0);
emit({ type: "agent_end" });
}用户追加输入进入 steering 队列,不会中止当前流式调用。runLoop 启动时先 poll 一次 getSteeringMessages(agent-loop.ts:167),之后每个 turn 结束再 poll(agent-loop.ts:259),取得的消息在下一轮流式请求前写入 context(agent-loop.ts:182-190)。用户在模型生成时输入的补充约束不会 abort 当前流,而是在 turn 边界生效,因此当前响应和工具结果仍会进入 transcript;代价是补充消息要等当前 turn 结束才生效。
以上图所示任务为例。T0,用户发送「重构 src/auth.ts」,消息进入 context,turn 1 开始请求模型。T1,模型正在流式分析并产生 read、edit 等 tool call。T2,用户补充「保留现有 export」。运行时把这条消息放入 steering 队列,不会向当前 AbortSignal 发 abort,也不会撤销已生成的内容。
T3,turn 1 的流自然结束,工具批仍按原来的参数执行并写入 toolResult,随后发出 turn_end。T4,loop 才 poll steering 队列、取出补充消息,并在 turn 2 的流式请求前将它写入 context。第二次请求因此同时带有 turn 1 的 assistant/toolResult 历史和「保留现有 export」约束;模型可以据此继续检查、补改或解释,turn 1 不会被半途截断。
外层 while 处理另一种时机:内层已没有待执行的 tool call 和 pending 消息,agent 原本将结束;此时 getFollowUpMessages() 如果返回消息,才把它放回 pending、重新进入内层。steering 是当前运行期间在 turn 边界追加的输入,followUp 是一次内层清空后决定是否再开一轮的输入;二者复用同一条 pending 管道,但 poll 时机不同。
队列由 PendingMessageQueue 管理(packages/agent/src/agent.ts:123),有两种 drain 模式:all 一次排空,one-at-a-time 每次只取最旧一条(agent.ts:139-152)。steering 和 followUp 两个队列默认都是 one-at-a-time(agent.ts:224-225)。每轮只注入最早的一条消息,模型可以分别响应各条补充输入,而不会在同一轮同时接收多条队列消息。
中止机制仅依赖 AbortSignal:同一个 signal 从 LLM 流传递至 tool 执行与 hooks。串行执行 tool 时,每完成一个调用都会检查 signal.aborted,以决定是否提前退出(agent-loop.ts:478)。
事件、状态与挂点
loop 的返回值是 EventStream<AgentEvent, AgentMessage[]>,agent_end 是终止事件。事件分四级生命周期(packages/agent/src/types.ts:422-437):agent(start/end)、turn(start/end)、message(start/update/end)、tool_execution(start/update/end)。UI 和扩展消费运行时的唯一通道就是这个事件流,loop 内部状态不对外暴露。
EventStream<T, R> 是一个同时提供事件迭代和最终结果的异步流容器(packages/ai/src/utils/event-stream.ts:4-66)。T 表示逐个产生的事件类型,R 表示运行结束时的结果类型。生产者通过 push(event) 写入事件;消费者可用 for await...of 按产生顺序读取。构造时传入的 isComplete 用于识别终止事件,extractResult 从该事件提取结果并兑现 result(): Promise<R>。在 agent loop 中,agent_end 同时作为事件流的最后一个事件和 AgentMessage[] 结果的来源,因此调用方可以一边更新 UI,一边等待 await stream.result() 取得本次运行产生的完整消息列表。
const stream = agentLoop(context, config);
for await (const event of stream) {
render(event);
}
const messages = await stream.result();四级事件的作用域与嵌套关系
这四类事件不是彼此独立的事件集,而是描述同一次 run 的嵌套作用域。agent 覆盖完整运行;turn 覆盖一次 assistant 响应及其请求的工具批次和工具结果;message 表示 transcript 中单条消息的增量更新与完成;tool_execution 表示单个工具调用的执行过程。包含工具调用的典型 turn 具有以下事件序列:
agent_start
turn_start
message_start(user) → message_end(user) // 仅首个 turn 的新 prompt
message_start(assistant)
message_update × n // text / thinking / tool-call 的流式增量
message_end(assistant)
tool_execution_start / update × n / end // 每个 tool call 一组
message_start(toolResult) → message_end(toolResult)
turn_end
agent_end下图展示了 agentLoop() 接收新 prompt 时的完整事件序列,其中 user 消息在 assistant 响应前发出 message_start 和 message_end。agentLoopContinue() 复用已有 context,不产生这对 user 消息事件;steering 和 followUp 消息则在后续 turn 注入时产生对应事件。
agent_start 和首个 turn_start 在 runAgentLoop() 调用 runLoop() 前同步发出(agent-loop.ts:109-116);agentLoopContinue() 也在进入 loop 前发出两者(:138-141)。后续 turn 不重新发 agent_start,而是在内层 while 下一轮开头发 turn_start(:175-179)。agent_end 有三条正常出口:assistant 的最终消息是 error 或 aborted(:196-200)、shouldStopAfterTurn 返回 true(:247-257),以及两个队列都空、外层循环退出(:262-274)。EventStream 把它定义为终止事件,终止值就是这次 run 产生的 messages(:145-149)。
turn_start 表示运行时准备发起一次 assistant 请求,不代表模型已经产出 token。turn 中先处理 pendingMessages:初始 prompt 在外层入口发 message 事件,steering 或 followUp 消息则在这里逐条发 message_start、message_end,随后加入 currentContext.messages(agent-loop.ts:181-190)。然后 streamAssistantResponse() 请求 provider。最终 assistant 消息落定后,loop 收集并执行 tool call,再以 assistant 消息和本批 toolResults 一起发 turn_end(:192-224)。因此订阅 turn_end 的 listener 能拿到一个可持久化的结算单元,而不用在每次 message_update 时写半截内容。
message 层覆盖 user、assistant、toolResult 三类 transcript 条目。user 与 toolResult 没有流式增量,分别只发 start/end;工具结果的发射在 emitToolResultMessage() 中完成。assistant 才有 message_update:provider 事件 start 到来时,loop 先把 partial assistant 消息放入 context 并发 message_start(agent-loop.ts:319-324);text_*、thinking_*、toolcall_* 事件每到一个,替换 context 尾部的 partial 消息,并转发为 message_update(:326-344);provider 发 done 或 error 时,调用 response.result() 取得最终消息,替换 partial 后发 message_end(:346-358)。如果 provider 在发出 start 前就结束,代码补发 message_start,再发 message_end(:354-357、:363-370),所以消费者总能得到成对的边界事件。
tool_execution 层不等同于 toolResult 消息。它的 start 在工具查找、参数预处理和 schema 校验之前发出,因此即使工具不存在、参数不合法或 beforeToolCall block 了调用,订阅方仍会收到 start 和带错误结果的 end(串行分支见 agent-loop.ts:444-474,并行 preflight 见 :499-515)。真正执行中的工具可以通过 onUpdate 发 tool_execution_update;执行体返回或抛错后,loop 经 afterToolCall 收尾,再发 tool_execution_end。并行模式下 end 按实际完成顺序发出;所有任务完成后,toolResult 的 message_start/end 才按 assistant 中原始 tool call 顺序发出(:522-547)。前者适合 UI 立即更新进度,后者保证 transcript 的确定性。
Agent 类在 loop 之外维护运行时状态。MutableAgentState(agent.ts:60)存 systemPrompt、model、tools、messages,外加 isStreaming、streamingMessage、pendingToolCalls 这些运行时态。每个事件先经过 processEvents() 这个 reducer 归约内部状态,再按注册顺序逐个 await listener(agent.ts:529-576)。具体来说,message_start/update 更新 streamingMessage,message_end 清空它并把最终消息加入 messages;tool_execution_start/end 往 pendingToolCalls 这个 Set 加入或删除 call ID;turn_end 从失败的 assistant 消息提取 errorMessage(:531-566)。所以 listener 运行时读到的是事件已经归约后的状态,而不是前一拍的状态。
有一项需要注意的语义:发出 agent_end 并不代表 run 已结束;只有全部异步 listener 完成、finishRun() 清理运行时状态后,Agent 才进入 idle 状态(agent.ts:522-528)。挂在 listener 上的持久化写入因此也计入这次 run 的结算,这一点第五篇讲 turn 边界落盘时会用到。
AgentLoopConfig 集中定义了扩展接口;本文涉及且后续文章会引用的接口包括:
| 挂点 | 作用 |
|---|---|
convertToLlm |
AgentMessage[] 投影为 Message[],发给 provider 前调用 |
transformContext |
请求前改写整个 context,compaction 挂这里 |
beforeToolCall |
工具执行前拦截,可 block |
afterToolCall |
工具执行后改写结果 |
prepareNextTurn |
turn 结束后更换模型、思考等级或 context |
getSteeringMessages / getFollowUpMessages |
两个队列的 poll 口 |
getApiKey |
每次请求动态取 key,为短寿命 OAuth token 设计(agent-loop.ts:305) |
shouldStopAfterTurn |
turn 结束后判定是否提前停止 loop |
never-throw 契约与失败消息
主线是 StreamFn 的 never-throw 契约。它的类型注释原文(types.ts:22-26):
Contract:
- Must not throw or return a rejected promise for request/model/runtime failures.
- Must return an AssistantMessageEventStream.
- Failures must be encoded in the returned stream via protocol events and a final AssistantMessage with stopReason “error” or “aborted” and errorMessage.
流函数不得抛出异常。请求失败、模型错误和运行时故障均编码到返回的事件流中:包含一条终止事件,以及一条 stopReason 为 error 或 aborted 的 AssistantMessage。pi-ai 侧 StreamFunction 的契约是同源表述(packages/ai/src/types.ts:314-319)。
该约定将失败表示为 assistant 消息的一种终止状态。loop 里 done 和 error 两个分支共用收尾代码,并把最终消息写进 context.messages(agent-loop.ts:346-358),因此失败会留在 transcript 中。后续模型请求或宿主策略可以据此决定重发请求、调整参数或停止。第六篇会看到 pi-ai 在更底层处理传输级重试,但最终失败仍以 error 消息上抛。
收尾逻辑可以概括成下面这样。统一收尾路径中的事件闭合比具体的 provider 异常类型更重要:消费事件流的 UI 和持久化 listener 因此不需要为「中途断流」另写一套恢复状态机。
const message = await streamFn(request, signal);
context.messages.push(message);
emit({ type: "message_end", message });
emit({ type: "turn_end", message });
if (message.stopReason === "error" || message.stopReason === "aborted") {
return; // agent_end 仍会在统一的收尾路径发出
}loop 层拿到 error/aborted 消息后,照发 turn_end 和 agent_end 收尾(agent-loop.ts:196-200),事件序列正常闭合。
未被前述契约覆盖的异常由一条兜底路径处理。Agent.runWithLifecycle 的 catch 调 handleRunFailure(agent.ts:496-512):人工合成一条失败消息(stopReason 按是否 abort 取 aborted 或 error),然后补发完整的事件序列:message_start、message_end、turn_end、agent_end,一个不少。即使 loop 内部哪里真的抛了异常,UI 和持久化层看到的也永远是一个闭合的事件序列,不会卡在「流开始了但没有结束」的悬挂状态。
相应地,hooks 也不得抛出异常。types.ts 里 convertToLlm、transformContext、getApiKey、shouldStopAfterTurn、两个队列 poll 函数的注释全是同一句「Contract: must not throw or reject」(types.ts:154、182、203、215、237、250),convertToLlm 和 shouldStopAfterTurn 两条还写了原因:抛异常会中断 loop 且产不出正常的事件序列。工具执行路径上的两个 hook 使用另一种处理方式:beforeToolCall 和 afterToolCall 的异常会被捕获,并转换为对应 tool call 的错误结果(agent-loop.ts:657-662、:743-746);其他调用仍按该批的调度规则处理。
截断即拒绝
第二种失败处理更激进。assistant 消息的 stopReason 为 “length” 时,说明输出被 token 上限切断,这条消息里每个 tool call 的参数都可能是半截。pi 的决策是整批拒绝执行(agent-loop.ts:211-214):
const executedToolBatch =
message.stopReason === "length"
? await failToolCallsFromTruncatedMessage(toolCalls, emit)
: await executeToolCalls(currentContext, message, config, signal, emit);failToolCallsFromTruncatedMessage 的注释(agent-loop.ts:374-380)说明了原因:流式 tool call 的参数由尽力修复的 JSON salvage parser 收尾,截断消息可能产生「能解析、也能通过 schema 校验,但内容不完整」的参数。edit 的 oldText 可能缺少后半段,bash 命令可能缺少管道符后的部分。为避免执行此类参数,运行时为每个 tool call 返回错误结果,说明响应撞上输出上限、参数可能被截断,并要求模型用完整参数重新发起(agent-loop.ts:396)。代价是增加一轮 round trip。
pi-ai 侧有一个配套的决策。parseStreamingJson(packages/ai/src/utils/json-parse.ts:104-124)做多级降级:先直接 parse(含 repairJson 修复),失败则用 partial-json 解析原串,再失败用 partial-json 解析修复后的串,最后兜底返回 {}。这个尽力解析的产物既喂流式展示,也在 tool call 收尾时成为最终参数(各 provider 的流式解析器都这么收尾,如 packages/ai/src/api/anthropic-messages.ts:695):完整消息里它等价于完整解析,截断消息里它就可能产出「能解析但悄悄不完整」的参数。所以执行层不能信它,stopReason=length 的整批拒绝正是堵这个口。两处决策共同构成分层处理:底层完成尽力解析,上层依据 stopReason 拒绝执行。
并行工具调度:preflight 串行、执行并发、落盘按源序
默认情况下,同一批 tool call 会并行执行,但需处理两个问题:多个 call 同时写入同一文件时如何协调,以及如何保证事件与 transcript 的顺序确定性。
执行管线分三段。prepare(agent-loop.ts:600-664):找工具、prepareArguments 预处理参数、schema 校验、跑 beforeToolCall hook,被 block 的 call 直接变成错误结果,不进入执行。execute(agent-loop.ts:666-707):调 tool.execute,带 onUpdate 回调发 tool_execution_update 流式部分结果,工具抛异常被捕获转成错误结果。finalize(agent-loop.ts:709-754):跑 afterToolCall hook,逐字段覆盖结果,content、details、usage、terminate、isError 各管各的,注释明确说没有深合并(types.ts:66-78)。
并行分支的调度分三层(agent-loop.ts:489-554)。第一层,preflight 串行:一个 for 循环按 assistant 消息里的源顺序逐个 prepare,能立刻出结果的(工具不存在、校验失败、被 block)当场 finalize,需要执行的不执行,包成 thunk 存进数组。第二层,Promise.all 并发执行所有 thunk,tool_execution_end 在每个工具 finalize 后立刻发出,按完成顺序。第三层,全部完成后按数组顺序(即 assistant 源顺序)逐个生成 toolResult 消息、发 message_start/end、落进 context。types.ts 的注释(:39-41)说明了这一解耦:tool_execution_end 按完成顺序发出,便于展示实际执行进度;tool-result 消息按源顺序写入 transcript,使同一 assistant 消息产生稳定的 toolResult 序列。
对应的伪代码如下。这里 results 虽然由并发任务填充,却始终以输入数组的下标作为位置;最后的 commit 循环不会受实际完成先后影响。
const prepared = toolCalls.map(prepareInSourceOrder);
const results = await Promise.all(
prepared.map(async item => {
if (item.readyResult) return item.readyResult;
const result = await item.execute();
emit({ type: "tool_execution_end", result }); // 完成序
return finalize(result);
}),
);
for (const result of results) {
const toolResultMessage = toToolResultMessage(result);
context.messages.push(toolResultMessage); // 源序
emit({ type: "message_end", message: toolResultMessage });
}默认并行执行依赖文件写入操作之间的互斥保证。edit 和 write 的执行体都包在 withFileMutationQueue 里(packages/agent/src/harness/tools/file-mutation-queue.ts:29,edit.ts:92、write.ts:28 调用),按 canonical path 维护 promise 链,同一路径的写操作串行,不同路径照旧并行。
此外还有两个执行控制项。工具可以声明 executionMode: "sequential",批里任何一个工具声明了,整批就走串行分支(agent-loop.ts:419-422)。早停的门槛是整批所有 tool result 都标了 terminate: true 才停(agent-loop.ts:582-584),单个工具想提前结束 run 是不够的,防止某个工具单方面掐断其他 call 的结果。
工具的防御性设计
最后考察具体工具的防御性实现。这些设计基于两个工程假设:输出规模可能超出上下文限制,输入参数可能不符合预期格式。
截断是双限,2000 行或 50KB,先到先触发(packages/agent/src/harness/utils/truncate.ts:11-13)。read 用 truncateHead 保留开头,bash 用 truncateTail 保留结尾(错误和最终结果通常在末尾),除 bash 尾部截断里单行自身就超字节限的边界情形外都不返回半行(truncate.ts:9 的注释写明这个例外)。字节计数优先用 Buffer.byteLength,没有 Buffer 的运行时降级到手写的 UTF-8 长度计算,外加未配对 surrogate 的替换(truncate.ts:54-110),保证截断不产出半个字符。
截断之后的引导是直接写给模型看的行动指令。bash 截断时全量输出落进临时文件,结果尾部附路径(packages/agent/src/harness/tools/bash.ts:131-140),模型想要完整内容可以自己 read。read 遇到单行就超 50KB 的情况,返回的不是错误而是建议:sed -n '123p' file | head -c 51200(packages/agent/src/harness/tools/read.ts:121)。普通截断则附 Use offset=N to continue(read.ts:128-136)。截断信息本身就是 prompt 的一部分。
edit 的 fuzzy 匹配用于处理模型重现代码时产生的格式偏差:normalizeForFuzzyMatch 做 NFKC 规范化、智能引号和破折号转 ASCII、去行尾空白(packages/coding-agent/src/core/tools/edit-diff.ts:33),exact 匹配失败后在规范化空间里再找(edit-diff.ts:206)。找到之后只重写被触及的行块,未改的行保留原始字节(edit-diff.ts:131-172),避免把全文都规范化掉、造成无谓的 diff。另外 prepareArguments 里有一个兼容 shim:有模型(注释点名 Opus 4.6、GLM-5.1)会把 edits 数组发成 JSON 字符串,这里尝试 parse 回数组(packages/coding-agent/src/core/tools/edit.ts:101-106)。
总结
packages/agent 的运行时将一次 agent 执行组织为可结算的 turn 序列。用户在运行过程中追加的输入通过 steering 或 followUp 队列在明确的边界进入 context;EventStream 则将 agent、turn、message 和 tool execution 的生命周期公开给 UI、扩展及持久化层。无论请求正常完成、被中止或发生 provider 错误,运行时都以闭合的事件序列和最终 assistant 消息结束本次执行。
工具调用是该结算模型的另一部分。输出因 length 截断时,运行时拒绝执行整批 tool call,避免不完整参数产生副作用;并行执行时,文件写入按路径互斥,tool execution 事件按完成顺序报告进度,toolResult 则按 assistant 中的源顺序写入 transcript。工具输出的行数和字节数限制、继续读取指引,以及 edit 的模糊匹配,共同限制了模型输入不完整或格式偏差时的影响范围。
这些约束将流式生成、工具执行、会话记录和界面更新统一到同一套边界规则中:每个 turn 都能被观察、持久化和恢复,失败与截断也具有与正常路径一致的可消费结果。

