并行结果排序
工具可以同时执行,不代表结果会按完成时间写入模型历史。当前 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 :: InFlightFuturecodex-rs/core/src/stream_events_utils.rs :: handle_output_item_done
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
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
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
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
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
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
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
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_groupedcodex-rs/core/tests/suite/tool_parallelism.rs :: shell_tools_start_before_response_completed_when_stream_delayedcodex-rs/core/tests/suite/tool_parallelism.rs :: read_file_tools_run_in_parallelcodex-rs/core/tests/suite/tool_parallelism.rs :: mixed_parallel_tools_run_in_parallel
这些测试证明执行时序和 history 分组,不证明所有失败/取消组合、provider 顺序或跨 Turn 结果排序。
9. 阅读练习
在 Codex 源码 workspace 中运行:
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然后尝试回答:
- 工具 future 为什么能在
drain_in_flight调用前执行?指出tokio::spawn的发生位置。 - call B 先完成但排在 call A 后面时,什么时候 B 的 output 才会写入 history?
- 为什么所有 call items 会位于所有 outputs 之前?这由哪两个记录时机共同保证?
- 普通
RespondToModel与 Fatal 在 FuturesOrdered 中各占据怎样的结果槽位? - 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 进度。稳定输出顺序来自有序交付与延后写回,不代表实际执行被串行化。
