Main Chain B · Step 13

executeToolCalls:barrier 与滚动池

packages/core/agent-loop/src/tool-calls.ts(290 行)调度一个 step 里的全部工具调用:exclusive 调用形成 barrier,parallel-safe 调用走受界滚动池;只有 dispatch/body 并发,策略、结果、结果上下文都按模型顺序提交。

11 step:467 executeToolCalls(...) → tool-calls.ts:59 executeToolCalls → :121 runGroup → (下一步) 14 ctx.tools[TOOL_RUNTIME_SCHEDULER] 分段接口

示例本次示例:bash 调用在调度器里的实况

示例轨迹 13-1 · 单工具调用的 runGroup
# 11 页 467 行 executeToolCalls(ctx, 1, 1, [bash 调用块], signal, acceptContext)
# planned = [{ block, exec:{ callId:'c1', name:'bash', arguments:{command:'ls -la'},
#                            agent, signal } }]
# 88 行:executionMode(bash) = 'parallel'(bash 并发安全)→ group = 全部
# runGroup:startCall(0)
#   168 行:落 tool/call → seq 8(11 页示例轨迹 11-1)
#   169 行:prepare → 14 页 pre-execute → 'dispatch'
#   173 行:dispatch(bash body) 并发执行
#   208 行:commitReady → 顺序提交
#   155 行:appendToolResult → seq 9,sourceEventSeqs:[8]
# 工具结果 content 里有目录列表;concluded=false
# → { concluded:false } → 11 页 471 行返回 null → step 2 继续
来源:13 页行级解读 + 11 页示例轨迹 11-1 的 seq 序列
示例轨迹 13-2 · 模型一次要了 3 个工具时的池
# 模型块:[read(a), bash(b), web_search(c)],全 parallel-safe
# maxParallelToolCalls=2(默认):池容量 2
# startCall(0): tool/call seq 10 → prepare → dispatch(inFlight=1)
# startCall(1): tool/call seq 11 → prepare → dispatch(inFlight=2,池满)
# c 等待;b 先 settle → commitReady 提交 b 的结果(顺序游标等 a)
# a settle → 提交 a → 填池 startCall(2)=c ...
# 结果的 tool/result 顺序 = 模型调用顺序(a,b,c),但执行是并发的
# —— 只有 dispatch/body 重叠;策略、结果、context 全部按序。
来源:13 页 fillPool/commitReady(198-230 行)+ 11 页调度入口

§1executeToolCalls:分组循环

packages/core/agent-loop/src/tool-calls.ts入口与分组59-110
59export async function executeToolCalls(ctx, turn, step, toolCalls, signal, acceptContext) {
67  const agent = ctx.agents.requireInitiator()
68  const { session } = agent
71  const planned: PlannedCall[] = toolCalls.map(block => ({
72    block,
73    exec: {
74      callId: block.id, name: block.name,
76      arguments: parseArguments(block.arguments),
77      agent, signal,
78    },
80  }))
82  let next = 0
84  while (next < planned.length) {
87    const first = planned[next]!
88    const mode = ctx.tools.executionMode(first.exec).kind
89    const group = mode === 'parallel' ? planned.slice(next) : [first]
90    const outcome = await runGroup(ctx, turn, step, group, mode, signal, acceptContext)
93    next += outcome.consumed
94    concluded ||= outcome.concluded
95    if (outcome.aborted) {
96      for (const call of planned.slice(next)) appendSkippedToolCall(session, turn, step, call.block)
97      return { concluded }
98    }
99  }
100  return { concluded }
101}
104function parseArguments(raw: string): unknown {
105  try { return raw ? JSON.parse(raw) : {} } catch { return raw }
106}
67

requireInitiator() 是归因的反面:05 页 wakeDriver 用 withInitiator 设置,这里 require——工具执行需要知道「发起这个 turn 的 agent」。没有 initiator 直接抛错(进程本地归因的消费点)。

71-80

planned:把模型块转成执行输入。注意 exec.signal 是共享的 step 信号——14 页会看到 tools/execute wrapper 可以替换它(但不能解除 caller 取消)。

88-89

分组规则:读当前调用的 executionMode(...)——exclusive 则组 = [这一个](barrier);parallel 则组 = 剩余全部(但 13 页 §2 的 fillPool 会逐个重查模式,中途变 exclusive 会切断池)。

90-98

组串行推进:跑完一组,consumed 推进游标。若组以 aborted 结束,剩余所有调用补合成结果(§3)后整体返回。

104-106

parseArguments:合法 JSON 解析;空串 → {};非法 JSON 保留原文(不抛——坏参数交给工具的 schema 校验去报,调度器不代劳)。

§2runGroup:滚动池与顺序提交

packages/core/agent-loop/src/tool-calls.tsrunGroup 的核心(节选)146-231
146  // `committed` advances only across contiguous model-order slots.
147  const commitReady = async (): Promise<void> => {
148    while (committed < group.length) {
149      const slot = slots[committed]
150      if (slot === undefined) break
151      const call = group[committed]
152      const result = slot.needsPost
153        ? await ctx.tools[TOOL_RUNTIME_SCHEDULER].finalize(slot.exec, slot.result)
154        : ctx.tools[TOOL_RUNTIME_SCHEDULER].finish(slot.exec, slot.result)
156      appendToolResult(session, turn, step, call!.block, result, callSeqs[committed]!)
157      for (const context of result.additionalContexts ?? []) acceptContext(context)
158      concluded ||= result.concludesTurn === true按模型顺序结算
159      committed++
160    }
161  }
163  const inFlight = new Map<number, Promise<number>>()
165  const startCall = async (index: number): Promise<void> => {
167    const call = group[index]!
168    callSeqs[index] = appendToolCall(session, turn, step, call.block)
169    started++
170    const prepared = await ctx.tools[TOOL_RUNTIME_SCHEDULER].prepare(call.exec)
199  const fillPool = async (): Promise<void> => {
200    while (!aborted && nextToStart < group.length && inFlight.size < maxParallelToolCalls) {
203      const nextCall = group[nextToStart]!
204      if (nextToStart > 0 && mode === 'parallel'exclusive 调用形成 barrier
205        && ctx.tools.executionMode(nextCall.exec).kind !== 'parallel') break
206      await startCall(nextToStart)
207      nextToStart++
209      await commitReady()
212      if (signal.aborted) aborted = true
213    }
214  }
219  try {
220    await fillPool()
221    while (inFlight.size > 0) {
222      const settledIndex = await Promise.race(inFlight.values())
223      inFlight.delete(settledIndex)
225      await commitReady()
229      if (signal.aborted) aborted = true
230      await fillPool()
231    }
145-160

顺序提交游标:committed 只跨连续 model 序 slot 前进——slot[committed] 没就绪就 break。慢的第 0 号调用未完成时,第 1 号即使完成也不提交结果。保证「结果的落盘顺序 = 模型的调用顺序」。

151-153

提交时区分 finalize(还要过 post-execute)与 finish(直接终局)——这两个就是 14 页分段接口暴露给调度器的钩子。

155-157

appendToolResult(269-290 行):tool/result 带 sourceEventSeqs: [callSeq] 引用它的 tool/call。additionalContexts 喂给 acceptContext(11 页 415 行 → next-step inbox)。concludesTurn 置位。

167

startCall 第一件事是 先落盘 tool/call(263 行:append 并返回 seq)——已启动的调用在日志里有 call 记录,结果引用它。

169-190

prepare 三态(14 页的分段接口):dispatch(pre 通过,等 body)、post-result(pre 直接给了结果,还要 post)、final-result(终局,连 post 都不走——如 UNKNOWN_TOOL)。

198-213

滚动池:池容量 = maxParallelToolCalls。两个停止条件:aborted、或重新分类——203-204 行对下一个未启动的调用重查 executionMode,若从 parallel 变成 exclusive 就 break,让它成为下一组的 barrier 起点。每次 startCall 后 commitReady——策略与结果在启动间隙就按序提交。

219-230

主驱动:填池 → race 等任一 settle → 顺序提交 → 继续填。调度器失败(231-235 行):停止新派发,排干已启动的,把第一个失败抛给 turn 边界——不伪造结果(与 abort 不同,见 §3)。

§3abort 的合成结果

packages/core/agent-loop/src/tool-calls.tsappendSkippedToolCall249-260
249/** Append the durable call/result pair for a model call skipped after cancellation. */
250function appendSkippedToolCall(session: Session, turn: number, step: number, block: ToolCallBlock): void {
251  const callSeq = appendToolCall(session, turn, step, block)
252  appendToolResult(session, turn, step, block, {
253    content: [{ type: 'text', text: 'Error: tool call aborted before dispatch' }],
254    isError: true,
255    error: {
256      message: 'tool call aborted before dispatch',
257      info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },
258    },
259  }, callSeq)
260}
248-257

为什么 abort 要合成未启动调用的结果?模型请求里出现过 N 个 tool-call,日志就必须有 N 对 call/result——否则 deriveMessages 投影出的 transcript 有一个悬空 tool-call 没有结果,下一次请求就是畸形历史。这是「模型可见 ⟺ 已记录」在取消路径上的强制。code 用 TOOL_ABORTED_BEFORE_DISPATCH(区别于 body 已启动后取消的 TOOL_ABORTED,14 页)。

250

注意:被跳过的调用也先落 tool/call 再落 tool/result——call/result 配对在日志里永远成对。