Appearance
第21章 双层循环——pi-agent-core
本章导读: 源码路线第 ③ 站(核心层)。从
pi-ai切到pi-agent-core包,读agent.ts和agent-loop.ts。这一站要看的是循环的真面目——模型输出和工具结果怎么一推一拉,把状态一轮轮往前送,以及那个容易被误读的"内外两层循环"到底是怎么回事。
小a读完流适配层,信心满满地跟老z说:“我懂了。模型流进来,循环就在那儿转——外层管模型轮次,内层管并行跑工具,对吧?”
老z没有直接回答,而是反问道:“你这句话里有几个‘想当然’?”
小a愣住了。
“流适配层只会产出事件;把事件变成上下文中的 assistant message、把 tool call 变成 tool result、再决定是否继续请求的,是 pi-agent-core 的 loop。”老z说,“你特别关心一个常见误读——源码中的内外循环是否就是‘外层模型轮次、内层并行工具循环’。答案是否定的,必须按状态变量和分支看。”
本章依据 pi-mono 提交 583f153d502aa8e958eefdb9af0fbd3344e68f95、@earendil-works/pi-agent-core 0.83.0。覆盖 runAgentLoop、runLoop、streamAssistantResponse 和 executeToolCalls;不讨论 Harness 的落盘、压缩或 Skills。
Agent 是带队列的宿主包装,不是循环本身
小a:“那 Agent 类呢?它不就是循环吗?”
“不是。”老z说,“packages/agent/src/agent.ts 的 Agent 持有 transcript state、listener 集、steering queue、follow-up queue 和当前 AbortController。它的 prompt() 最终调用 runAgentLoop;continue() 在最后消息不是 assistant 时调用 runAgentLoopContinue,而 assistant 结尾时会先尝试排空已排队的 steering 或 follow-up。”
ts
// packages/agent/src/agent.ts,Agent 的队列接口(583f153d,节选)
/** Queue a message to be injected after the current assistant turn finishes. */
steer(message: AgentMessage): void {
this.steeringQueue.enqueue(message);
}
/** Queue a message to run only after the agent would otherwise stop. */
followUp(message: AgentMessage): void {
this.followUpQueue.enqueue(message);
}“注意这两个队列的时机不同,不能因为都叫‘用户消息’而合并。”老z说,“PendingMessageQueue.drain() 还受 QueueMode 控制:all 取全部,否则取一条;默认是 one-at-a-time。调用 prompt() 时若已有 active run,源码直接抛错并提示使用 steer() 或 followUp()——避免两个 loop 并发改同一 state。”
| API/状态 | 已确认的作用 | 不能从中推出 |
|---|---|---|
prompt | 新增用户消息并启动一段 run | 可以并行开始多个 run |
continue | 从 user/tool result 结尾继续;assistant 结尾时只消费已排队的 steering/follow-up | assistant 结尾也可无条件继续 |
steer | 当前 assistant turn 后、下次请求前注入 | 立刻打断正在到来的 token |
followUp | Agent 否则会停止时再注入 | 每一轮都会插入 |
abort | abort 当前 controller | 已经执行的外部副作用自动撤销 |
subscribe | 顺序 await listener,传入 signal | listener 是持久化机制 |
“subscribe 的注释还说明 agent_end 是事件上的末尾,但 Agent 要等该事件的所有 listener settle 后才真正 idle。这个等待关系是观察一致性的一部分;并不等同 session 已经写盘,后者是 Harness 的职责。”
runAgentLoop 先建立本次增量
“那循环函数进来第一件事干嘛?”小a问。
“先建立本次增量。”老z打开 agent-loop.ts:
低层 runAgentLoop 接收 prompts、已有 context、config、event sink、signal、stream function。它克隆了本次 newMessages,再用历史 messages 加 prompts 构造 currentContext,发出 agent_start、首个 turn_start 和每条 prompt 的 start/end 事件,然后进入 runLoop。
ts
// packages/agent/src/agent-loop.ts,runAgentLoop(583f153d,节选)
const newMessages: AgentMessage[] = [...prompts];
const currentContext: AgentContext = {
...context,
messages: [...context.messages, ...prompts],
};
await emit({ type: "agent_start" });
await emit({ type: "turn_start" });
for (const prompt of prompts) {
await emit({ type: "message_start", message: prompt });
await emit({ type: "message_end", message: prompt });
}
await runLoop(currentContext, newMessages, config, signal, emit, streamFn ?? getDefaultStreamFn());“注意这里有两个列表,别混。”老z强调,“currentContext.messages 是下一次模型调用可见的完整工作上下文,newMessages 是本次 run 新增加、最后随 agent_end 返回的增量。混淆二者会错误解释为什么 prompts 既被 emit 又会进入模型上下文。”
“runAgentLoopContinue 的边界也明确:空 context 会 throw;最后一条是 assistant 时,只有两个队列都为空才 throw。这样不会把 assistant 回答直接接到下一次 provider 请求,同时仍允许已排队消息启动下一段 run;真正的 convertToLlm 只在每个 turn 调用时才能验证 provider 接受的消息形式。”
外层和内层循环分别解决什么
小a:“那内外两层 while 到底各管什么?”
“固定快照中 runLoop 在进入外层前读取一次 getSteeringMessages。外层 while (true) 的注释是‘queued follow-up messages arrive after agent would stop’;内层 while (hasMoreToolCalls || pendingMessages.length > 0) 的注释是‘process tool calls and steering messages’。”老z说,“所以不是‘内层并行工具循环’。并行只在后面的工具执行辅助函数中发生。”
ts
// packages/agent/src/agent-loop.ts,runLoop(583f153d,节选)
while (true) {
let hasMoreToolCalls = true;
while (hasMoreToolCalls || pendingMessages.length > 0) {
// 处理 pendingMessages,然后 streamAssistantResponse(...)
const message = await streamAssistantResponse(currentContext, config, signal, emit, streamFunction);
// tool result 加回 currentContext.messages 与 newMessages
pendingMessages = (await config.getSteeringMessages?.()) || [];
}
const followUpMessages = (await config.getFollowUpMessages?.()) || [];
if (followUpMessages.length > 0) { pendingMessages = followUpMessages; continue; }
break;
}正常一次工具轮的状态序列如下:
text
初始 prompts → emit user start/end → assistant stream
→ start/delta 更新 partial assistant → done 得最终 assistant
→ 找 content 内 toolCall
→ 执行/拒绝每个 tool call,产生 ToolResultMessage
→ 将 tool results 追加到 currentContext 与 newMessages
→ emit turn_end
→ prepareNextTurn 可替换 context、model、thinking
→ 取 steering;有 tool result 或 steering 时内层再请求模型
→ 否则取 follow-up;有则回到内层,无则 emit agent_end“assistant 没有 tool call 时,hasMoreToolCalls 会变为 false;只有 steering 能让内层继续。assistant 之后的 follow-up 则要等内层准备退出才读取。这正对应两类队列的语义,而非一个消息优先级列表。”
逐段跟 streamAssistantResponse
小a:“那模型调用这一层呢?循环怎么把上下文交给模型?”
“看 streamAssistantResponse——它是 AgentMessage[] 与 AI Message[] 的边界。”老z说,“先可选 transformContext,再 convertToLlm,用 system prompt、转换后 messages 和 tools 建 Context;接着刷新可过期的 API key,并调用 injected streamFunction。”
ts
// packages/agent/src/agent-loop.ts,streamAssistantResponse(583f153d,节选)
let messages = context.messages;
if (config.transformContext) messages = await config.transformContext(messages, signal);
const llmMessages = await config.convertToLlm(messages);
const llmContext: Context = { systemPrompt: context.systemPrompt, messages: llmMessages, tools: context.tools };
const resolvedApiKey =
(config.getApiKey ? await config.getApiKey(config.model.provider) : undefined) || config.apiKey;
const response = await streamFunction(config.model, llmContext, { ...config, apiKey: resolvedApiKey, signal });“这段的收益是 loop 不直接绑定 ModelsImpl:测试或宿主可注入 stream function,且 context 转换在每一 turn 都可以重算。代价是 transform/转换/key resolver 都可能失败;这段函数本身不 catch 它们,错误由调用它的 run 与 Harness/Agent 宿主处理。”
“下面是事件消费的核心分支。节选保留了状态改变,删去了全部 delta 类型的枚举。”
ts
// packages/agent/src/agent-loop.ts,streamAssistantResponse(583f153d,节选)
case "start":
partialMessage = event.partial;
context.messages.push(partialMessage);
addedPartial = true;
await emit({ type: "message_start", message: { ...partialMessage } });
break;
// text/thinking/toolcall delta:替换 context 的末条并 emit message_update
case "done":
case "error": {
const finalMessage = await response.result();
if (addedPartial) context.messages[context.messages.length - 1] = finalMessage;
else context.messages.push(finalMessage);
await emit({ type: "message_end", message: finalMessage });
return finalMessage;
}“先放 partial,再以 final message 替换末条,保证后续工具执行看的是已完成的 assistant content。即便没收到 start,done/error 分支仍会 append final message 并补发 message_start;迭代自然结束也会调用 response.result()。这覆盖了 adapter 的不同事件形态,不能据此假定所有 adapter 一定先发 start。”
停止原因先于工具执行
小a:“那模型说要调工具,就一定执行?”
“不一定。”老z回到 runLoop,“若 assistant 的 stopReason 是 error 或 aborted,它发 turn_end 和 agent_end 后返回,不执行工具。若 content 有 tool calls 且 stop reason 为 length,也不会执行任何一个,而由 failToolCallsFromTruncatedMessage 为每个 call 生成错误结果;注释说明解析器可能把被截断的 JSON 尽力修复,仍不安全执行。”
| assistant 状态 | 有 tool call 时的动作 | 是否继续内层 |
|---|---|---|
| 正常非 length | executeToolCalls | 取决于所有结果是否 terminate: true |
length | 每个 call 生成“arguments may be truncated”错误结果 | 是,让模型重新发完整调用 |
error | 不执行工具 | 否,立即 agent_end |
aborted | 不执行工具 | 否,立即 agent_end |
| 无 tool call | 无 tool result | 只可能由 steering 继续 |
“这至少给出两个中断/边界路径:模型中止停止 loop;输出长度截断保留 tool result 证据但拒绝副作用。把 length 当成‘可安全用已解析参数继续’会违背源码注释。”
executeToolCalls 决定顺序或并行
小a:“那并行呢?我一开始以为内层就是并行跑工具。”
“回到代码看。”老z打开 executeToolCalls:
ts
// packages/agent/src/agent-loop.ts,executeToolCalls(583f153d,节选)
const toolCalls = assistantMessage.content.filter((c) => c.type === "toolCall");
const hasSequentialToolCall = toolCalls.some(
(tc) => currentContext.tools?.find((t) => t.name === tc.name)?.executionMode === "sequential",
);
if (config.toolExecution === "sequential" || hasSequentialToolCall) {
return executeToolCallsSequential(currentContext, assistantMessage, toolCalls, config, signal, emit);
}
return executeToolCallsParallel(currentContext, assistantMessage, toolCalls, config, signal, emit);“它先搜出 assistant content 的所有 toolCall。**只要全局 toolExecution 为 sequential,或其中一个被找到的工具声明 sequential,整个 batch 走顺序路径;否则走并行路径。**因此,‘看到多个 tool calls 就一定并发’是错误前提。”
“两条路径都先发每个 tool_execution_start,并经过 prepareToolCall。准备会检查工具存在、预处理参数、validateToolArguments、可选 beforeToolCall hook 与 abort;任何异常、找不到工具、被 hook block 或 abort 都变成立即 error result,而不是抛掉整个 batch。”
“并行路径先按 assistant 中的顺序收集立即结果或 async thunk,随后 Promise.all 并发执行 prepared thunk,却按照 orderedFinalizedCalls 的原始顺序再生成 tool-result messages。这个选择的收益是独立工具可以同时运行而上下文顺序稳定;代价是共享资源工具必须显式标 executionMode: "sequential" 或自行串行化。”
“executePreparedToolCall 将 tool 的 progress callback 转为 tool_execution_update,等待已发出的 update event 后才返回。其 catch 把工具抛出的异常包成 error result;afterToolCall hook 还能修改结果,hook 自己抛错则也生成错误结果。最后 shouldTerminateToolBatch 只有在 batch 非空且所有 result 都 terminate === true 时停止下一次模型请求。”
| 工具边界 | 来源 | loop 可见的结果 | 后续效果 |
|---|---|---|---|
| 工具名不存在 | prepareToolCall | error tool result | 仍交模型决定下一步 |
| 参数无效/before hook block | 准备阶段 | error tool result | 不执行工具 |
| signal aborted | 准备或顺序循环中 | error result 或停止后续顺序项 | 不承诺撤销已启动工具 |
| tool.execute 抛错 | 执行阶段 catch | error tool result | 仍形成上下文证据 |
| 所有 result terminate | shouldTerminateToolBatch | terminate: true | 内层不因工具继续 |
收益、代价与未覆盖范围
基于源码的推断: 将低层 loop 的 stream function、消息转换、hook 和队列都做成注入项,使同一循环可服务不同宿主;事件又让 UI/持久化层观察每个状态变化。代价是状态分散于 context、newMessages、队列和 listener,扩展实现必须严肃处理取消、写入顺序和 listener 失败。
本章没有覆盖 Harness 怎样持久化 message_end,也没有证明任何 tool 的幂等性或取消能力。工具调用的实际副作用取决于具体 tool;并行执行不是安全保证。第22章 调度台——pi-agent-core会把事件接到 session repository、JSONL、compaction 和 skills。
小结
这一章回答的是:模型流进来之后,循环到底怎么"转"起来。答案是双层结构——runAgentLoop 建立本次增量和工作上下文;runLoop 的内层处理工具结果与 steering,外层只在原本停止后才读取 follow-up。streamAssistantResponse 用 partial 占位、最终消息替换来消费流,保证模型还没结束时,上下文里不会出现半截的 assistant 消息;executeToolCalls 依配置和工具声明选择顺序或并行,并把多个失败转换为可见的 tool result,而不是让异常静默吞掉。
- 边界:length、aborted、error 都不是普通成功结束——它们各自映射不同的恢复路径。被截断消息里的 tool call 不会进入
executeToolCalls,这是防止半截调用被执行的硬边界。
记住:外层轮次和内层工具循环不是一回事。 内层处理"这一次模型响应引发的所有工具结果",外层才处理"下一轮该不该继续"。把两层混成一个循环,是理解 pi-mono 循环设计最常见的误读。
源码走查
- 在
packages/agent/src/agent.ts阅读PendingMessageQueue.drain、steer、followUp;分别把 mode 设想为all与默认值,写下每次 drain 的消息数。 - 从
runAgentLoop读到runLoop;跟踪同一 prompt 如何同时进入currentContext.messages与newMessages。 - 在
runLoop给内外两个while标注各自继续条件;验证 follow-up 只在内层退出后读取,steering 在每个 turn 尾部读取。 - 逐行读
streamAssistantResponse的 start、done/error 与循环结束后 result 分支;确认没有 start 时也会把 final message 写入 context。 - 搜索
failToolCallsFromTruncatedMessage;核对length的 tool call 不会进入executeToolCalls。 - 比较
executeToolCallsSequential与executeToolCallsParallel;确认 parallel 用Promise.all执行 prepared thunk,却按原 call 顺序追加 tool-result message。