Skip to content

Unified Exec数据结构

追踪 Unified Exec 的请求、进程句柄、状态快照与进程存储,解释并发访问、容量回收和资源所有权。

基于rust-v0.150.0
CodexRustExecution

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

rust
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

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

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

rust
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

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,
}

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

rust
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

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;

这个顺序形成三道屏障:进程已结束、尾部输出已聚合、同一会话没有并发交互。缺少第一道会过早结束,缺少第二道会丢最后几块输出,缺少第三道则可能与 write_stdin 同时消费状态或发布终态。

3. 状态快照 ​

3.1 四个字段 ​

ProcessState 是本地与远程进程的共同投影。退出码可以为空;失败消息与 sandbox denial 不会因后续退出通知被清除。

源码位置:codex-rs/core/src/unified_exec/process_state.rs :: ProcessState

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

rust
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

rust
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

rust
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

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

rust
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

rust
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

rust
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

rust
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

rust
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_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

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

text
cd codex-rs
cargo test -p codex-core --lib 'unified_exec::process_manager::tests::pruning_' -- --test-threads=1

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

text
cd codex-rs
cargo test -p codex-core --lib unified_exec::tests::terminating_initial_exec_command_rechecks_initial_response_state -- --test-threads=1

9. 故障定位 ​

当 write_stdin 返回 UnknownProcessId,先区分 ID 仍在 reserved_process_ids、已经进入 processes,还是已被释放;条目存在但 Arc::ptr_eq 失败,说明异步操作拿到旧句柄。达到容量时,应同时检查退出状态、last_used、最近 8 个保护集和 interaction_lock,不要按数字 ID 直接删除。

text
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 只能被一条终态路径消费,回收不能打断正在交互或发布终态的进程。