Skip to content

ToolOrchestrator执行流程

基于当前 ToolCallRuntime 与 Registry 源码,追踪工具调用的 admission、并行锁、Hook、handler、取消、lifecycle 和结果回灌。

基于rust-v0.150.0
CodexRustToolsRuntime

ToolOrchestrator执行流程 ​

archive 中曾经存在一个名为 ToolOrchestrator 的单体设计,但当前版本的编排责任已经拆开: stream_events_utils 接收模型响应并排队 tool future,ToolCallRuntime 负责 readiness、并行 gate、取消和结果 适配,ToolRegistry 负责 lookup、Hook、lifecycle 与 handler 调用。本文使用“编排流程”描述这组真实协作,不把 旧的单体类型当作当前源码事实。

本文面向已经读过ToolRouter解析与分派、ToolOutput与错误模型 和ToolRegistry数据结构的读者。本文回答“一次本地工具调用如何从模型响应走到 结果或取消终态”,不展开 approval/sandbox 的具体策略,也不重复 schema 和 handler 业务。读完后,读者应能定位 dispatch waiting、handler execution、Hook feedback、取消 teardown 和 terminal outcome 的所有者。

1. 编排入口 ​

1.1 响应排队 ​

模型流收到一个可本地分派的响应项后,handle_output_item_done 记录原始 response item,创建 child cancellation token,并把 ToolCallRuntime::handle_tool_call 放入 in-flight future。这里还没有执行 handler,只有把调用交给 后续运行时。

源码位置:codex-rs/core/src/stream_events_utils.rs :: handle_output_item_done

rust
match ToolRouter::build_tool_call(item.clone()) {
    Ok(Some(call)) => {
        let payload_preview = tool_log_payload(&call.payload, &call.direct_source());
        tracing::info!("ToolCall: {} {}", call.tool_name, payload_preview);

        record_completed_response_item(ctx.sess.as_ref(), ctx.turn_context.as_ref(), &item)
            .await;

        let cancellation_token = ctx.cancellation_token.child_token();
        let tool_future: InFlightFuture<'static> = Box::pin(
            ctx.tool_runtime
                .clone()
                .handle_tool_call(call, cancellation_token),
        );
        output.needs_follow_up = true;
        output.tool_future = Some(tool_future);
    }
    Ok(None) => { /* 其他 stream item 走普通消费者 */ }
    Err(error) => return Err(error.into()),
}

record_completed_response_item 记录的是模型发出的调用,不是执行成功证明;真正的执行结果要等 future 完成后 再由 Turn loop 消费。child token 让单个工具可以被当前 Turn 取消,而不会创建新的全局取消源。

1.2 Runtime快照 ​

ToolCallRuntime 保存 session、产生工具表的 StepContext、共享 diff tracker 和并行执行锁。它不从全局重新 读取 Router,因为工具调用可能在模型响应之后才真正 admission。

源码位置:codex-rs/core/src/tools/parallel.rs :: ToolCallRuntime

rust
pub(crate) struct ToolCallRuntime {
    session: Arc<Session>,
    // Tool calls may run later, so retain the step whose tool list advertised them.
    step_context: Arc<StepContext>,
    tracker: SharedTurnDiffTracker,
    parallel_execution: Arc<RwLock<()>>,
}

2. Admission阶段 ​

2.1 能力查询 ​

handle_tool_call_with_source 一开始从 Router 查询三件事:工具是否支持并行、runtime 是否存在、取消时是否需要 等待 runtime teardown。它随后建立计时 guard、terminal outcome flag 和异步 dispatch task。

源码位置:codex-rs/core/src/tools/parallel.rs :: ToolCallRuntime::handle_tool_call_with_source

rust
let router = &self.step_context.tool_router;
let supports_parallel = router.tool_supports_parallel(&call);
let tool_runtime = router.tool_runtime(&call);
let wait_for_runtime_cancellation = router.tool_waits_for_runtime_cancellation(&call);
let terminal_outcome_reached = Arc::new(AtomicBool::new(false));

没有 runtime 不会在这里立刻 panic;Registry dispatch 会返回 RespondToModel。parallel 和 cancellation metadata 在 admission 前读取,保证后续 gate 决策与调用所属 Step 一致。

2.2 Readiness ​

异步 task 先等待 runtime readiness,再竞争共享 RwLock。MCP runtime 可以等待 server ready;没有 readiness 的 runtime 直接继续。read lock 代表允许并行的调用,write lock 保护不支持并行的调用。

源码位置:codex-rs/core/src/tools/parallel.rs :: handle_tool_call_with_source

rust
let mut dispatch_handle = AbortOnDropHandle::new(tokio::spawn(async move {
    if let Some(tool_runtime) = tool_runtime
        && let Some(readiness) = tool_runtime.wait_until_ready(&session)
    {
        readiness.await;
    }

    let _guard = if supports_parallel {
        Either::Left(lock.read().await)
    } else {
        Either::Right(lock.write().await)
    };

    router
        .dispatch_tool_call_with_terminal_outcome(
            session,
            step_context,
            invocation_cancellation_token,
            tracker,
            dispatch_call,
            source,
            dispatch_terminal_outcome_reached,
        )
        .await
}));

2.3 Admission时刻 ​

进入读写锁之后,源码把 execution_started_at 写入 OnceLock。计时因此区分“等待 readiness/锁的 dispatch duration” 和真正 handler execution duration;取消发生在 gate 前时,不会被报告成 handler 已启动。

3. Registry阶段 ​

3.1 Invocation组装 ​

Router dispatch 将 ToolCall 消费为 ToolInvocation;Registry 接到的已经是带 session、turn、StepContext、source、 cancel token 和 tracker 的执行上下文。

源码位置:codex-rs/core/src/tools/router.rs :: dispatch_tool_call_with_code_mode_result_inner

rust
let ToolCall {
    tool_name,
    call_id,
    payload,
    ..
} = call;

let invocation = ToolInvocation {
    session,
    turn: Arc::clone(&step_context.turn),
    step_context,
    cancellation_token,
    tracker,
    call_id,
    tool_name,
    source,
    payload,
};

self.registry
    .dispatch_any_with_terminal_outcome(invocation, terminal_outcome_reached)
    .await

3.2 前置顺序 ​

Registry dispatch 先增加 active-turn tool call 计数,再 lookup、kind check、lifecycle start 和 PreToolUse hook; handler 只有在这些步骤通过后才被调用。

源码位置:codex-rs/core/src/tools/registry.rs :: dispatch_any_with_terminal_outcome

rust
let tool = match self.tool(&tool_name) {
    Some(tool) => tool,
    None => {
        let message = unsupported_tool_call_message(&invocation.payload, &tool_name);
        return Err(FunctionCallError::RespondToModel(message));
    }
};

if !tool.matches_kind(&invocation.payload) {
    let message = format!("tool {tool_name} invoked with incompatible payload");
    return Err(FunctionCallError::Fatal(message));
}

notify_tool_start(&invocation).await;

3.3 Hook反馈 ​

PreToolUse 可以 block 或改写 invocation;handler 成功后才可能创建 PostToolUse payload。PostToolUse block 不会 撤销已经完成的 handler,只会让最终模型投影变成 feedback。

源码位置:codex-rs/core/src/tools/registry.rs :: dispatch_any_with_terminal_outcome

rust
if let Some(pre_tool_use_payload) = tool.pre_tool_use_payload(&invocation) {
    match run_pre_tool_use_hooks(/* ... */).await {
        PreToolUseHookResult::Blocked(message) => {
            return Err(FunctionCallError::RespondToModel(message));
        }
        PreToolUseHookResult::Continue {
            updated_input: Some(updated_input),
        } => {
            invocation = tool.with_updated_hook_input(invocation, updated_input)?;
        }
        PreToolUseHookResult::Continue {
            updated_input: None,
        } => {}
    }
}

4. 结果与终态 ​

4.1 正常结果 ​

handler 返回 ToolOutput 后,Registry 先读取 preview 和 success,构造 post-use payload,运行 PostToolUse, 最后记录完成 trace 并返回 AnyToolResult。Direct 调用的外层再把它转换为 ResponseInputItem。

相关源码:

  • codex-rs/core/src/tools/registry.rs :: handle_any_tool
  • codex-rs/core/src/tools/registry.rs :: ToolRegistry::dispatch_any_with_terminal_outcome
rust
let output = tool.handle(invocation.clone()).await?;
let post_tool_use_payload =
    CoreToolRuntime::post_tool_use_payload(tool, &invocation, output.as_ref());
Ok(AnyToolResult {
    call_id,
    payload,
    result: output,
    post_tool_use_payload,
})

4.2 取消竞争 ​

外层使用 tokio::select! 同时等待 dispatch task 和 cancellation token。如果 task 已完成或 terminal outcome 已经 被 Registry 设置,取消分支会等待已有结果,避免重复发 aborted lifecycle。否则根据 runtime policy:需要 teardown 的 runtime 等待任务收尾,不需要的直接 abort,再合成 AbortedToolOutput 和 aborted lifecycle。

源码位置:codex-rs/core/src/tools/parallel.rs :: handle_tool_call_with_source

rust
tokio::select! {
    res = &mut dispatch_handle => {
        res.map_err(Self::tool_task_join_error)?
    }
    _ = cancellation_token.cancelled() => {
        if terminal_outcome_reached.load(Ordering::Acquire)
            || dispatch_handle.is_finished()
        {
            dispatch_handle.await.map_err(Self::tool_task_join_error)?
        } else if wait_for_runtime_cancellation {
            dispatch_handle.await.map_err(Self::tool_task_join_error)?;
            Self::aborted_response(&call, started.elapsed().as_secs_f32().max(0.1))
        } else {
            dispatch_handle.abort();
            Self::aborted_response(&call, started.elapsed().as_secs_f32().max(0.1))
        }
    }
}

上面省略了完整的 notify_tool_aborted 和 join error 分支,但保留了真实选择条件。取消路径的核心不变量是: 一个调用只能有一个 terminal lifecycle outcome。

4.3 错误适配 ​

FunctionCallError::RespondToModel 会转成 payload 对应的失败 response,Fatal 则升级为 CodexErr::Fatal。因此 “handler 返回错误”不是一个统一状态:是否还能让模型修复,取决于 error variant。

相关源码:

  • codex-rs/core/src/tools/parallel.rs :: ToolCallRuntime::handle_tool_call
  • codex-rs/core/src/tools/parallel.rs :: ToolCallRuntime::failure_response
rust
match future.await {
    Ok(response) => Ok(response.into_response()),
    Err(FunctionCallError::Fatal(message)) => Err(CodexErr::Fatal(message)),
    Err(other) => Ok(Self::failure_response(error_call, other)),
}

5. 终态竞争 ​

5.1 Terminal flag ​

terminal_outcome_reached 是跨 Registry 和取消分支共享的 AtomicBool。Registry 在 lifecycle finish 前通过 notify_tool_finish_if_unclaimed 竞争设置它;取消分支也会先检查或设置它。谁先获得 terminal outcome,谁负责 发出 finish/aborted lifecycle,另一方只能复用结果或退出。

相关源码:

  • codex-rs/core/src/tools/parallel.rs :: terminal_outcome_reached
  • codex-rs/core/src/tools/parallel.rs :: notify_tool_finish_if_unclaimed
rust
async fn notify_tool_finish_if_unclaimed(
    invocation: &ToolInvocation,
    terminal_outcome_reached: Option<&AtomicBool>,
    outcome: ToolCallOutcome,
) -> bool {
    if terminal_outcome_reached.is_some_and(|reached| reached.swap(true, Ordering::AcqRel)) {
        return false;
    }

    notify_tool_finish(invocation, outcome).await;
    true
}

这个竞争关系解释了为什么 lifecycle finish 和 aborted 不会各自独立发送:同一个原子标记决定唯一 terminal owner。

5.2 Timing ​

ToolCallTimingGuard 只为 Direct/DirectPlaintextMessage 创建;Code Mode 嵌套调用不创建独立 timing event,因为外层 direct Code Mode 调用已经覆盖它们。计时 guard 的 Drop 同时计算 total、dispatch 和 handler duration。

相关源码:

  • codex-rs/core/src/tools/parallel.rs :: ToolCallTimingGuard::capture
  • codex-rs/core/src/tools/parallel.rs :: Drop for ToolCallTimingGuard

6. 测试路径 ​

6.1 取消测试 ​

cancellation_before_dispatch_admission_logs_dispatch_only_timing 持有 execution gate 后取消调用,断言 response 仍能返回,并且 timing 只报告 dispatch duration,没有 execution start 或 handler duration。

cancellation_after_handler_finishes_preserves_completed_lifecycle 让 handler 已完成但 lifecycle finish 尚未结束, 此时取消不会把已完成调用改写成 aborted lifecycle。

相关测试:

  • codex-rs/core/src/tools/parallel.rs :: cancellation_before_dispatch_admission_logs_dispatch_only_timing
  • codex-rs/core/src/tools/parallel.rs :: cancellation_after_handler_finishes_preserves_completed_lifecycle

6.2 生命周期测试 ​

dispatch_uses_canonical_tool_names_for_lifecycle_contributors 验证成功但 success: false 的 handler 仍记录 Completed outcome,而真正返回 error 的 handler 记录 Failed;这说明 lifecycle outcome 不等同于 Rust Result 的简单 布尔值。

相关测试:

  • codex-rs/core/src/tools/registry_tests.rs :: dispatch_uses_canonical_tool_names_for_lifecycle_contributors
  • codex-rs/core/src/tools/registry_tests.rs :: post_tool_use_feedback_output_keeps_code_mode_result_typed

6.3 Timing测试 ​

tool_call_timing_guard_ignores_code_mode_source 断言 Direct 创建 timing guard,Code Mode source 不创建重叠事件。 cancellation_after_handler_finishes_preserves_completed_lifecycle 验证 handler 已完成时取消不会把完成结果改写成 aborted lifecycle。

这些测试证明取消竞争、lifecycle ownership 和 timing 统计边界,不证明进程在 kill -9 后的清理或所有平台 runtime。

7. 阅读练习 ​

在 Codex 源码 workspace 中运行:

bash
cargo test -p codex-core cancellation_before_dispatch_admission_logs_dispatch_only_timing
cargo test -p codex-core cancellation_after_handler_finishes_preserves_completed_lifecycle
cargo test -p codex-core dispatch_uses_canonical_tool_names_for_lifecycle_contributors
cargo test -p codex-core tool_call_timing_guard_ignores_code_mode_source
cargo test -p codex-core cancellation_after_handler_finishes_preserves_completed_lifecycle

然后尝试回答:

  1. 一个工具在 readiness 等待期间被取消,为什么不能报告 handler 已经执行?
  2. waits_for_runtime_cancellation 为 true 与 false 时,取消分支分别由谁拥有 teardown?
  3. PostToolUse block 为什么可以改变模型反馈,却不能把 lifecycle outcome 简单改成“未执行”?
  4. terminal_outcome_reached 为什么必须是 AtomicBool,而不是普通局部布尔变量?
  5. Code Mode 为什么不创建独立 timing event?这对诊断嵌套调用有什么影响?

8. 边界 ​

当前编排流程负责 admission、执行 gate、Registry dispatch、取消竞争和结果适配,不负责:

  • 工具是否注册以及 exposure 如何计算;
  • 具体 approval/sandbox 策略和 OS 进程实现;
  • handler 业务参数、MCP 协议和 patch grammar;
  • provider stream 的重试、模型采样和最终历史压缩。

排查“工具卡住、取消后仍有进程、重复 finish event 或模型收到错误反馈”时,应按 stream queue → readiness → parallel lock → Registry → Hook → handler → terminal flag → response 顺序定位, 不要把所有等待都归类为 handler 执行,也不要把一次 aborted response 当作 runtime 已完成 teardown 的充分证明。