Skip to content

AsyncWatcher后台监控

追踪 Unified Exec 的输出读取、UTF-8 分块、delta 事件、退出收尾和成功失败事件发布。

基于rust-v0.150.0
CodexRustExecution

AsyncWatcher后台监控 ​

Unified Exec 启动后至少有两条后台路径:输出 task 持续消费 stdout/stderr,并把字节写入 transcript 与实时 delta;exit watcher 等待进程退出、输出 drain 和网络拒绝监控,最后只发布一次成功或失败结束事件。它们共享 UnifiedExecProcess 的 token、output close 通知、transcript 和交互锁,但职责不同。

本文承接Unified Exec轮询与等待和Unified Exec会话状态机,面向理解 Tokio task、broadcast channel、UTF-8 边界和 Arc<Mutex<_>> 的读者。范围是 async_watcher.rs 的后台任务和测试,不展开 head-tail 缓冲策略。读完后,你应能解释为什么退出 token 后还要等待输出、broadcast lag 会损失哪一层观察结果,以及成功/失败事件如何选择聚合文本和收尾 metrics。

1. 输出接管 ​

1.1 两个消费者 ​

start_streaming_output 从 UnifiedExecProcess 订阅 broadcast receiver,复制 output close 通知和取消 token,再启动 Tokio task。实时 delta 与 transcript 都由这个 task 产生;轮询使用的共享 output buffer 是另一条生产/消费路径。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: start_streaming_output

rust
let mut receiver = process.output_receiver();
let output_drained = process.output_drained_notify();
let exit_token = process.cancellation_token();
let OutputHandles {
    output_closed,
    output_closed_notify,
    ..
} = process.output_handles().clone();

let emitter = Emitter {
    remaining_deltas: MAX_EXEC_OUTPUT_DELTAS_PER_CALL,
    session: Arc::clone(&context.session),
    turn: Arc::clone(&context.step_context.turn),
    call_id: context.call_id.clone(),
};

tokio::spawn(async move {
    use tokio::sync::broadcast::error::RecvError;

    let mut output: Buffer = Buffer {
        pending: Vec::new(),
        transcript,
        emitter,
    };

    let mut grace_sleep: Option<Pin<Box<Sleep>>> = None;
    let output_closed_notified = output_closed_notify.notified();
    tokio::pin!(output_closed_notified);
    let mut output_complete = false;
});

Emitter 捕获 Session、StepContext 中的 Turn 和 call ID,保证异步发送事件时上下文仍可用;Buffer 则统一持有 pending、transcript 与事件预算。这个拆分让“如何切帧”与“是否还允许发事件”成为两个独立对象。

1.2 广播与 lag ​

reader 通过 broadcast 发送 chunk。watcher 收到 Lagged 时直接跳过,不把它当作进程失败。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: start_streaming_output

rust
received = receiver.recv() => {
    let chunk = match received {
        Ok(chunk) => chunk,
        Err(RecvError::Lagged(_)) => continue,
        Err(RecvError::Closed) => {
            output_complete = true;
            break;
        }
    };
    output.push(chunk).await;
}

这里必须准确区分两份缓冲:AsyncWatcher 的 transcript 由 output.push 写入,因此 lagged chunk 不会进入该 transcript,也不会生成 delta;UnifiedExecProcess 内另有供 write_stdin/首次响应 drain 的有界 output buffer。两者来自同一 producer,但消费进度独立。结束事件优先使用 watcher transcript,工具响应使用轮询 buffer,所以极端 lag 下两种观察面可能不完全相同。

2. UTF-8分块 ​

2.1 transcript优先 ​

Buffer::push 收到 producer chunk 后,第一步就把整块原始字节写入 transcript;切帧和事件预算只影响 delta。即使后续没有剩余事件额度,这一块仍保留在 transcript 中。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: Buffer::push

rust
let Self {
    pending,
    transcript,
    emitter,
} = self;

transcript.lock().await.push_chunk(&bytes);

if pending.is_empty() && bytes.len() <= MAX_BYTES {
    emitter
        .emit(|| {
            let complete = utf8_boundary(&bytes);
            pending.extend(bytes.drain(complete..));
            bytes
        })
        .await;
    return;
}

快速路径复用 producer 的 Vec<u8>,只把末尾不完整 UTF-8 scalar 留进 pending,避免为常见的小块输出再次分配。MAX_BYTES 默认 8192,并在编译期要求至少能容纳一个四字节 scalar。

2.2 通用切帧 ​

当 pending 非空或输入大于单事件上限时,next_chunk 先用 pending 填充 frame,再从新输入取剩余空间。utf8_boundary 只保留最多三个字节的不完整后缀;完整字符和 malformed bytes 都可以进入当前 frame。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: Buffer::push、utf8_boundary

rust
let mut bytes = bytes.as_slice();
let mut next_chunk = || {
    let space = MAX_BYTES.saturating_sub(pending.len());
    let (prefix, rest) = bytes.split_at_checked(space).unwrap_or((bytes, &[]));
    let mut chunk = Vec::with_capacity(pending.len().saturating_add(prefix.len()));
    chunk.append(pending);
    chunk.extend_from_slice(prefix);
    bytes = rest;

    let complete = utf8_boundary(&chunk);
    pending.extend(chunk.drain(complete..));
    chunk
};
while emitter.emit(&mut next_chunk).await {}

当前实现不把非法字节逐个转成 replacement character。utf8_boundary 遇到确定非法序列时跨过其 error_len,只在 error_len == None 时认定末尾是不完整 scalar。delta 的 chunk 仍是原始字节,最终文本展示才使用 lossy 转换。

2.3 delta预算 ​

Emitter::emit 在构造 frame 之前检查剩余额度。额度耗尽后 closure 不再执行,Buffer::push 停止为这一块构造更多 delta;但整块数据已经在第一步写入 transcript。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: Emitter::emit

rust
let Some(remaining) = remaining_deltas.checked_sub(1) else {
    return false;
};
let chunk = make_chunk();
let emit = !chunk.is_empty();
if emit {
    let event = ExecCommandOutputDeltaEvent {
        call_id: call_id.clone(),
        stream: ExecOutputStream::Stdout,
        chunk,
    };
    session
        .send_event(turn.as_ref(), EventMsg::ExecCommandOutputDelta(event))
        .await;
    *remaining_deltas = remaining;
}
emit

默认预算是每次命令 10,000 个 delta。它限制事件数量,不限制 transcript 的 head-tail 容量;这两套上限作用于不同观察面。

3. 退出协作 ​

3.1 trailing grace ​

输出 task 收到 exit token 后进入 TRAILING_OUTPUT_GRACE 100ms 窗口;如果 output close 先到,可以提前结束。这样 child exit 与 reader close 的时间差不会丢掉尾部输出。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: start_streaming_output

rust
_ = exit_token.cancelled(), if grace_sleep.is_none() => {
    let deadline = Instant::now() + TRAILING_OUTPUT_GRACE;
    grace_sleep.replace(Box::pin(tokio::time::sleep_until(deadline)));
}
_ = async {
    if let Some(sleep) = grace_sleep.as_mut() {
        sleep.as_mut().await;
    }
}, if grace_sleep.is_some() => break,
_ = &mut output_closed_notified, if grace_sleep.is_some() => {
    output_closed_notified.set(output_closed_notify.notified());
}

3.2 最终try_recv ​

确认 close 后再 try_recv,把生产者设置 closed 前发布的剩余 broadcast chunk 最后消费一次;随后通知 output_drained。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: start_streaming_output

rust
output_complete |= output_closed.load(Ordering::Acquire);
if output_complete {
    loop {
        let chunk = match receiver.try_recv() {
            Ok(chunk) => chunk,
            Err(tokio::sync::broadcast::error::TryRecvError::Lagged(_)) => continue,
            Err(
                tokio::sync::broadcast::error::TryRecvError::Empty
                | tokio::sync::broadcast::error::TryRecvError::Closed,
            ) => break,
        };
        output.push(chunk).await;
    }
}
output.finish().await;
output_drained.notify_one();

finish 会尝试发送 pending 中最后的不完整字节;如果 delta 预算已耗尽,它不会发事件,但原始字节此前已经进入 transcript。output_drained 只表示输出 task 完成可见数据处理,不表示 child 仍存活。

4. 结束事件 ​

4.1 单次结束 ​

exit watcher 等待 cancellation、output drain、网络 denial monitor 和 interaction lock,之后取得共享 metrics sidecar,并根据 failure_message 选择 Failure 或 Success emitter。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: spawn_exit_watcher

rust
exit_token.cancelled().await;
output_drained.notified().await;
if let Some(network_denial_monitor) = network_denial_monitor {
    let _ = network_denial_monitor.await;
}
let _interaction_guard = interaction_lock.lock_owned().await;

let duration = Instant::now().saturating_duration_since(started_at);
let plugin_metrics_sidecar = plugin_metrics_sidecar
    .as_ref()
    .and_then(take_plugin_metrics_sidecar);
if let Some(message) = process.failure_message() {
    drop(plugin_metrics_sidecar);
    emit_failed_exec_end_for_unified_exec(
        session_ref,
        turn_ref,
        call_id,
        command,
        cwd,
        Some(process_id.to_string()),
        plugin_attribution,
        transcript,
        String::new(),
        message,
        duration,
    )
    .await;
} else {
    let exit_code = process.exit_code().unwrap_or(-1);
    finish_and_track_measurements(
        plugin_metrics_sidecar,
        exit_code,
        &session_ref,
        &turn_ref,
        &call_id,
    )
    .await;
    emit_exec_end_for_unified_exec(
        session_ref,
        turn_ref,
        call_id,
        command,
        cwd,
        Some(process_id.to_string()),
        plugin_attribution,
        transcript,
        String::new(),
        exit_code,
        duration,
    )
    .await;
}

失败分支不会把 -1 当作正常插件 exit code 上报,而是丢弃 sidecar 并保留失败消息。成功分支先用真实 exit code 完成测量,再发布结束事件。sidecar 通过 Option::take 与首次响应路径共享,保证只被一方消费。

4.2 聚合来源 ​

成功事件优先从 transcript 读取带 omission marker 的聚合文本;失败事件把 failure message 追加到输出或作为 stderr。

源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: resolve_aggregated_output

rust
async fn resolve_aggregated_output(
    transcript: &Arc<Mutex<HeadTailBuffer>>,
    fallback: String,
) -> String {
    let guard = transcript.lock().await;
    if guard.retained_bytes() == 0 {
        return fallback;
    }
    String::from_utf8_lossy(&guard.to_bytes_with_omission_marker()).to_string()
}

5. 四组验证 ​

5.1 跨chunk字符 ​

streaming_output_preserves_multibyte_characters_across_chunks 把 é 的两个字节拆成两次 broadcast,随后关闭 producer 并退出。断言只产生一个完整 é delta,transcript 也是完整的两个字节,没有把不完整前导字节提前发出。

源码位置:codex-rs/core/src/unified_exec/async_watcher_tests.rs :: streaming_output_preserves_multibyte_characters_across_chunks

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

5.2 关闭与grace ​

streaming_output_finishes_on_close_without_waiting_for_grace 在 exit 后 50ms 发送尾部输出并关闭 producer,断言收尾早于 100ms grace;streaming_output_keeps_grace_as_fallback_without_close 不关闭 producer,断言 task 等满 grace,同时 pending 的不完整字节仍进入 transcript 和最终 delta。

源码位置:

  • codex-rs/core/src/unified_exec/async_watcher_tests.rs :: streaming_output_finishes_on_close_without_waiting_for_grace
  • codex-rs/core/src/unified_exec/async_watcher_tests.rs :: streaming_output_keeps_grace_as_fallback_without_close
text
cd codex-rs
cargo test -p codex-core --lib streaming_output_finishes_on_close_without_waiting_for_grace -- --test-threads=1
cargo test -p codex-core --lib streaming_output_keeps_grace_as_fallback_without_close -- --test-threads=1

5.3 非法字节与预算 ​

两个 utf8_boundary 单测覆盖完整多字节字符、不完整后缀和 malformed sequence。streaming_output_bounds_invalid_bytes_and_keeps_the_full_transcript 使用 8 字节 frame 与仅 2 个 delta 额度,输入包含非法字节、emoji 和不完整 é;断言事件只有两帧且不切断 emoji,而 transcript 保留全部原始字节,包括预算耗尽后的输入。

源码位置:

  • codex-rs/core/src/unified_exec/async_watcher_tests.rs :: utf8_boundary_preserves_complete_characters
  • codex-rs/core/src/unified_exec/async_watcher_tests.rs :: utf8_boundary_batches_malformed_output
  • codex-rs/core/src/unified_exec/async_watcher_tests.rs :: streaming_output_bounds_invalid_bytes_and_keeps_the_full_transcript
text
cd codex-rs
cargo test -p codex-core --lib utf8_boundary_ -- --test-threads=1
cargo test -p codex-core --lib streaming_output_bounds_invalid_bytes_and_keeps_the_full_transcript -- --test-threads=1

5.4 延迟拒绝 ​

exit_watcher_waits_for_late_network_denial_before_classifying_end 先发送 exit,再在 10ms 后调用 fail_and_terminate("LATE_DENIAL")。断言最终 CommandExecution 为 Failed、exit code 为 -1、聚合文本包含拒绝原因,证明 watcher 不会在 network monitor 完成前发布成功事件。

源码位置:codex-rs/core/src/unified_exec/async_watcher_tests.rs :: exit_watcher_waits_for_late_network_denial_before_classifying_end

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

这些测试没有主动制造 broadcast receiver lag,因此不能证明 lag 场景下 transcript 完整;metrics 分支也传入 None,没有验证真实 plugin sidecar 文件的完成行为。相关结论来自当前分支所有权和错误处理源码。

6. 源码排查 ​

遇到实时 delta 缺失,先区分 broadcast Lagged、UTF-8 pending、单事件 8192 字节上限和全局 delta 数量上限;若结束事件 transcript 同样缺片段,检查 watcher receiver 是否 lagged,而不是假定轮询 buffer 会自动回填 transcript。遇到结束事件缺尾,检查 exit token、trailing grace、output closed 和最终 try_recv;遇到成功/失败分类错误,检查 network monitor、failure message 和 metrics sidecar 的取得顺序。

text
rg -n "start_streaming_output|struct Buffer|struct Emitter|utf8_boundary" codex-rs/core/src/unified_exec/async_watcher.rs
rg -n "TRAILING_OUTPUT_GRACE|output_drained|spawn_exit_watcher" codex-rs/core/src/unified_exec
rg -n "async_watcher::tests|MAX_EXEC_OUTPUT_DELTAS|plugin_metrics_sidecar" codex-rs/core/src/unified_exec

本篇的主线是:未 lag 的输出 chunk 先完整写入 transcript,再由 Buffer 按 UTF-8 边界和事件预算生成 delta;退出后等待尾部输出和关闭标记;exit watcher 在 output drain、网络分类与交互锁之后消费 metrics sidecar,并选择唯一的成功或失败结束事件。