Skip to content

ExecServer架构

沿真实源码解释 Environment统一执行面、远程client恢复、exec-server连接分派、session registry、进程owner与资源清理边界。

基于rust-v0.150.0
CodexRustExecutionExecServer

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

rust
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

rust
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

rust
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

rust
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、EnvironmentCapabilities
  • codex-rs/exec-server/src/client.rs :: ExecServerClient::environment_info
rust
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

rust
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_websocket
  • codex-rs/exec-server/src/client/accepted.rs :: Inner::accept_replacement_connection
rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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.rs
  • codex-rs/exec-server/src/server/request_dispatcher_tests.rs
  • codex-rs/exec-server/src/server/transport_tests.rs
  • codex-rs/exec-server/src/server.rs :: telemetry_entrypoint_emits_root_span
text
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=1

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

text
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=1

8.3 Accepted替换 ​

accepted WebSocket 的 5 项测试覆盖初始 resume 拒绝、initialize metadata、立即 ready、重叠/失败 replacement,以及 replacement 后恢复运行中 process 和 output。

源码位置:codex-rs/exec-server/tests/accepted_websocket.rs

text
cd codex-rs
cargo test -p codex-exec-server --test accepted_websocket -- --test-threads=1

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

text
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. 源码排查 ​

text
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/server

Exec-server 的主线是:EnvironmentManager 为 Core 选择统一能力;远程 Environment 通过 lazy client 建立或恢复连接;server processor 为每条 transport 建立 handler 和 dispatcher;router 把方法投影到 connection-local 资源或跨连接 process session;断连时关闭连接资源并 detach session,TTL 内可由新 transport 恢复。下一篇将展开 initialize 握手与连接状态。