ExecServer架构
ExecServer crate 不是一个只负责“转发 JSON-RPC”的薄壳,也不只是一支 server binary。对 Core 来说,它首先提供 EnvironmentManager 与 Environment,把本地和远程的进程、文件系统与 HTTP 能力统一起来;远程分支才进一步使用 lazy/recoverable client、JSON-RPC transport 和独立 server 进程。服务端连接内部仍由 router、dispatcher、handler 与跨连接 SessionRegistry 协作。
本文面向已读过Codex执行体系总览、非交互Exec CLI和Unified Exec创建进程并理解 JSON-RPC、Tokio task 的读者。范围是 exec-server crate 的总体所有权和连接生命周期,不展开握手字段、全部消息模型或单个进程 RPC。读完后,你应能从 Core 的 Environment 走到 local/remote backend,再从远程 client 走到 server handler,并指出连接、session 与 process 各自的 owner。
1. 执行环境统一面
1.1 Manager选择环境
EnvironmentManager 保存 environment ID 到稳定 Arc<Environment> 的映射,还单独记录 default 与 local environment。provider snapshot 可以同时配置多个远程 transport 和可选 local 环境;构造完成后远程环境在后台开始连接,但真正使用仍由后续选择触发。
源码位置:codex-rs/exec-server/src/environment.rs :: EnvironmentManager::from_snapshot
let local_environment = if include_local {
let local_runtime_paths = local_runtime_paths.clone().ok_or_else(|| {
ExecServerError::Protocol(
"local environment requires configured runtime paths".to_string(),
)
})?;
let local_environment = Arc::new(Environment::local(
local_runtime_paths,
http_client_factory.clone(),
));
environment_map.insert(
LOCAL_ENVIRONMENT_ID.to_string(),
Arc::clone(&local_environment),
);
Some(local_environment)
} else {
None
};
for (id, transport) in environments {
let environment = Environment::remote_with_transport(
transport,
/*local_runtime_paths*/ None,
http_client_factory.clone(),
);
environment_map.insert(id, Arc::new(environment));
}
for environment in environment_map.values() {
environment.start_connecting();
}1.2 环境能力组合
本地 Environment 直接组合 LocalProcess、LocalFileSystem 与 route-aware HTTP client;远程 Environment 用同一个 lazy client 分别构造 RemoteProcess、RemoteFileSystem 和 HTTP client。Core 因此依赖 trait,而不是在每个调用点判断 transport。
源码位置:codex-rs/exec-server/src/environment.rs :: Environment::local、remote_with_client
pub(crate) fn remote_with_client(
client: LazyRemoteExecServerClient,
local_runtime_paths: Option<ExecServerRuntimePaths>,
) -> Self {
let exec_backend: Arc<dyn ExecBackend> = Arc::new(RemoteProcess::new(client.clone()));
let filesystem: Arc<dyn ExecutorFileSystem> =
Arc::new(RemoteFileSystem::new(client.clone()));
Self {
remote_client: Some(client.clone()),
ready_info: Arc::new(ArcSwapOption::empty()),
provisioning_status_tx: None,
startup_task: Arc::new(Mutex::new(None)),
exec_backend,
filesystem,
http_client: Arc::new(client),
local_runtime_paths,
}
}2. Server入口
2.1 Crate边界
exec-server/src/lib.rs 把 client、environment、process、filesystem、server、transport 和 telemetry 分成模块,并只导出跨 crate 的稳定接口。Core 使用 ExecServerClient、Environment 和 ExecBackend;服务端启动使用 run_main。
源码位置:codex-rs/exec-server/src/lib.rs :: module declarations、run_main
mod client;
mod client_transport;
mod environment;
mod environment_config;
mod local_process;
mod process;
mod remote_process;
mod server;
#[cfg(unix)]
mod shell_snapshot;
mod telemetry;
pub use client::ExecServerClient;
pub use client::ExecServerError;
pub use environment::Environment;
pub use environment::EnvironmentManager;
pub use process::ExecBackend;
pub use process::ExecProcess;
pub use server::run_main;模块声明是内部实现边界;pub use 才是 Core 或 CLI 能依赖的 crate API。
2.2 Telemetry入口
run_main_with_telemetry 只负责建立根 span 并把监听地址、runtime paths、HTTP client、dispatch mode 交给 transport。它不创建某个具体进程。
源码位置:codex-rs/exec-server/src/server.rs :: run_main_with_telemetry
pub async fn run_main_with_telemetry(
listen_url: &str,
runtime_paths: ExecServerRuntimePaths,
telemetry: ExecServerTelemetry,
http_client_factory: HttpClientFactory,
request_dispatch_mode: RequestDispatchMode,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
transport::run_transport(
listen_url,
runtime_paths,
telemetry,
http_client_factory,
request_dispatch_mode,
)
.await
}3. Client与恢复
3.1 初始化元数据
initialize response 现在可以直接携带 EnvironmentInfo。它包括 executor shell、cwd、临时目录和 capability flags;旧 server 不返回时,client 延迟调用 environment/info。结果保存在 client lifetime 的 OnceCell 中,后续 clone 不重复 RPC。
源码位置:
codex-rs/exec-server-protocol/src/protocol.rs :: InitializeResponse、EnvironmentInfo、EnvironmentCapabilitiescodex-rs/exec-server/src/client.rs :: ExecServerClient::environment_info
pub struct InitializeResponse {
pub session_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub environment_info: Option<EnvironmentInfo>,
}
pub struct EnvironmentCapabilities {
pub network_proxy_launch: bool,
pub capability_discovery_sandbox: bool,
pub environment_config_read: bool,
pub http_header_env_vars: bool,
pub sandboxed_file_streaming: bool,
pub shell_snapshot_v2: bool,
}新请求字段必须由 capability gate 决定,不能用版本号猜测。例如 ShellSnapshotV2、executor-local network proxy 和 sandboxed file stream 都有独立 flag。
3.2 Lazy重连
LazyRemoteExecServerClient 缓存首次连接结果与最新成功 client。WebSocket、Deferred 和 Noise transport 支持重连;stdio process 不自动重连。相同 outage 的调用者共享一个 ConnectionAttempt,失败完成后清空 attempt,使未来操作仍可重试。
源码位置:codex-rs/exec-server/src/client.rs :: LazyRemoteExecServerClient::initial_client、reconnect、can_reconnect
async fn reconnect(&self) -> Result<ExecServerClient, ExecServerError> {
let attempt = {
let mut reconnect = self
.reconnect
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(client) = self.connected_client() {
return Ok(client);
}
reconnect
.get_or_insert_with(|| Arc::new(ConnectionAttempt::new()))
.clone()
};
let result = attempt
.get_or_init(|| async {
let result = self.connect_once().await;
if let Ok(client) = &result {
*self
.current_client
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(client.clone());
}
result
})
.await;
let mut reconnect = self
.reconnect
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if reconnect
.as_ref()
.is_some_and(|current| Arc::ptr_eq(current, &attempt))
{
*reconnect = None;
}
result.clone().map_err(ExecServerError::ConnectionAttempt)
}3.3 Accepted连接
accepted 模式用于宿主已经完成 WebSocket 接受与认证的场景。EnvironmentManager::from_accepted_websocket 把初始 socket 包装成 remote Environment;替换时宿主只交付新 socket,client 负责关闭旧 transport、用保存的 session ID initialize/resume,并恢复既有 process。
源码位置:
codex-rs/exec-server/src/environment/accepted.rs :: EnvironmentManager::from_accepted_websocket、replace_accepted_websocketcodex-rs/exec-server/src/client/accepted.rs :: Inner::accept_replacement_connection
let client = ExecServerClient::connect_accepted_websocket(websocket, options).await?;
let client = LazyRemoteExecServerClient::from_connected(client, http_client_factory.clone());
let environment = Arc::new(Environment::remote_with_client(
client, /*local_runtime_paths*/ None,
));4. 连接所有权
4.1 连接处理器
processor 跨连接持有 SessionRegistry、runtime paths、telemetry、HTTP factory 和 dispatch mode。每个连接调用 run_connection,但 registry 是共享的,允许环境和 session 在连接恢复时继续存在。
源码位置:codex-rs/exec-server/src/server/processor.rs :: ConnectionProcessor
pub(crate) struct ConnectionProcessor {
session_registry: Arc<SessionRegistry>,
runtime_paths: ExecServerRuntimePaths,
telemetry: ExecServerTelemetry,
http_client_factory: HttpClientFactory,
request_dispatch_mode: RequestDispatchMode,
}
pub(crate) async fn run_connection(
&self,
connection: JsonRpcConnection,
transport: ConnectionTransport,
) {
run_connection(
connection,
Arc::clone(&self.session_registry),
self.runtime_paths.clone(),
self.telemetry.clone(),
self.http_client_factory.clone(),
transport,
self.request_dispatch_mode,
).await;
}4.2 连接Handler
每条连接创建一个 handler、outbound channel 和 dispatcher。handler 不是全局单例,因为它绑定该连接的通知发送器和 request sender;但它访问的 session registry 来自 processor 共享实例。
源码位置:codex-rs/exec-server/src/server/processor.rs :: run_connection
let (outgoing_tx, mut outgoing_rx) =
mpsc::channel::<RpcServerOutboundMessage>(CHANNEL_CAPACITY);
let notifications = RpcNotificationSender::new(outgoing_tx.clone());
let requests = notifications.request_sender();
let handler = Arc::new(ExecServerHandler::new(
session_registry,
notifications,
runtime_paths,
http_client_factory,
));
let mut dispatcher = RequestDispatcher::new(
Arc::new(build_router()),
Arc::clone(&handler),
outgoing_tx.clone(),
disconnected_rx.clone(),
requests.clone(),
telemetry,
request_dispatch_mode,
);5. 请求路由
5.1 Router注册
build_router 把 initialize、environment、exec、filesystem、HTTP 和 process 控制方法绑定到 handler。Router 只选择函数,不持有进程 map;具体资源由 handler 下沉到 registry 或 backend。
源码位置:codex-rs/exec-server/src/server/registry.rs :: build_router
let mut router = RpcRouter::new();
router.request(
INITIALIZE_METHOD,
|handler: Arc<ExecServerHandler>, params: InitializeParams| async move {
handler.initialize(params).await
},
);
router.request(
EXEC_METHOD,
|handler: Arc<ExecServerHandler>, params: ExecParams| async move {
handler.exec(params).await
},
);
router.request(
FS_READ_FILE_METHOD,
|handler: Arc<ExecServerHandler>, params: FsReadFileParams| async move {
handler.fs_read_file(params).await
},
);5.2 Dispatch模式
dispatcher 在 inline 模式下按输入顺序等待 route;concurrent 模式使用 ordinary/control lanes 和 semaphore。initialize 完成前,即使启用并发,也不会让普通请求观察到未初始化 session。
源码位置:codex-rs/exec-server/src/server/request_dispatcher.rs :: RequestDispatcher::new、dispatch_request
let lanes = match mode {
RequestDispatchMode::Inline => None,
RequestDispatchMode::Concurrent { max_concurrent_requests } => {
Some(RequestLanes {
ordinary: Arc::new(Semaphore::new(max_concurrent_requests.get())),
control: Arc::new(Semaphore::new(max_concurrent_requests.get())),
})
}
};6. Session与Process
6.1 SessionRegistry
SessionRegistry 跨连接保存 SessionEntry,每个 entry 拥有 ProcessHandler 和连接 attachment 状态。新 initialize 创建 session;resume 只能附着一个已 detach 且未过期的 session。断连 detach 后生产环境保留 30 秒,超时才 shutdown process。
源码位置:codex-rs/exec-server/src/server/session_registry.rs :: SessionRegistry::attach
if let Some(session_id) = resume_session_id {
let entry = sessions
.get(&session_id)
.cloned()
.ok_or_else(|| invalid_request(format!("unknown session id {session_id}")))?;
if entry.is_expired(tokio::time::Instant::now()) {
let entry = sessions.remove(&session_id).ok_or_else(|| {
invalid_request(format!("unknown session id {session_id}"))
})?;
Ok(AttachOutcome::Expired { session_id, entry })
} else if entry.has_active_connection() {
Err(session_already_attached(format!(
"session {session_id} is already attached to another connection"
)))
} else {
entry.process.set_notification_sender(Some(notifications));
entry.attach(connection_id);
Ok(AttachOutcome::Attached(entry))
}
}detach 会把 notification sender 清空并安排 TTL expiration;resume attach 新 sender 后,process output/read 继续使用原 ProcessHandler。只有 ConnectionProcessor::shutdown 才会清空整个 registry 并终止所有 session。
6.2 Handler下沉
handler 接收 router 的 typed params,调用对应的 process handler、file system handler 或 environment provider。该分层让请求协议和资源所有权分开:router 可替换,registry 仍能维持跨连接状态。
源码位置:codex-rs/exec-server/src/server/handler.rs :: ExecServerHandler
pub(crate) struct ExecServerHandler {
session_registry: Arc<SessionRegistry>,
notifications: RpcNotificationSender,
session: StdMutex<Option<SessionHandle>>,
active_body_stream_ids: Mutex<HashSet<String>>,
background_task_shutdown: CancellationToken,
background_tasks: TaskTracker,
file_system: FileSystemHandler,
runtime_paths: ExecServerRuntimePaths,
http_client: RouteAwareHttpClient,
initialize_requested: AtomicBool,
initialized: AtomicBool,
}initialize 通过 registry attach SessionHandle;process RPC 使用 session.process(),filesystem 与 HTTP 则使用 handler 自己的 connection-local owner。这使 process 可以跨 transport 恢复,而打开的文件 stream 和 HTTP stream 随连接关闭。
7. 断连清理
连接循环同时监听 request task 完成、transport disconnected 和 incoming message。断连后先完成已经排队的 client response,关闭 request sender,等待 dispatcher 与 connection-local handler 收尾,再终止 transport task。handler.shutdown() 会关闭 background/file resources并 detach session;process session 暂留 registry 的 TTL 内等待 resume。
源码位置:codex-rs/exec-server/src/server/processor.rs :: run_connection
if *disconnected_rx.borrow() {
complete_queued_client_responses(&requests, &mut incoming_rx);
}
requests.close();
dispatcher.shutdown().await;
handler.shutdown().await;
drop(handler);
drop(requests);
drop(outgoing_tx);
for task in connection_tasks {
task.abort();
let _ = task.await;
}
let _ = outbound_task.await;8. 验证
8.1 Server门槛
initialize 的 2 项测试验证 session ID 与 initialized 前请求拒绝;dispatcher 的 4 项测试覆盖并发上限解析、request span 和 semaphore admission;transport 的 7 项测试覆盖 stdio/websocket URL 解析与 stdio initialize。telemetry 入口还断言生成 codex.exec_server 根 span。
源码位置:
codex-rs/exec-server/tests/initialize.rscodex-rs/exec-server/src/server/request_dispatcher_tests.rscodex-rs/exec-server/src/server/transport_tests.rscodex-rs/exec-server/src/server.rs :: telemetry_entrypoint_emits_root_span
cd codex-rs
cargo test -p codex-exec-server --test initialize -- --test-threads=1
cargo test -p codex-exec-server --lib 'server::request_dispatcher::tests::' -- --test-threads=1
cargo test -p codex-exec-server --lib 'server::transport::transport_tests::' -- --test-threads=1
cargo test -p codex-exec-server --lib telemetry_entrypoint_emits_root_span -- --test-threads=18.2 Metadata与重试
environment_info_is_cached 分别模拟 initialize 已携带 metadata 和 legacy server 延迟 environment/info,断言连接关闭后 clone 仍使用缓存。retryable_startup_failure_does_not_burn_environment 先让 WebSocket handshake 返回 500,再提供 replacement server,断言后续 readiness/get 共享恢复后的 client。
源码位置:codex-rs/exec-server/src/client.rs :: environment_info_is_cached、retryable_startup_failure_does_not_burn_environment
cd codex-rs
cargo test -p codex-exec-server --lib environment_info_is_cached -- --test-threads=1
cargo test -p codex-exec-server --lib retryable_startup_failure_does_not_burn_environment -- --test-threads=18.3 Accepted替换
accepted WebSocket 的 5 项测试覆盖初始 resume 拒绝、initialize metadata、立即 ready、重叠/失败 replacement,以及 replacement 后恢复运行中 process 和 output。
源码位置:codex-rs/exec-server/tests/accepted_websocket.rs
cd codex-rs
cargo test -p codex-exec-server --test accepted_websocket -- --test-threads=18.4 Session与Process
handler 的 5 项测试覆盖重复 process ID 竞争、退出后 terminate、resume 使旧 long poll 失败、active session resume 拒绝,以及 notification receiver 关闭后 output/exit 仍可 read。这些断言直接体现 registry 与 connection owner 的分离。
源码位置:codex-rs/exec-server/src/server/handler/tests.rs
cd codex-rs
cargo test -p codex-exec-server --lib 'server::handler::tests::' -- --test-threads=1这些测试不覆盖所有 filesystem、HTTP、Noise registry、shell snapshot 或 managed-network 业务路径;它们验证的是总体 owner、连接恢复和 dispatch 边界。
9. 源码排查
rg -n "EnvironmentManager|remote_with_client|environment_info" codex-rs/exec-server/src
rg -n "LazyRemoteExecServerClient|reconnect|accepted_websocket" codex-rs/exec-server/src
rg -n "build_router|RequestDispatcher|SessionRegistry|handler.shutdown" codex-rs/exec-server/src/serverExec-server 的主线是:EnvironmentManager 为 Core 选择统一能力;远程 Environment 通过 lazy client 建立或恢复连接;server processor 为每条 transport 建立 handler 和 dispatcher;router 把方法投影到 connection-local 资源或跨连接 process session;断连时关闭连接资源并 detach session,TTL 内可由新 transport 恢复。下一篇将展开 initialize 握手与连接状态。
