ExecServer进程RPC
exec-server 的进程 RPC 是一个由 session 持有的生命周期对象,而不是五个彼此独立的函数。请求先经过 ExecServerHandler 找到当前 SessionHandle,再交给 ProcessHandler 和 LocalProcess;本地 backend 把 child 的 stdout、stderr、exit 状态写入同一张进程表,远程 backend 则通过同名 RPC 转发到 executor。 read 读取 retained output,事件订阅读取 live/replay 流,write 以 write_id 抑制重试造成的重复 输入,signal 与 terminate 修改同一个进程所有者。
本文承接ExecServer消息模型与Unified Exec创建进程。 前者解释 envelope 和 pending,后者解释调用方怎样选择本地或远程 backend;本文只讲 exec-server 进程 RPC 的字段、状态和消费者,不展开 sandbox 命令构造。读完后,你应能从 process/start 找到 child 所有权,解释 read 的序号与等待窗口,判断 write 重试为何不会重复写入,并区分 Exited、Closed 和 Failed 三种结束语义。
1. 调用边界
1.1 handler到backend
协议注册把五个方法映射到 ExecServerHandler。handler 只负责初始化门槛和 session 查找;真正的进程 操作由 session 中的 ProcessHandler 继续转交给 LocalProcess。这个分层让 session resume 可以保留 同一个 process owner,也让同一套 ExecProcess trait 同时承载本地和远程实现。
源码位置:
codex-rs/exec-server/src/server/registry.rs::register_process_routescodex-rs/exec-server/src/server/handler.rs::ExecServerHandler::exec、exec_read、exec_write、signal、terminatecodex-rs/exec-server/src/server/process_handler.rs::ProcessHandler
// server/registry.rs
router.request(EXEC_METHOD, |handler: Arc<ExecServerHandler>, params: ExecParams| async move {
handler.exec(params).await
});
router.request(EXEC_READ_METHOD, |handler: Arc<ExecServerHandler>, params: ReadParams| async move {
handler.exec_read(params).await
});
router.request(EXEC_WRITE_METHOD, |handler: Arc<ExecServerHandler>, params: WriteParams| async move {
handler.exec_write(params).await
});
router.request(EXEC_SIGNAL_METHOD, |handler: Arc<ExecServerHandler>, params: SignalParams| async move {
handler.signal(params).await
});
router.request(EXEC_TERMINATE_METHOD, |handler: Arc<ExecServerHandler>, params: TerminateParams| async move {
handler.terminate(params).await
});源码位置:codex-rs/exec-server/src/server/handler.rs :: ExecServerHandler::exec、ExecServerHandler::exec_read
// server/handler.rs
pub(crate) async fn exec(&self, params: ExecParams) -> Result<ExecResponse, JSONRPCErrorError> {
let session = self.require_initialized_for("exec")?;
session.process().exec(params).await
}
pub(crate) async fn exec_read(
&self,
params: ReadParams,
) -> Result<ReadResponse, JSONRPCErrorError> {
let session = self.require_initialized_for("exec")?;
let response = session.process().exec_read(params).await?;
self.require_session_attached()?;
Ok(response)
}exec_read 在读取完成后再次检查 session attachment,因为长等待窗口期间连接可能已经被替换或分离。 这不是对每个方法都重复的装饰,而是 read 这类可阻塞调用的生命周期边界。
1.2 统一trait
ExecProcess 暴露的是调用方真正需要的能力:读取 retained output、订阅事件、写入 stdin、发送 interrupt 和终止。它不泄露 LocalProcess 的 map、PTY 或 channel 细节。
源码位置:codex-rs/exec-server/src/process.rs :: ExecProcess、ExecBackend、StartedExecProcess
pub trait ExecProcess: Send + Sync {
fn process_id(&self) -> &ProcessId;
fn subscribe_wake(&self) -> watch::Receiver<u64>;
fn subscribe_events(&self) -> ExecProcessEventReceiver;
fn read(
&self,
after_seq: Option<u64>,
max_bytes: Option<usize>,
wait_ms: Option<u64>,
) -> ExecProcessFuture<'_, ReadResponse>;
fn write(&self, chunk: Vec<u8>) -> ExecProcessFuture<'_, WriteResponse>;
fn signal(&self, signal: ProcessSignal) -> ExecProcessFuture<'_, ()>;
fn terminate(&self) -> ExecProcessFuture<'_, ()>;
}
pub trait ExecBackend: Send + Sync {
fn start(&self, params: ExecParams) -> ExecBackendFuture<'_>;
}2. 启动所有权
2.1 参数与占位
ExecParams 同时描述 child 启动和后续输入能力:tty 决定 PTY 输出,pipe_stdin 决定非 TTY 是否 保留 stdin,cwd、env_policy、sandbox 和 network proxy 由 executor 在启动时解释。process_id 是 调用方选择的逻辑句柄,不是 OS PID。
源码位置:codex-rs/exec-server-protocol/src/protocol.rs :: ExecParams、ExecResponse
pub struct ExecParams {
pub process_id: ProcessId,
pub argv: Vec<String>,
pub cwd: PathUri,
pub env_policy: Option<ExecEnvPolicy>,
pub shell_snapshot: Option<ShellSnapshotRequest>,
pub env: HashMap<String, String>,
pub tty: bool,
pub pipe_stdin: bool,
pub arg0: Option<String>,
pub sandbox: Option<FileSystemSandboxContext>,
pub enforce_managed_network: bool,
pub managed_network: Option<ManagedNetworkSandboxContext>,
pub network_proxy: Option<RemoteNetworkProxyLaunchConfig>,
}
pub struct ExecResponse {
pub process_id: ProcessId,
pub sandbox_type: Option<ProcessSandboxType>,
}LocalProcess::start_process 先把 ProcessEntry::Starting(Arc<ProcessStart>) 写入 map,再调用 codex_sandboxing::spawn_process。如果 spawn 失败,只删除仍然指向同一 ProcessStart 的占位,避免 并发取消误删后来复用相同 ID 的新启动。
源码位置:codex-rs/exec-server/src/local_process.rs :: LocalProcess::start_process
let start = Arc::new(ProcessStart);
{
let mut process_map = self.inner.processes.lock().await;
if process_map.contains_key(&process_id) {
return Err(invalid_request(format!("process {process_id} already exists")));
}
process_map.insert(
process_id.clone(),
ProcessEntry::Starting(Arc::clone(&start)),
);
}
let spawned_result = codex_sandboxing::spawn_process(SpawnRequest {
command: &prepared.command,
cwd: prepared.cwd.as_path(),
env: &prepared.env,
arg0: &prepared.arg0,
sandbox: prepared.sandbox,
windows_sandbox: prepared.windows_sandbox_spawn_request(),
tty: params.tty,
stdin_open: params.tty || params.pipe_stdin,
inherited_fds: &[],
})
.await;2.2 Running转换
spawn 成功后,源码再次取得 map 锁,确认占位仍是同一个 Arc,然后换成 RunningProcess。此时才 初始化 next_seq = 1、retained output、两个输出流计数、事件日志、metrics guard 和可选 network policy shutdown,并启动两个 stream_output task 与一个 watch_exit task。
源码位置:codex-rs/exec-server/src/local_process.rs :: RunningProcess、LocalProcess::start_process
process_map.insert(
process_id.clone(),
ProcessEntry::Running(Box::new(RunningProcess {
session: spawned.session,
tty: params.tty,
pipe_stdin: params.pipe_stdin,
accepted_stdin_write_ids: Arc::new(Mutex::new(
AcceptedStdinWriteIds::default(),
)),
output: VecDeque::new(),
retained_bytes: 0,
next_seq: 1,
exit_code: None,
wake_tx: wake_tx.clone(),
events: events.clone(),
output_notify: Arc::clone(&output_notify),
open_streams: 2,
closed: false,
metrics: Some(self.inner.telemetry.process_started()),
termination_requested: false,
sandbox: prepared.sandbox,
sandbox_denied: false,
network_proxy_handle: prepared.network_proxy_handle,
network_policy_shutdown,
})),
);3. 输出读取
3.1 retained queue
stream_output 对每个 child chunk 在同一把进程表锁下完成四件事:分配递增 seq、写入 retained queue、更新 wake_tx/事件日志,并构造 process/output 通知。队列同时受 1 MiB 字节上限和 50,000 chunk 上限约束;淘汰只影响迟到的 read/replay,不影响已经发送的 live 通知。
源码位置:codex-rs/exec-server/src/local_process.rs :: stream_output
let seq = process.next_seq;
process.next_seq += 1;
process.retained_bytes += chunk.len();
process.output.push_back(RetainedOutputChunk {
seq,
stream,
chunk: chunk.clone(),
});
while process.retained_bytes > RETAINED_OUTPUT_BYTES_PER_PROCESS
|| process.output.len() > RETAINED_OUTPUT_CHUNKS_PER_PROCESS
{
let Some(evicted) = process.output.pop_front() else { break };
process.retained_bytes = process
.retained_bytes
.saturating_sub(evicted.chunk.len());
}
let _ = process.wake_tx.send(seq);3.2 read等待
exec_read 从 after_seq 之后扫描 retained queue。max_bytes 限制本次响应,wait_ms 只在没有 新 chunk、没有 closed 且没有终态事件时等待 output_notify;超时返回当前快照,而不是错误。没有 指定 max_bytes 时,next_seq 保持进程当前序号,避免把分页语义误解为“只返回一页”。
源码位置:codex-rs/exec-server/src/local_process.rs :: LocalProcess::exec_read
let after_seq = params.after_seq.unwrap_or(0);
let max_bytes = params.max_bytes.unwrap_or(usize::MAX);
let wait = Duration::from_millis(params.wait_ms.unwrap_or(0));
let deadline = tokio::time::Instant::now() + wait;
for retained in process.output.iter().filter(|chunk| chunk.seq > after_seq) {
let chunk_len = retained.chunk.len();
if !chunks.is_empty() && total_bytes + chunk_len > max_bytes {
break;
}
total_bytes += chunk_len;
chunks.push(ProcessOutputChunk {
seq: retained.seq,
stream: retained.stream,
chunk: retained.chunk.clone().into(),
});
}3.3 event replay
进程事件有四种:Output 携带 stdout/stderr/pty chunk,Exited 携带退出码和 sandbox 判断,Closed 表示所有输出流都结束,Failed 表示 client 侧无法继续维持 process session。ExecProcessEventLog 新 订阅者先得到 bounded replay,再接入 broadcast live channel;live receiver 若 Lagged,应回到 read(last_seq, ...) 补齐,而不是假设事件永不丢失。
源码位置:codex-rs/exec-server/src/process.rs :: ExecProcessEvent、ExecProcessEventLog
pub enum ExecProcessEvent {
Output(ProcessOutputChunk),
Exited {
seq: u64,
exit_code: i32,
sandbox_denied: Option<bool>,
},
Closed { seq: u64 },
Failed(String),
}
pub fn subscribe(&self) -> ExecProcessEventReceiver {
let history = self.inner.history.lock().unwrap();
let live_rx = self.inner.live_tx.subscribe();
let replay = history.events.iter().cloned().collect();
ExecProcessEventReceiver { replay, live_rx, _keepalive: None }
}4. stdin幂等
4.1 状态返回
process/write 先拒绝空 writeId,然后在进程表锁下区分 UnknownProcess、Starting 和 StdinClosed。 只有 Running 且 TTY 或 pipe_stdin 开启时,才进入 writer channel。WriteStatus 是业务响应,不是 RPC error;调用方可以据此决定等待启动、停止写入或报告句柄失效。
源码位置:codex-rs/exec-server/src/protocol.rs :: WriteParams、WriteStatus、WriteResponse
pub struct WriteParams {
pub process_id: ProcessId,
pub chunk: ByteChunk,
pub write_id: String,
}
pub enum WriteStatus {
Accepted,
UnknownProcess,
StdinClosed,
Starting,
}4.2 reserve再记忆
Running 进程为已接受的 ID 保存一个最多 4096 项的集合。实现先快速检查,再 reserve writer channel, 再次持锁检查竞态,最后同步 permit.send 并立即 remember。记录发生在任何后续 await 之前,因此 handler 被取消后重试同一个 ID,只会得到 Accepted,不会把同一字节写入 child 两次。
源码位置:codex-rs/exec-server/src/local_process.rs :: LocalProcess::exec_write、AcceptedStdinWriteIds
if accepted_stdin_write_ids.lock().await.contains(¶ms.write_id) {
return Ok(WriteResponse { status: WriteStatus::Accepted });
}
let permit = writer_tx
.reserve()
.await
.map_err(|_| internal_error("failed to write to process stdin".to_string()))?;
let mut accepted_stdin_write_ids = accepted_stdin_write_ids.lock().await;
if accepted_stdin_write_ids.contains(¶ms.write_id) {
return Ok(WriteResponse { status: WriteStatus::Accepted });
}
permit.send(params.chunk.into_inner());
accepted_stdin_write_ids.remember(params.write_id);
Ok(WriteResponse { status: WriteStatus::Accepted })5. 终止与清理
5.1 信号与终止
signal_process 对 Starting 或缺失进程返回空 SignalResponse,对已退出 Running 进程也不重复发送 signal;只有仍在运行的 session 才调用 PTY 的 Interrupt。terminate_process 则取消 network policy, 标记 termination_requested 并终止 child;Starting 占位会直接移除,已退出进程返回 running: false。
源码位置:codex-rs/exec-server/src/local_process.rs :: signal_process、terminate_process
match process_map.get(¶ms.process_id) {
Some(ProcessEntry::Running(process)) => {
if process.exit_code.is_some() {
return Ok(SignalResponse {});
}
process
.session
.signal(pty_process_signal(params.signal))
.map_err(|err| internal_error(format!("failed to signal process: {err}")))?;
}
Some(ProcessEntry::Starting(_)) | None => {}
}5.2 Exited到Closed
watch_exit 先记录退出码并完成 metrics;sandbox 进程还会等待很短窗口,让迟到的输出进入 retained queue。随后发布 Exited。两个输出 receiver 都结束后,maybe_emit_closed 才把 closed 置真、取消 network policy、关闭 proxy、发布 Closed,并在测试 25ms、生产 30s 的 retention 后删除进程表条目。
源码位置:codex-rs/exec-server/src/local_process.rs :: watch_exit、finish_output_stream、maybe_emit_closed
if process.closed || process.open_streams != 0 || process.exit_code.is_none() {
return;
}
process.closed = true;
if let Some(network_policy_shutdown) = process.network_policy_shutdown.take() {
network_policy_shutdown.cancel();
}
let seq = process.next_seq;
process.next_seq += 1;
process.events.publish(ExecProcessEvent::Closed { seq });5.3 本地与远程
RemoteProcess 不复制本地状态机,而是把 ExecProcess 方法委托给 Session。session 负责通过 LazyRemoteExecServerClient 建立 RPC,断连时可以恢复;RemoteExecProcess::Drop 取消 network policy 决策并异步 unregister。这样本地 child 的 owner 仍在 executor,client 侧只持有可恢复的逻辑句柄。
源码位置:codex-rs/exec-server/src/remote_process.rs :: RemoteProcess::start、RemoteExecProcess、Drop for RemoteExecProcess
async fn start(
&self,
params: ExecParams,
network_policy_decider: Option<Arc<dyn NetworkPolicyDecider>>,
) -> Result<StartedExecProcess, crate::ExecServerError> {
let client = self.client.get().await?;
let session = client.start_process(params, network_policy_decider).await?;
let sandbox_type = sandbox_type_from_protocol(session.sandbox_type());
Ok(StartedExecProcess {
process: Arc::new(RemoteExecProcess { session }),
sandbox_type,
})
}6. 源码验证
端到端进程测试通过 ExecBackend 同时覆盖 local 和 remote 参数化路径。启动/退出测试检查 process handle 和 exit 状态;输出测试检查 read 的 seq、stdout/stderr 和 Closed;写入测试先启动带 stdin 的 shell,再写入并从 read 取回;无 pipe_stdin 测试断言 StdinClosed;signal 测试验证 interrupt 能让 进程离开运行态。
源码位置:codex-rs/exec-server/tests/exec_process.rs :: exec_process_starts_and_exits、exec_process_streams_output、exec_process_write_then_read、exec_process_rejects_write_without_pipe_stdin、exec_process_signal_interrupts_process
cd codex-rs
cargo test -p codex-exec-server --test exec_process exec_process_starts_and_exits -- --test-threads=1
cargo test -p codex-exec-server --test exec_process exec_process_streams_output -- --test-threads=1
cargo test -p codex-exec-server --test exec_process exec_process_write_then_read -- --test-threads=1
cargo test -p codex-exec-server --test exec_process exec_process_rejects_write_without_pipe_stdin -- --test-threads=1
cargo test -p codex-exec-server --test exec_process exec_process_signal_interrupts_process -- --test-threads=1事件日志的单测使用 8 个事件、3 字节 retained budget,发布一个更大的 Output 后再发布 Exited 和 Closed,断言迟到订阅者仍能看到终态;这证明 replay 有界且终态不依赖输出字节预算。
这些测试不证明每一种平台 sandbox、network proxy 决策或 Windows-only 终止分支都已执行;它们证明的是 进程 RPC 在 local/remote backend 之间共享的字段和生命周期 contract。
源码位置:codex-rs/exec-server/src/process.rs :: event_history_replay_is_bounded_by_retained_bytes
cargo test -p codex-exec-server --lib process::tests::event_history_replay_is_bounded_by_retained_bytes -- --test-threads=1沿源码排查“输出不见了”时,先确认 stream_output 是否为该 process_id 分配了 seq,再检查 retained queue 是否因字节或 chunk 上限淘汰,最后判断消费者是应继续 read(after_seq) 还是处理 broadcast 的 Lagged。排查“写入重复”时,检查同一 write_id 是否在 reserve 前后两次命中集合,以及进程是否 真的允许 stdin;排查“终止后仍占用句柄”时,沿 watch_exit → maybe_emit_closed → retention cleanup 检查 exit_code、open_streams 和 closed 三个字段。
下一篇ExecServer文件系统RPC会沿用相同的 handler/session/RPC 分层, 分析 PathUri、sandbox context 和文件句柄的生命周期。
