Skip to content

并行结果排序

追踪工具 future 的即时启动、FuturesOrdered 收集、stream 完成后的有序 drain,以及失败和取消如何影响历史输出顺序。

基于rust-v0.150.0
CodexRustToolsConcurrency

并行结果排序 ​

工具可以同时执行,不代表结果会按完成时间写入模型历史。当前 Turn loop 在看到每个 tool call 时立即创建执行 future,ToolCallRuntime 内部又马上 tokio::spawn dispatch task;但这些 future 被放入 FuturesOrdered,直到 当前 response stream 收束后才按插入顺序 drain。于是“执行开始顺序”“完成顺序”和“历史写入顺序”是三条不同 时间线。

本文面向已经读过工具并行判定和 ToolOrchestrator执行流程的读者。前文解释 read/write gate,本文只研究 future 收集、join、错误隔离和输出排序,不展开 runtime 是否 parallel-safe 的具体判定。读完后,读者应能解释为什么 后完成的工具可能先完成却后写入历史,以及为什么模型下一次请求会看到“所有 calls 在前、所有 outputs 在后”。

1. Future诞生 ​

1.1 Stream回调 ​

handle_output_item_done 对每个本地 ToolCall 创建一个 InFlightFuture。函数先记录模型发出的 call item,再将 ToolCallRuntime::handle_tool_call 包装为 boxed future,返回给 Turn loop。

相关源码:

  • codex-rs/core/src/stream_events_utils.rs :: InFlightFuture
  • codex-rs/core/src/stream_events_utils.rs :: handle_output_item_done
rust
pub(crate) type InFlightFuture<'f> =
    Pin<Box<dyn Future<Output = Result<ResponseInputItem>> + Send + 'f>>;

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);

这里记录的 call item 与 tool output 分开:call 在 stream item 完成时立即持久化,output 要等 future drain。因此即使 工具执行非常快,历史中仍先存在原始 call。

1.2 Dispatch立即启动 ​

handle_tool_call 返回 future,但它调用的 handle_tool_call_with_source 在构造返回 future 之前已经执行 tokio::spawn。因此将 future 放进容器并不意味着工具要等到 drain 才开始;dispatch task 已经由 Tokio 调度。

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

rust
let mut dispatch_handle: AbortOnDropHandle<Result<AnyToolResult, FunctionCallError>> =
    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
    }));

async move {
    tokio::select! {
        res = &mut dispatch_handle => res.map_err(Self::tool_task_join_error)?,
        _ = cancellation_token.cancelled() => { /* cancellation path */ },
    }
}

2. 有序容器 ​

2.1 FuturesOrdered ​

每次 sampling request 创建一个新的 FuturesOrdered<BoxFuture<...>>。tool call item 按 stream 到达顺序调用 push_back,这个顺序成为之后的输出顺序锚点。

源码位置:codex-rs/core/src/session/turn.rs :: run_sampling_request

rust
let mut in_flight: FuturesOrdered<
    BoxFuture<'static, CodexResult<ResponseInputItem>>,
> = FuturesOrdered::new();

// 每个 OutputItemDone 完成后:
if let Some(tool_future) = output_result.tool_future {
    in_flight.push_back(tool_future);
}

FuturesOrdered 会推进多个 future,但 next() 只按 insertion order 交付结果。后加入的 future 即使先完成,也会 在容器内部等待前面的 future 可交付;这提供历史顺序稳定性,但会带来 head-of-line blocking。

2.2 三种顺序 ​

顺序决定者是否稳定
Call 出现顺序Responses stream 的 OutputItemDone是,按 stream item
Handler 完成顺序readiness、RwLock、I/O 和 runtime否,可交错
History output 顺序FuturesOrdered::next()是,按 push_back

3. Drain屏障 ​

3.1 Stream先收束 ​

Turn loop 在 Responses stream 的 Completed、错误或关闭分支结束后,才调用 drain_in_flight。不过 spawned tasks 可能早已开始甚至完成;这里的“之后”指结果写入屏障,不是执行启动屏障。

源码位置:codex-rs/core/src/session/turn.rs :: run_sampling_request

rust
drop(sampling_timing_guard);

let tool_blocking_timing_guard = if in_flight.is_empty() {
    None
} else {
    Some(turn_context.turn_timing_state.begin_tool_blocking())
};

drain_in_flight(&mut in_flight, sess.clone(), turn_context.clone()).await?;
drop(tool_blocking_timing_guard);

这保证当前 response 中所有 call items 已记录完毕后,再集中写入 tool outputs。tool_blocking timing 统计的是 stream 收束后仍需等待的工具时间;已经提前完成的 future 不会贡献额外 handler 执行时间,但仍按顺序被 drain。

3.2 有序写回 ​

drain_in_flight 循环调用 next().await,把每个 ResponseInputItem 转成 ResponseItem 并逐个记录。写入动作本身 是顺序 await,因此 conversation history 的 output 顺序与 FuturesOrdered 交付顺序一致。

源码位置:codex-rs/core/src/session/turn.rs :: drain_in_flight

rust
async fn drain_in_flight(
    in_flight: &mut FuturesOrdered<BoxFuture<'static, CodexResult<ResponseInputItem>>>,
    sess: Arc<Session>,
    turn_context: Arc<TurnContext>,
) -> CodexResult<()> {
    while let Some(res) = in_flight.next().await {
        match res {
            Ok(response_input) => {
                let response_item = response_input.into();
                sess.record_conversation_items(
                    &turn_context,
                    std::slice::from_ref(&response_item),
                )
                .await;
                mark_thread_memory_mode_polluted_if_external_context(
                    sess.as_ref(),
                    turn_context.as_ref(),
                    &response_item,
                )
                .await;
            }
            Err(err) => {
                error_or_panic(format!(
                    "in-flight tool future failed during drain: {err}"
                ));
            }
        }
    }
    Ok(())
}

4. 分组不变量 ​

4.1 Calls先于Outputs ​

模型 call item 在收到时立即记录,outputs 在 stream 结束后 drain。结果是下一次模型请求中,同一 sampling response 的所有 function calls 位于所有 function call outputs 之前。

4.2 CallId配对 ​

outputs 不按 handler 完成顺序排列,而是与 calls 的出现顺序 zip 对齐。call_id 是配对依据;不同工具输出内容 再相似,也不能依赖数组位置之外的猜测。

源码位置:codex-rs/core/tests/suite/tool_parallelism.rs :: tool_results_grouped

rust
assert_eq!(function_calls.len(), 3);
assert_eq!(function_call_outputs.len(), 3);

for (index, _) in &function_calls {
    for (output_index, _) in &function_call_outputs {
        assert!(
            *index < *output_index,
            "all function calls must come before outputs"
        );
    }
}

for (call, output) in function_calls.iter().zip(function_call_outputs.iter()) {
    assert_eq!(
        call.1.get("call_id").and_then(Value::as_str),
        output.1.get("call_id").and_then(Value::as_str),
    );
}

5. 错误隔离 ​

5.1 可回复错误 ​

每个 ToolCallRuntime::handle_tool_call 会把 RespondToModel 转成对应 payload 的失败 ResponseInputItem。因此一个 普通工具失败仍占据自己在 FuturesOrdered 中的位置,后面的结果不会因为它失败而丢失顺序。

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

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.2 Fatal与Join错误 ​

Fatal 和 task join error 会让该 in-flight future 返回 Err。drain_in_flight 对 Err 调用 error_or_panic:debug 构建会 panic,release 构建记录 error 后继续 drain。它不是普通模型可见失败,因此不能声称所有工具错误都被隔离 成独立 output。

源码位置:codex-rs/core/src/util.rs :: error_or_panic

rust
pub(crate) fn error_or_panic(message: impl ToString) {
    if cfg!(debug_assertions) {
        panic!("{}", message.to_string());
    } else {
        error!("{}", message.to_string());
    }
}

5.3 Head-of-line阻塞 ​

如果 call A 很慢、call B 很快,B 的 handler 可以先完成,但 FuturesOrdered 不会先 yield B。history 写入必须等 A 先交付。这是稳定顺序的成本;它不阻止 B 执行,只延迟 B 的结果可见时间。

6. 取消传播 ​

每个 tool future 使用从 Turn cancellation token 派生的 child token。父 token 取消会传播到所有 pending calls; 每个 ToolCallRuntime 根据自身 cleanup policy 产生正常完成、aborted output 或 Fatal join error。FuturesOrdered 仍按 call 顺序 drain,所以前面需要长 teardown 的调用会延迟后面已经取消完成的结果写入。

取消后 Turn loop 仍先 drain in-flight futures,再检查父 token 并返回 TurnAborted。这样已经产生的工具终态可以先 完成 lifecycle/记录处理,但调用方最终仍知道 Turn 被取消。

7. 流式并发 ​

shell_tools_start_before_response_completed_when_stream_delayed 使用分块 SSE:第一块发出四个 shell calls,第二块 的 response.completed 被 gate 阻塞。测试在释放 completed gate 前轮询文件,确认四个命令都已写入时间戳。这直接 证明 dispatch task 在 stream 完成前启动。

源码位置:codex-rs/core/tests/suite/tool_parallelism.rs :: shell_tools_start_before_response_completed_when_stream_delayed

测试随后释放 completed gate,等待 TurnComplete,并断言所有命令时间戳早于或等于 streaming response completed 时间。它证明执行启动与结果 drain 是分离的,不证明 outputs 在 stream 完成前写入模型历史。

8. 测试路径 ​

8.1 分组与顺序 ​

tool_results_grouped 证明 calls 全部在 outputs 之前,并按 call_id 顺序配对。它不要求 handler 按同一顺序完成, 验证的是下一次模型请求的 history shape。

8.2 提前启动 ​

shell_tools_start_before_response_completed_when_stream_delayed 证明工具在 delayed response.completed 之前已经执行。 read_file_tools_run_in_parallel 和 mixed_parallel_tools_run_in_parallel 则用 barrier/耗时证明多个 spawned tasks 可重叠。

相关测试:

  • codex-rs/core/tests/suite/tool_parallelism.rs :: tool_results_grouped
  • codex-rs/core/tests/suite/tool_parallelism.rs :: shell_tools_start_before_response_completed_when_stream_delayed
  • codex-rs/core/tests/suite/tool_parallelism.rs :: read_file_tools_run_in_parallel
  • codex-rs/core/tests/suite/tool_parallelism.rs :: mixed_parallel_tools_run_in_parallel

这些测试证明执行时序和 history 分组,不证明所有失败/取消组合、provider 顺序或跨 Turn 结果排序。

9. 阅读练习 ​

在 Codex 源码 workspace 中运行:

bash
cargo test -p codex-core tool_results_grouped
cargo test -p codex-core shell_tools_start_before_response_completed_when_stream_delayed
cargo test -p codex-core read_file_tools_run_in_parallel
cargo test -p codex-core mixed_parallel_tools_run_in_parallel

然后尝试回答:

  1. 工具 future 为什么能在 drain_in_flight 调用前执行?指出 tokio::spawn 的发生位置。
  2. call B 先完成但排在 call A 后面时,什么时候 B 的 output 才会写入 history?
  3. 为什么所有 call items 会位于所有 outputs 之前?这由哪两个记录时机共同保证?
  4. 普通 RespondToModel 与 Fatal 在 FuturesOrdered 中各占据怎样的结果槽位?
  5. Turn 取消后为什么仍要 drain in-flight futures,再返回 TurnAborted?

10. 边界 ​

本篇解释的是单次 sampling response 中的 in-flight future 和 history 顺序,不负责:

  • runtime 是否拿 read lock 或 write lock;
  • Registry lifecycle、Hook 和 handler 内部副作用;
  • 跨 Turn、resume/fork 或 rollout reconstruction 的长期排序;
  • provider 是否按调用顺序发送 OutputItemDone。

排查“工具已经执行但结果迟迟没进下一次请求”时,应区分 spawned task 状态、FuturesOrdered 前序 future、response stream 是否结束和 history drain 进度。稳定输出顺序来自有序交付与延后写回,不代表实际执行被串行化。