Unified Exec轮询与等待
Unified Exec 的等待不是单一 sleep(yield_time_ms)。每次轮询都要在输出 buffer、输出通知、退出 token、关闭通知、暂停状态和 deadline 之间选择先发生者;进程退出后还要给 reader 一小段关闭窗口,避免最后的 stdout 被遗漏。
本文承接Unified Exec写入stdin和Unified Exec会话状态机,面向理解 Tokio select!、Notify、watch 和异步 mutex 的读者。范围是 collect_output_until_deadline 及其调用者,不重复输入分支、状态 store 或 head-tail 算法。读完后,你应能解释空 poll 为什么等待更久、暂停为何延长 deadline、退出后为何还有最多 50ms 的关闭等待。
1. 轮询输入
1.1 输出句柄
轮询函数只接收共享 OutputHandles,不直接接触 child handle。它从 buffer drain 当前批次,并使用通知对象等待下一批数据。
源码位置:codex-rs/core/src/unified_exec/process.rs :: OutputHandles
#[derive(Clone)]
pub(crate) struct OutputHandles<const MAX_BYTES: usize = UNIFIED_EXEC_OUTPUT_MAX_BYTES> {
pub(crate) output_buffer: Arc<Mutex<HeadTailBuffer<MAX_BYTES>>>,
pub(crate) output_notify: Arc<Notify>,
pub(crate) output_closed: Arc<AtomicBool>,
pub(crate) output_closed_notify: Arc<Notify>,
pub(crate) cancellation_token: CancellationToken,
}buffer 保存数据,output_notify 唤醒“有新数据”的等待者,output_closed 发布 reader 已关闭,output_closed_notify 让退出后的等待可以立即结束,取消 token 则表示进程生命周期已结束。它们分别表达数据、唤醒、关闭和生命周期,不能互换。
1.2 drain优先
collect_output_until_deadline 每轮先取得 buffer mutex 并 drain,再决定是否等待。这样通知在检查前到达也不会导致丢数据;数据已经存在时不会无意义睡眠。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: collect_output_until_deadline
let mut collected = HeadTailBuffer::default();
let mut exit_signal_received = cancellation_token.is_cancelled();
let mut post_exit_deadline: Option<Instant> = None;
loop {
Self::extend_deadlines_while_paused(
&mut pause_state,
&mut deadline,
&mut post_exit_deadline,
)
.await;
let drained_output: HeadTailBuffer<MAX_BYTES>;
let has_drained_output: bool;
let mut wait_for_output = None;
{
let mut guard = output_buffer.lock().await;
drained_output = std::mem::take(&mut *guard);
has_drained_output =
drained_output.retained_bytes() > 0 || drained_output.omitted_bytes() > 0;
if !has_drained_output {
wait_for_output = Some(output_notify.notified());
}
}通知 future 在释放 mutex 前创建,避免“检查为空后、创建等待前”错过通知的竞态;释放锁后才进入 select!,所以 reader 可以继续写入。
std::mem::take 用同类型的空 HeadTailBuffer<MAX_BYTES> 替换共享 buffer,因此 producer 能立即继续写入。被取出的 buffer 随后通过 collected.push_buffer(drained_output) 合并到本轮结果;合并不是简单拼接,而是继续使用同一个 head/tail 容量预算,并累加已经发生的 omission。
1.3 多轮容量预算
一次 poll 可能经历多次共享 buffer drain。push_buffer 必须让这些分段表现得像连续写入同一个缓冲:保留最早的 head、最新的 tail,并把源 buffer 已记录的 omission 加到目标中。
源码位置:codex-rs/core/src/unified_exec/head_tail_buffer.rs :: HeadTailBuffer::push_buffer
pub(crate) fn push_buffer(&mut self, buffer: Self) {
let Self {
head,
tail,
omitted_bytes,
} = buffer;
self.omitted_bytes = self.omitted_bytes.saturating_add(omitted_bytes);
let overflow = if self.head.is_empty() {
self.head = head;
&[]
} else {
self.fill_head(&head)
};
if tail.len() == Self::TAIL_BUDGET {
self.omitted_bytes = self
.omitted_bytes
.saturating_add(self.tail.len())
.saturating_add(overflow.len());
self.tail = tail;
} else {
self.push_tail(overflow);
if self.tail.is_empty() {
self.tail = tail;
} else {
let (first, second) = tail.as_slices();
self.push_tail(first);
self.push_tail(second);
}
}
}这也是 original_token_count 和 omission marker 能跨多次唤醒保持一致的基础。如果每次 drain 都重置容量预算,长时间后台命令可以通过频繁输出绕过 1 MiB 上限。
2. 截止时间
2.1 普通等待
没有输出且进程尚未退出时,等待者同时监听新输出、退出 token、deadline 和暂停状态。任何一个分支都能重新进入循环或返回。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: collect_output_until_deadline
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining == Duration::ZERO {
break;
}
let notified = wait_for_output.unwrap_or_else(|| output_notify.notified());
tokio::pin!(notified);
let exit_notified = cancellation_token.cancelled();
tokio::pin!(exit_notified);
tokio::select! {
_ = &mut notified => {}
_ = &mut exit_notified => exit_signal_received = true,
_ = tokio::time::sleep(remaining) => break,
_ = Self::wait_for_pause_change(pause_state.as_ref()) => {}
}deadline 是本次调用的等待边界,不是进程总寿命。首轮 exec_command、非空写入和空 poll 会先用不同策略计算它,然后共享这一收集循环。
2.2 等待上限
write_stdin 对空输入和非空输入采用不同上限:空 poll 至少等待 MIN_EMPTY_YIELD_TIME_MS,上限由 manager 配置;非空写入最多 MAX_YIELD_TIME_MS。这样后台输出可以长轮询,交互输入不会被任意大的模型参数拖住。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
let yield_time_ms = {
let time_ms = request.yield_time_ms.max(MIN_YIELD_TIME_MS);
if request.input.is_empty() {
time_ms.clamp(MIN_EMPTY_YIELD_TIME_MS, self.max_write_stdin_yield_time_ms)
} else {
time_ms.min(MAX_YIELD_TIME_MS)
}
};
let deadline = Instant::now() + Duration::from_millis(yield_time_ms);
let collected_output =
Self::collect_output_until_deadline(&output, pause_state, deadline).await;这段代码只决定“本次工具响应等多久”,不会杀死超过等待时间的 child;进程仍活跃时响应会带 process_id,后续可以继续 poll。
2.3 首次等待
首次 exec_command 不使用空 poll 的 5 秒下限,而是调用 clamp_yield_time。所有平台都限制在 250ms 到 30s;Windows 还把低于 10s 的请求提升到 10s,为进程创建和 ConPTY 初始化留出更稳定的首轮窗口。
源码位置:codex-rs/core/src/unified_exec/mod.rs :: clamp_yield_time
pub(crate) fn clamp_yield_time(yield_time_ms: u64) -> u64 {
let yield_time_ms = if cfg!(windows) {
yield_time_ms.max(WINDOWS_INITIAL_EXEC_YIELD_TIME_FLOOR_MS)
} else {
yield_time_ms
};
yield_time_ms.clamp(MIN_YIELD_TIME_MS, MAX_YIELD_TIME_MS)
}因此三个窗口应分开理解:首次 exec 使用 clamp_yield_time,非空 stdin 写入最多 30s,空 poll 至少 5s 且可提高到 Session 配置的后台终端上限。它们都只控制工具调用何时返回,不控制进程寿命。
3. 暂停语义
3.1 延长deadline
暂停状态由 watch::Receiver<bool> 提供。若进入轮询时已暂停,函数等待状态变为非暂停,并把暂停持续时间加回普通 deadline 和退出后的 deadline。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: extend_deadlines_while_paused
let Some(receiver) = pause_state.as_mut() else {
return;
};
if !*receiver.borrow() {
return;
}
let paused_at = Instant::now();
while *receiver.borrow() {
if receiver.changed().await.is_err() {
break;
}
}
let paused_for = paused_at.elapsed();
*deadline += paused_for;
if let Some(post_exit_deadline) = post_exit_deadline.as_mut() {
*post_exit_deadline += paused_for;
}暂停不是取消,也不是把 deadline 设为无限;它只把暂停期间从“有效等待时间”中扣除。watch channel 关闭时循环结束,调用仍会继续按当前 deadline 收尾。
3.2 等待状态变化
暂停发生在普通 select! 等待期间时,wait_for_pause_change 只负责唤醒循环;下一轮再由 extend_deadlines_while_paused 统一修正时间。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: wait_for_pause_change
async fn wait_for_pause_change(pause_state: Option<&watch::Receiver<bool>>) {
match pause_state {
Some(pause_state) => {
let mut receiver = pause_state.clone();
let _ = receiver.changed().await;
}
None => std::future::pending::<()>().await,
}
}没有 pause receiver 时使用 pending future,等价于从 select! 中移除这一事件源,而不是创建一个不断唤醒的空任务。
4. 退出后窗口
4.1 退出信号
进程退出时 cancellation token 先被取消,但 output reader 可能还有尾部数据。轮询检测到 exit_signal_received 后不再使用完整 deadline,而是建立最多 50ms 的 post_exit_deadline。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: collect_output_until_deadline
const POST_EXIT_CLOSE_WAIT_CAP: Duration = Duration::from_millis(50);
if exit_signal_received {
let now = Instant::now();
let close_wait_deadline = *post_exit_deadline
.get_or_insert_with(|| now + remaining.min(POST_EXIT_CLOSE_WAIT_CAP));
let close_wait_remaining = close_wait_deadline.saturating_duration_since(now);
if close_wait_remaining == Duration::ZERO {
break;
}退出后的等待不是为了让 child 继续运行,而是等待 reader 把已经产生的数据推入共享 buffer。
4.2 关闭通知
退出窗口内同时监听新输出、output_closed_notify、关闭窗口 timer 和 pause change。reader 关闭时可立即结束;只有 token 被取消但 reader 尚未关闭时,才使用 50ms 上限。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: collect_output_until_deadline
let notified = wait_for_output.unwrap_or_else(|| output_notify.notified());
let closed = output_closed_notify.notified();
tokio::pin!(notified);
tokio::pin!(closed);
tokio::select! {
_ = &mut notified => {}
_ = &mut closed => {}
_ = tokio::time::sleep(close_wait_remaining) => break,
_ = Self::wait_for_pause_change(pause_state.as_ref()) => {}
}如果 output task 已将 output_closed 置真,下一轮会在没有数据时直接 break;关闭通知解决了 reader 关闭但没有新字节的情况。
producer 通过 OutputTaskGuard::drop 统一发布关闭。它先用 Release 写入原子标记,再唤醒等待者;collector 使用 Acquire 读取标记,因此在观察到 closed 后可以安全执行最后一次 drain。
源码位置:codex-rs/core/src/unified_exec/process.rs :: OutputTaskGuard::drop
impl Drop for OutputTaskGuard {
fn drop(&mut self) {
self.output_closed.store(true, Ordering::Release);
self.output_closed_notify.notify_waiters();
}
}5. 结果形成
5.1 合并输出
每轮 drain 的 HeadTailBuffer 会合并进本次 collected;循环结束后再生成 omission marker、token 统计和响应正文。轮询函数不直接决定 process ID,状态判断仍由 write_stdin 的 refresh 逻辑完成。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: write_stdin
let original_token_count = usize::try_from(approx_tokens_from_byte_count(
collected_output.total_bytes(),
))
.unwrap_or(usize::MAX);
let output_omitted_bytes = NonZeroUsize::new(collected_output.omitted_bytes());
let collected = collected_output.to_bytes_with_omission_marker();
let chunk_id = generate_chunk_id();
let response = ExecCommandToolOutput {
event_call_id,
chunk_id,
wall_time,
raw_output: collected,
truncation_policy: request.truncation_policy,
max_output_tokens: request.max_output_tokens,
process_id,
exit_code,
original_token_count: Some(original_token_count),
output_omitted_bytes,
hook_command: Some(hook_command),
};total_bytes() 包含已经省略的中段字节,所以 original token count 不会只统计最终保留的 head/tail;output_omitted_bytes 则单独告诉格式化层有多少原始字节未保留。因此“等待超时”不是“命令失败”:只要没有 failure message,响应可以带部分输出、omission metadata 和仍存活的 process ID。
5.2 超时后续取
unified_exec_timeouts 在交互 shell 中执行“等待 5 秒后输出变量”,第一次非空写入只给 10ms,断言结果没有变量;7 秒后用空输入再次 poll,断言取得延迟输出。它直接证明 yield 到期只结束本次工具等待,不终止进程,也不丢弃之后到达的输出。
源码位置:codex-rs/core/src/unified_exec/mod_tests.rs :: unified_exec_timeouts
cd codex-rs
cargo test -p codex-core --lib unified_exec_timeouts -- --test-threads=15.3 暂停测试
unified_exec_pause_blocks_yield_timeout 让命令延迟输出,在 poll 期间暂停 elicitation,断言墙钟时间超过原始 yield 且最终收到输出。它证明暂停会延长有效 deadline,而不是证明所有 UI 暂停来源都相同。
源码位置:codex-rs/core/src/unified_exec/mod_tests.rs :: unified_exec_pause_blocks_yield_timeout
cd codex-rs
cargo test -p codex-core --lib unified_exec_pause_blocks_yield_timeout -- --test-threads=15.4 多轮合并
output_collection_stays_bounded_across_repeated_drains 使用 10 字节容量,分四轮写入并等待每轮共享 buffer 被取空,断言最终结果等于连续写入同一个 10 字节 buffer。output_collection_preserves_omissions_from_drained_buffer 则先制造 omission,再让 collector 取走已关闭 buffer,断言 omission metadata 没有因 mem::take 丢失。
源码位置:
codex-rs/core/src/unified_exec/process_manager_tests.rs :: output_collection_stays_bounded_across_repeated_drainscodex-rs/core/src/unified_exec/process_manager_tests.rs :: output_collection_preserves_omissions_from_drained_buffer
cd codex-rs
cargo test -p codex-core --lib unified_exec::process_manager::tests::output_collection_ -- --test-threads=1这两项验证容量与遗漏连续性,不单独证明 50ms close window;close window 的正确性来自 collector 分支与 producer Release/Acquire 关闭协议。
5.5 首次窗口
当前非 Windows 测试断言首次 1000ms 保持不变,而 1ms 被提升到 250ms。Windows 使用独立测试断言低于 10 秒的首次请求提升到 10 秒。平台条件意味着一次构建只执行其中一组。
源码位置:codex-rs/core/src/unified_exec/process_manager_tests.rs :: initial_exec_yield_time_has_no_platform_floor、initial_exec_yield_time_uses_windows_floor
cd codex-rs
cargo test -p codex-core --lib initial_exec_yield_time_has_no_platform_floor -- --test-threads=16. 源码排查
遇到轮询提前结束,先检查 deadline 是否已到、cancellation token 是否已取消、output_closed 是否已发布;遇到输出缺尾,检查 reader 是否在退出 token 后仍有 drain,以及 50ms close window 是否被耗尽。遇到暂停后仍超时,检查 pause receiver 是否真的发生 changed(),以及 paused_for 是否同时加回两个 deadline。
rg -n "collect_output_until_deadline|POST_EXIT_CLOSE_WAIT_CAP" codex-rs/core/src/unified_exec
rg -n "extend_deadlines_while_paused|wait_for_pause_change|output_closed_notify" codex-rs/core/src/unified_exec
rg -n "pause_blocks_yield_timeout|output_collection_" codex-rs/core/src/unified_exec本篇的主线是:先 drain,再决定等待;普通阶段等待输出、退出、暂停或 deadline;退出后切换到关闭窗口;最后把已收集字节交给响应和状态判断。下一篇将展开后台 AsyncWatcher 如何持续读取 stdout/stderr、发布 delta 并完成清理。
