Unified Exec数据结构
上一文把命令执行拆成 handler、runtime、统一进程和 backend。本篇深入统一执行层的数据结构:谁冻结本次调用上下文,谁保存启动请求,谁拥有实际进程,谁保存输出与退出状态,谁把可恢复会话放进 ProcessStore。本文不展开 spawn 参数、完整输出截断算法和 exec-server RPC 字段,它们属于后续专题。
读者需要理解 Rust 的 Arc、Weak、Mutex、watch 和 CancellationToken。建议先读Codex执行体系总览,以免把统一层 process_id 与 OS pid 混淆。读完后,你应能从 write_stdin 找到对应 ProcessEntry,解释 last_used 的更新时机,并判断容量达到 64 时哪个条目可能被回收。
1. 调用与请求
1.1 调用上下文
UnifiedExecContext 把一次工具调用必须跨越 handler、manager、orchestrator、runtime 和 watcher 的信息集中起来。它不再只保存 TurnContext,而是保存完整 StepContext 与调用取消令牌。
源码位置:codex-rs/core/src/unified_exec/mod.rs :: UnifiedExecContext
pub(crate) struct UnifiedExecContext {
pub session: Arc<Session>,
pub step_context: Arc<StepContext>,
pub cancellation_token: CancellationToken,
pub call_id: String,
}
impl UnifiedExecContext {
pub fn new(
session: Arc<Session>,
step_context: Arc<StepContext>,
cancellation_token: CancellationToken,
call_id: String,
) -> Self {
Self {
session,
step_context,
cancellation_token,
call_id,
}
}
}StepContext 的价值不只是多包了一层 Turn。环境选择、审批上下文和工具生命周期都必须使用同一步骤快照,不能在命令启动过程中重新读取可能已经变化的 Session 状态。调用取消令牌则属于工具分派:它会传入 ToolCtx 和 apply-patch 拦截路径。
这里要区分两个同名但不同作用域的 token:UnifiedExecContext::cancellation_token 表示本次工具调用被取消;后文 OutputHandles::cancellation_token 表示底层进程已经退出或被终止,供输出收集器与 exit watcher 收尾。把两者混为一谈,会误以为 Turn 取消后后台进程一定已经退出。
1.2 启动请求
ExecCommandRequest 是 handler 交给 manager 的一次启动描述。它既包含命令和统一 ID,也包含当前 Turn 的环境、shell 模式、网络与权限快照。
源码位置:codex-rs/core/src/unified_exec/mod.rs :: ExecCommandRequest
#[derive(Debug)]
pub(crate) struct ExecCommandRequest {
pub command: Vec<String>,
pub shell_type: ShellType,
pub hook_command: String,
pub process_id: i32,
pub yield_time_ms: u64,
pub max_output_tokens: Option<usize>,
pub cwd: PathUri,
pub sandbox_cwd: PathUri,
pub turn_environment: TurnEnvironment,
pub shell_mode: UnifiedExecShellMode,
pub network: Option<NetworkProxy>,
pub tty: bool,
pub sandbox_permissions: SandboxPermissions,
pub additional_permissions: Option<AdditionalPermissionProfile>,
pub additional_permissions_preapproved: bool,
pub justification: Option<String>,
pub prefix_rule: Option<Vec<String>>,
}这是“本次启动”的所有权边界,不是后台会话的长期状态。启动成功后,manager 只把 process_id、cwd、hook 命令、TTY 和事件关联字段复制进 ProcessEntry;shell 模式和权限只参与当前运行尝试。
1.3 交互请求
WriteStdinRequest 不重复 shell、权限和环境字段。它只描述要查找的统一 ID、本次输入、等待上限和输出截断方式。
源码位置:codex-rs/core/src/unified_exec/mod.rs :: WriteStdinRequest
#[derive(Debug)]
pub(crate) struct WriteStdinRequest<'a> {
pub process_id: i32,
pub input: &'a str,
pub yield_time_ms: u64,
pub max_output_tokens: Option<usize>,
pub truncation_policy: TruncationPolicy,
pub interaction_event: Option<WriteStdinInteractionEvent<'a>>,
}因此 write_stdin 不会重新解释启动参数。输入为空时,它仍可作为轮询请求;失败首先来自 ID 查找、已退出状态或底层 stdin,而不是 shell 参数解析。
2. 进程所有权
2.1 传输句柄
ProcessHandle 是统一层唯一需要知道本地与远程差异的枚举。本地分支拥有 PTY/pipe 会话,远程分支持有 exec-server 代理。
源码位置:codex-rs/core/src/unified_exec/process.rs :: ProcessHandle
enum ProcessHandle {
Local(Box<ExecCommandSession>),
ExecServer(Arc<dyn ExecProcess>),
}上层不能直接依赖枚举分支,而应使用 UnifiedExecProcess 的 write、terminate、interrupt、has_exited 和 exit_code。本地用 Box 独占会话析构权,远程用 Arc 支持输出 task 和交互调用共享 proxy。
2.2 输出句柄
OutputHandles 把数据、唤醒、关闭标记和取消信号分开保存。四者用途不同,不能相互替代。
源码位置: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,
}MAX_BYTES 默认是 1 MiB,使“输出上限”成为缓冲类型的一部分;测试可以用 HeadTailBuffer<10> 构造极小容量,验证多次 drain 后仍保持全局 head/tail 摘要。buffer 用异步 mutex 保护内容,output_notify 表示可能有新数据,output_closed 保存可重复读取的关闭事实,output_closed_notify 避免关闭后仍等待到超时。
这里的 cancellation token 更准确地说是进程终态信号。它在 child 退出、远程失败或显式 terminate 时取消,让 collector 开始最后一次 drain;输出 producer 是否已经完全关闭,仍要由 output_closed 单独确认。只看 token 会漏掉退出后的尾部输出,只看 notify 又可能错过先于订阅发生的关闭。
2.3 统一对象
源码位置:codex-rs/core/src/unified_exec/process.rs :: UnifiedExecProcess
pub(crate) struct UnifiedExecProcess {
process_handle: ProcessHandle,
output_tx: broadcast::Sender<Vec<u8>>,
output: OutputHandles,
output_drained: Arc<Notify>,
interaction_lock: Arc<Mutex<()>>,
state_tx: watch::Sender<ProcessState>,
state_rx: watch::Receiver<ProcessState>,
output_task: Option<JoinHandle<()>>,
sandbox_type: SandboxType,
_spawn_lifecycle: Option<SpawnLifecycleHandle>,
}output_tx 面向实时订阅,buffer 面向迟到的轮询;interaction_lock 串行化写入、轮询、终止和回收;watch channel 发布最新状态快照。_spawn_lifecycle 保存的是启动附属资源,而不只是 fd 列表。例如 zsh-fork 的 after_spawn 会关闭父进程持有的 client socket,但 EscalationSession 仍需随统一进程存活,直到对象析构时完成剩余清理。
2.4 终态握手
进程退出与输出读完是两个事件。streaming task 先等待进程 token;收到退出信号后进入短暂尾部输出窗口,并通过 output_closed 的 Release/Acquire 配对完成最终 drain,最后才通知 output_drained。exit watcher 必须同时等到 token 与 output_drained,才能生成完整的结束事件。
源码位置:codex-rs/core/src/unified_exec/async_watcher.rs :: start_streaming_output、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;这个顺序形成三道屏障:进程已结束、尾部输出已聚合、同一会话没有并发交互。缺少第一道会过早结束,缺少第二道会丢最后几块输出,缺少第三道则可能与 write_stdin 同时消费状态或发布终态。
3. 状态快照
3.1 四个字段
ProcessState 是本地与远程进程的共同投影。退出码可以为空;失败消息与 sandbox denial 不会因后续退出通知被清除。
源码位置:codex-rs/core/src/unified_exec/process_state.rs :: ProcessState
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub(crate) struct ProcessState {
pub(crate) has_exited: bool,
pub(crate) exit_code: Option<i32>,
pub(crate) failure_message: Option<String>,
pub(crate) sandbox_denied: bool,
}
pub(crate) fn exited(&self, exit_code: Option<i32>) -> Self {
Self {
has_exited: true,
exit_code,
failure_message: self.failure_message.clone(),
sandbox_denied: self.sandbox_denied,
}
}状态更新采用完整快照替换。这样“先记录网络失败,再收到 child exit”不会把失败降级为普通退出,也避免多个 task 分别改字段形成不可观察的中间态。
3.2 本地与远程
本地进程可直接查询 ExecCommandSession,远程进程只能依赖 exec-server 事件更新 watch 状态。
源码位置:codex-rs/core/src/unified_exec/process.rs :: has_exited、exit_code
pub(super) fn has_exited(&self) -> bool {
let state = self.state_rx.borrow().clone();
match &self.process_handle {
ProcessHandle::Local(process_handle) => state.has_exited || process_handle.has_exited(),
ProcessHandle::ExecServer(_) => state.has_exited,
}
}
pub(super) fn exit_code(&self) -> Option<i32> {
let state = self.state_rx.borrow().clone();
match &self.process_handle {
ProcessHandle::Local(process_handle) => {
state.exit_code.or_else(|| process_handle.exit_code())
}
ProcessHandle::ExecServer(_) => state.exit_code,
}
}所以本地分支可能在 watcher 发布前已经观察到退出;远程分支必须等待协议事件。要求严格终态时,调用方还要等待确认终止或输出关闭,不能只读取一次 has_exited。
4. 存储不变量
4.1 ProcessEntry
后台条目除进程 Arc 外,还保存事件关联、展示信息、网络审批和并发保护字段。
源码位置:codex-rs/core/src/unified_exec/mod.rs :: ProcessEntry
struct ProcessEntry {
process: Arc<UnifiedExecProcess>,
plugin_metrics_sidecar: Option<SharedPluginMetricsSidecar>,
call_id: String,
process_id: i32,
cwd: PathUri,
initial_exec_command_active: Arc<std::sync::atomic::AtomicBool>,
hook_command: String,
tty: bool,
network_approval: Option<DeferredNetworkApproval>,
session: Weak<Session>,
last_used: tokio::time::Instant,
}Weak<Session> 防止后台 shell 反向保活整个 Session;需要订阅暂停状态或结束网络审批时才 upgrade。initial_exec_command_active 保护首轮工具响应:process ID 尚未交给模型时,终止路径不能提前删条目。last_used 是最近交互时间,不等于启动时间。
4.2 Metrics所有权
插件命令可能携带 PluginMetricsSidecar。短命令由首次调用路径完成测量,长命令则把 sidecar 放进 ProcessEntry,同时交给 exit watcher。由于 manager 可能在首次 yield 后立即观察到进程退出,watcher 也可能先收到终态,两条路径会竞争同一份资源。
源码位置:codex-rs/core/src/unified_exec/mod.rs :: SharedPluginMetricsSidecar、take_plugin_metrics_sidecar
type SharedPluginMetricsSidecar = Arc<std::sync::Mutex<Option<PluginMetricsSidecar>>>;
fn take_plugin_metrics_sidecar(
sidecar: &SharedPluginMetricsSidecar,
) -> Option<PluginMetricsSidecar> {
sidecar
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
}这里使用 Arc<Mutex<Option<T>>>,不是为了反复共享 sidecar,而是为了实现“多个终态观察者、唯一消费者”。第一条收尾路径通过 Option::take 取得所有权,后到的路径只会得到 None,从而避免同一插件命令重复上报测量或重复关闭输出文件。
4.3 双集合
ProcessStore 同时维护已注册条目和已预留 ID。分配 ID 发生在 spawn 前,因此两者不能合并成一个 map。
源码位置:codex-rs/core/src/unified_exec/mod.rs :: ProcessStore
#[derive(Default)]
pub(crate) struct ProcessStore {
processes: HashMap<i32, ProcessEntry>,
reserved_process_ids: HashSet<i32>,
}
impl ProcessStore {
fn remove(&mut self, process_id: i32) -> Option<ProcessEntry> {
self.reserved_process_ids.remove(&process_id);
self.processes.remove(&process_id)
}
}remove 必须同时清理两个集合。只删 map 会泄漏保留 ID;只删保留集合可能允许新进程复用仍有旧条目的数字 ID。
4.4 Manager
源码位置:codex-rs/core/src/unified_exec/mod.rs :: UnifiedExecProcessManager
pub(crate) struct UnifiedExecProcessManager {
process_store: Mutex<ProcessStore>,
max_write_stdin_yield_time_ms: u64,
}
pub(crate) fn new(max_write_stdin_yield_time_ms: u64) -> Self {
Self {
process_store: Mutex::new(ProcessStore::default()),
max_write_stdin_yield_time_ms: max_write_stdin_yield_time_ms
.max(MIN_EMPTY_YIELD_TIME_MS),
}
}manager 由 Session services 共享,handler 只借用。锁内只做查找、替换和时间戳更新;I/O、等待输出、终止确认和网络审批清理都应在锁外,否则一个慢进程会阻塞所有会话。
5. 交互借用
prepare_process_handles 在锁内取得条目、校验进程身份、更新 last_used,然后复制一次交互所需的进程、输出、暂停状态、Session、网络审批和展示字段。它不会把 store 锁带入轮询。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: prepare_process_handles
let entry = store
.processes
.get_mut(&process_id)
.ok_or(UnifiedExecError::UnknownProcessId { process_id })?;
if !Arc::ptr_eq(&entry.process, expected_process) {
return Err(UnifiedExecError::UnknownProcessId { process_id });
}
entry.last_used = Instant::now();
let output = entry.process.output_handles().clone();
let pause_state = entry
.session
.upgrade()
.map(|session| session.subscribe_elicitation_pause_state());
let session = entry.session.upgrade();
Ok(PreparedProcessHandles {
process: Arc::clone(&entry.process),
output,
pause_state,
session,
network_approval: entry.network_approval.clone(),
call_id: entry.call_id.clone(),
hook_command: entry.hook_command.clone(),
process_id: entry.process_id,
tty: entry.tty,
})Arc::ptr_eq 防止一个旧异步任务持有相同数字 ID,却操作已经替换的新进程。pause_state 让用户交互暂停期间的等待 deadline 同步后移;弱 Session 无法升级时,进程仍可被读取和清理,只是不再依赖 Session 事件通道。复制这些句柄后释放 mutex,等待输出期间其他调用仍可查询和分配 ID。
6. 容量回收
6.1 选择策略
统一执行最多保留 64 个进程条目。达到上限时,算法保护最近使用的 8 个 ID,优先选择保护集之外的已退出进程,否则回退到最久未使用条目。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: process_id_to_prune_from_meta
let mut by_recency = meta.to_vec();
by_recency.sort_by_key(|(_, last_used, _)| Reverse(*last_used));
let protected: HashSet<i32> = by_recency
.iter()
.take(8)
.map(|(process_id, _, _)| *process_id)
.collect();
let mut lru = meta.to_vec();
lru.sort_by_key(|(_, last_used, _)| *last_used);
if let Some((process_id, _, _)) = lru
.iter()
.find(|(process_id, _, exited)| !protected.contains(process_id) && *exited)
{
return Some(*process_id);
}返回候选后还要尝试取得 interaction_lock。若已退出进程正在发布终端事件,算法宁可暂时超过软上限,也不因为它被锁住就回收一个活进程。
6.2 锁保护
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: prune_processes_if_needed
if let Some(interaction_lock) = candidate_process
.as_ref()
.map(|process| process.interaction_lock())
&& let Ok(_interaction_guard) = interaction_lock.try_lock_owned()
{
return store.remove(process_id);
}回收只在交互锁可立即取得时发生。被移出的条目在释放 store mutex 后才注销网络审批并终止进程,避免持锁 await。
7. 终止竞态
单个终止采用两阶段操作:先复制进程 Arc 并释放 store 锁,等待 backend 确认;随后重新加锁,校验 map 中仍是同一个 Arc,再判断首轮保护位是否允许删除。
源码位置:codex-rs/core/src/unified_exec/process_manager.rs :: terminate_process
let (process, already_exited) = {
let store = self.process_store.lock().await;
let Some(entry) = store.processes.get(&process_id) else {
return false;
};
(Arc::clone(&entry.process), entry.process.has_exited())
};
if !already_exited && process.terminate_confirmed().await.is_err() {
return false;
}重新加锁后使用 Arc::ptr_eq 防止误删替换后的新进程;initial_exec_command_active 为真时保留条目,等首轮响应完成后再清理。
8. 三组验证
8.1 有界输出
output_collection_stays_bounded_across_repeated_drains 使用容量仅 10 字节的 OutputHandles<10>,分四次写入共 29 字节,并强制 collector 每轮都把共享 buffer 清空。最终结果必须等于把四块数据连续压入同一个 HeadTailBuffer<10> 的摘要。这证明“drain 共享 buffer”不会重置全局 head/tail 预算。
相邻的 output_collection_preserves_omissions_from_drained_buffer 先在共享 buffer 中制造 omission,再执行一次收集,断言 omitted byte 计数仍被带入结果。它验证容量与遗漏元数据,不验证模型侧 token 截断。
源码位置:
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=18.2 回收选择
pruning_prefers_exited_processes_outside_recently_used 与相邻测试构造达到容量的 store,设置退出状态和 last_used,断言保护集外的已退出条目优先;没有退出条目时才选择 LRU。它证明选择策略,不证明真实 OS process 已被终止。
源码位置:codex-rs/core/src/unified_exec/process_manager_tests.rs :: pruning_prefers_exited_processes_outside_recently_used
cd codex-rs
cargo test -p codex-core --lib 'unified_exec::process_manager::tests::pruning_' -- --test-threads=18.3 首轮保护
terminating_initial_exec_command_rechecks_initial_response_state 使用阻塞终止的 fake process:终止开始时保护位为 true,等待期间改为 false,最后断言条目被删除。它证明终止路径会重新检查首轮状态,不证明真实 PTY 或远程 backend 的延迟。
源码位置:codex-rs/core/src/unified_exec/mod_tests.rs :: terminating_initial_exec_command_rechecks_initial_response_state
cd codex-rs
cargo test -p codex-core --lib unified_exec::tests::terminating_initial_exec_command_rechecks_initial_response_state -- --test-threads=19. 故障定位
当 write_stdin 返回 UnknownProcessId,先区分 ID 仍在 reserved_process_ids、已经进入 processes,还是已被释放;条目存在但 Arc::ptr_eq 失败,说明异步操作拿到旧句柄。达到容量时,应同时检查退出状态、last_used、最近 8 个保护集和 interaction_lock,不要按数字 ID 直接删除。
rg -n "struct ProcessEntry|struct ProcessStore|prepare_process_handles" codex-rs/core/src/unified_exec
rg -n "initial_exec_command_active|Arc::ptr_eq|reserved_process_ids" codex-rs/core/src/unified_exec
rg -n "prune_processes_if_needed|process_id_to_prune_from_meta" codex-rs/core/src/unified_exec下一篇将沿这些字段追踪 ExecCommandRequest 如何变成真实本地或远程进程。本篇需要记住的不是字段数量,而是四条不变量:调用取消与进程退出属于不同作用域,保留 ID 与 map 必须一起释放,metrics sidecar 只能被一条终态路径消费,回收不能打断正在交互或发布终态的进程。
