Skip to content

Unified Exec轮询与等待

追踪 Unified Exec 的输出轮询、截止时间、暂停状态、退出后关闭窗口与等待通知。

基于rust-v0.150.0
CodexRustExecution

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

rust
#[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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

text
cd codex-rs
cargo test -p codex-core --lib unified_exec_timeouts -- --test-threads=1

5.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

text
cd codex-rs
cargo test -p codex-core --lib unified_exec_pause_blocks_yield_timeout -- --test-threads=1

5.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_drains
  • codex-rs/core/src/unified_exec/process_manager_tests.rs :: output_collection_preserves_omissions_from_drained_buffer
text
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

text
cd codex-rs
cargo test -p codex-core --lib initial_exec_yield_time_has_no_platform_floor -- --test-threads=1

6. 源码排查 ​

遇到轮询提前结束,先检查 deadline 是否已到、cancellation token 是否已取消、output_closed 是否已发布;遇到输出缺尾,检查 reader 是否在退出 token 后仍有 drain,以及 50ms close window 是否被耗尽。遇到暂停后仍超时,检查 pause receiver 是否真的发生 changed(),以及 paused_for 是否同时加回两个 deadline。

text
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 并完成清理。