Skip to content

第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-upassistant 结尾也可无条件继续
steer当前 assistant turn 后、下次请求前注入立刻打断正在到来的 token
followUpAgent 否则会停止时再注入每一轮都会插入
abortabort 当前 controller已经执行的外部副作用自动撤销
subscribe顺序 await listener,传入 signallistener 是持久化机制

“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 时的动作是否继续内层
正常非 lengthexecuteToolCalls取决于所有结果是否 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 可见的结果后续效果
工具名不存在prepareToolCallerror tool result仍交模型决定下一步
参数无效/before hook block准备阶段error tool result不执行工具
signal aborted准备或顺序循环中error result 或停止后续顺序项不承诺撤销已启动工具
tool.execute 抛错执行阶段 catcherror tool result仍形成上下文证据
所有 result terminateshouldTerminateToolBatchterminate: 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 循环设计最常见的误读。

源码走查 ​

  1. 在 packages/agent/src/agent.ts 阅读 PendingMessageQueue.drain、steer、followUp;分别把 mode 设想为 all 与默认值,写下每次 drain 的消息数。
  2. 从 runAgentLoop 读到 runLoop;跟踪同一 prompt 如何同时进入 currentContext.messages 与 newMessages。
  3. 在 runLoop 给内外两个 while 标注各自继续条件;验证 follow-up 只在内层退出后读取,steering 在每个 turn 尾部读取。
  4. 逐行读 streamAssistantResponse 的 start、done/error 与循环结束后 result 分支;确认没有 start 时也会把 final message 写入 context。
  5. 搜索 failToolCallsFromTruncatedMessage;核对 length 的 tool call 不会进入 executeToolCalls。
  6. 比较 executeToolCallsSequential 与 executeToolCallsParallel;确认 parallel 用 Promise.all 执行 prepared thunk,却按原 call 顺序追加 tool-result message。