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
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
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
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
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
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
_ = 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
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
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
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
cd codex-rs
cargo test -p codex-core --lib streaming_output_preserves_multibyte_characters_across_chunks -- --test-threads=15.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_gracecodex-rs/core/src/unified_exec/async_watcher_tests.rs :: streaming_output_keeps_grace_as_fallback_without_close
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=15.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_characterscodex-rs/core/src/unified_exec/async_watcher_tests.rs :: utf8_boundary_batches_malformed_outputcodex-rs/core/src/unified_exec/async_watcher_tests.rs :: streaming_output_bounds_invalid_bytes_and_keeps_the_full_transcript
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=15.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
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 的取得顺序。
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,并选择唯一的成功或失败结束事件。
