本文是「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 消息,并沿普通消息路径写入会话历史,后续模型请求可以读取失败原因。
主循环如何读取 steering 和 followUp
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用于没有 tool call 或 pending 消息、agent 原本将结束时追加输入。两者最终都写入 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 消息后自动发起请求。
下面的伪代码省略事件细节,只保留两层循环与队列读取边界。运行时只在 turn 开始前和结束后读取 steering 队列,因此 steering 不会中断当前流式调用。
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)。用户在模型生成时输入的补充约束会在 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 |
StreamFn 如何编码失败
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 消息上抛。
收尾逻辑可以概括为: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 拒绝整批调用:底层完成尽力解析,上层拒绝执行截断消息中的参数。
并行工具调度如何保持 transcript 源序
默认情况下,同一批 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)。先按 assistant 消息中的源顺序串行 preflight:一个 for 循环逐个 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 向 UI、扩展和持久化层公开 agent、turn、message 与 tool execution 的生命周期。请求正常完成、被中止或 provider 发生错误时,运行时都以闭合的事件序列和最终 assistant 消息结束本次执行。
输出因 length 截断时,运行时拒绝整批 tool call,避免不完整参数产生副作用。并行执行时,文件写入按路径互斥,tool execution 事件按完成顺序报告进度,toolResult 按 assistant 中的源顺序写入 transcript。工具输出的行数和字节数限制、继续读取指引,以及 edit 的模糊匹配,限制了模型输入不完整或格式偏差时的影响范围。
这些边界使流式生成、工具执行、会话记录和界面更新遵循同一套结算规则:每个 turn 都可被观察、持久化和恢复,失败与截断也会产生可消费的结果。
