Skip to content

ExecServer连接与握手

追踪 exec-server 的 stdio/WebSocket transport、initialize/initialized 两阶段握手、session attach/resume 与断连保活。

基于rust-v0.150.0
CodexRustExecutionExecServer

ExecServer连接与握手 ​

exec-server 的连接状态不是“socket 打开就能执行”。transport 建立 JsonRpcConnection 后,dispatcher 先要求 initialize,handler 再把连接 attach 到 SessionRegistry;收到 initialized 后才允许 exec、filesystem 和 environment 请求。v0.150.0 的 initialize response 还可以携带 executor metadata 与 capability flags,远程 client 会缓存这些信息并据此决定是否发送新字段。连接断开时,session 进入 detached 状态,在 TTL 内可以由另一条连接 resume,超过 TTL 才销毁资源。

本文承接ExecServer架构,面向理解 JSON-RPC、Tokio watch 和连接生命周期的读者。范围是 transport、两阶段握手、metadata capability、session attach/resume、lazy recovery 和 accepted WebSocket handoff,不展开完整消息字段或 process RPC。读完后,你应能解释为什么握手前请求被拒绝、为什么 stdio 不能 resume、为什么旧 server 仍可被兼容,以及 detached session 何时真正清理。

1. Transport入口 ​

1.1 地址解析 ​

服务端只接受 stdio、stdio:// 或 ws://IP:PORT。WebSocket listener 负责接受连接,stdio 只服务一条连接。

源码位置:codex-rs/exec-server/src/server/transport.rs :: parse_listen_url、run_transport

rust
pub(crate) fn parse_listen_url(
    listen_url: &str,
) -> Result<ExecServerListenTransport, ExecServerListenUrlParseError> {
    if matches!(listen_url, "stdio" | "stdio://") {
        return Ok(ExecServerListenTransport::Stdio);
    }
    if let Some(socket_addr) = listen_url.strip_prefix("ws://") {
        return socket_addr
            .parse::<SocketAddr>()
            .map(ExecServerListenTransport::WebSocket)
            .map_err(|_| ExecServerListenUrlParseError::InvalidWebSocketListenUrl(
                listen_url.to_string(),
            ));
    }
    Err(ExecServerListenUrlParseError::UnsupportedListenUrl(
        listen_url.to_string(),
    ))
}

1.2 connection建立 ​

stdio 和 WebSocket 最终都交给同一个 ConnectionProcessor::run_connection;差异只保留在 JsonRpcConnection 的读写任务和 telemetry transport 标签。

源码位置:codex-rs/exec-server/src/server/transport.rs :: run_stdio_connection_with_io、websocket_upgrade_handler

rust
processor
    .run_connection(
        JsonRpcConnection::from_stdio(reader, writer, "exec-server stdio".to_string()),
        ConnectionTransport::Stdio,
    )
    .await;

WebSocket listener 还拒绝带 Origin header 的请求,避免浏览器跨源页面把本地 exec-server 当成普通 HTTP 服务访问;/readyz 只表示 listener 可达,不表示某个 JSON-RPC session 已 initialize。

源码位置:codex-rs/exec-server/src/server/transport.rs :: reject_requests_with_origin_header、readiness_handler

rust
if request.headers().contains_key(ORIGIN) {
    warn!(
        method = %request.method(),
        uri = %request.uri(),
        "rejecting exec-server websocket listener request with Origin header"
    );
    Err(StatusCode::FORBIDDEN)
} else {
    Ok(next.run(request).await)
}

2. 两阶段握手 ​

2.1 initialize ​

handler 用原子 initialize_requested 保证每条连接只初始化一次;initialize 成功后把 session handle 保存到连接 handler,并返回 session ID。

源码位置:codex-rs/exec-server/src/server/handler.rs :: initialize

rust
if self.initialize_requested.swap(true, Ordering::SeqCst) {
    return Err(invalid_request(
        "initialize may only be sent once per connection".to_string(),
    ));
}
let session = match self.session_registry.attach(
    params.resume_session_id.clone(),
    self.notifications.clone(),
    self.runtime_paths.clone(),
).await {
    Ok(session) => session,
    Err(error) => {
        self.initialize_requested.store(false, Ordering::SeqCst);
        return Err(error);
    }
};
let session_id = session.session_id().to_string();
*self
    .session
    .lock()
    .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(session);
Ok(InitializeResponse {
    session_id,
    environment_info: Some(EnvironmentInfo::local()),
})

initialize 成功时不仅返回 session ID,还返回 EnvironmentInfo。本地实现会报告 shell、cwd、临时目录和 capability;旧 peer 可以省略该字段,client 随后再请求 environment/info。initialize_requested 在 attach 失败时恢复为 false,允许调用方修正 session 参数后重试同一连接。

源码位置:codex-rs/exec-server-protocol/src/protocol.rs :: InitializeResponse、EnvironmentInfo、EnvironmentCapabilities

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

2.2 initialized ​

initialized 是一个 notification,不返回 response;handler 检查 initialize 已请求且 session 仍 attached,才把 initialized 标记为真。dispatcher 也维护自己的 initialized,用于并发请求门槛。

源码位置:codex-rs/exec-server/src/server/handler.rs :: initialized、request_dispatcher.rs :: handle_notification

rust
pub(crate) fn initialized(&self) -> Result<(), String> {
    if !self.initialize_requested.load(Ordering::SeqCst) {
        return Err("received `initialized` notification before `initialize`".into());
    }
    self.require_session_attached().map_err(|error| error.message)?;
    self.initialized.store(true, Ordering::SeqCst);
    Ok(())
}

dispatcher 在 concurrent 模式下还有第二道门:initialize 本身和 initialized 之前的请求保持 inline,只有握手完成后才进入 ordinary/control semaphore。control lane 预留给 environment status、signal、terminate、filesystem close 等清理类操作,避免普通长轮询耗尽所有 permit。

源码位置:codex-rs/exec-server/src/server/request_dispatcher.rs :: dispatch_request

rust
let Some(RequestLanes { ordinary, control }) = &self.lanes else {
    return task.await;
};
if method == INITIALIZE_METHOD || !self.initialized {
    return task.await;
}
let admission = if matches!(
    method,
    ENVIRONMENT_INFO_METHOD
        | ENVIRONMENT_STATUS_METHOD
        | EXEC_SIGNAL_METHOD
        | EXEC_TERMINATE_METHOD
        | FS_CLOSE_METHOD
) {
    Arc::clone(control)
} else {
    Arc::clone(ordinary)
};

3. Session attach ​

3.1 新session ​

没有 resume ID 时,registry 生成 UUID,创建 SessionEntry 和 ProcessHandler,并把 entry 放入 sessions map。后台 process 资源因此归 registry entry 所有,而不是归连接 task 所有。

源码位置:codex-rs/exec-server/src/server/session_registry.rs :: attach

rust
let session_id = Uuid::new_v4().to_string();
let entry = Arc::new(SessionEntry::new(
    session_id.clone(),
    ProcessHandler::new(notifications, self.telemetry.clone(), runtime_paths),
    connection_id,
));
sessions.insert(session_id, Arc::clone(&entry));

3.2 resume ​

带 resume ID 时,registry 要求 session 存在、未过期且没有其他 active connection;成功后替换 notification sender 并更新 connection ID。已被其他连接占用会返回 session-already-attached,而不是创建副本。

源码位置:codex-rs/exec-server/src/server/session_registry.rs :: attach

rust
let entry = sessions
    .get(&session_id)
    .cloned()
    .ok_or_else(|| invalid_request(format!("unknown session id {session_id}")))?;
if entry.is_expired(Instant::now()) {
    return Err(invalid_request(format!("unknown session id {session_id}")));
} else if entry.has_active_connection() {
    return Err(session_already_attached(format!(
        "session {session_id} is already attached to another connection"
    )));
}
entry.process.set_notification_sender(Some(notifications));
entry.attach(connection_id);

resume 的互斥是 connection 级别,不是 process 级别:同一 session 同时被两个 connection attach 时,后者返回 session_already_attached。旧 connection 关闭时,handler 只在 attachment ID 仍匹配时执行 detach,避免旧 transport 的延迟清理误 detach 新 connection。

3.3 Accepted连接 ​

accepted 模式由宿主先接受并认证 WebSocket,再交给 EnvironmentManager。初始 accepted connection 禁止携带 resume ID;替换 socket 时 client 关闭旧 transport,把 replacement 排入专用 channel,恢复任务使用保存的 session ID 完成 initialize/resume。

源码位置:

  • codex-rs/exec-server/src/environment/accepted.rs :: EnvironmentManager::from_accepted_websocket、replace_accepted_websocket
  • codex-rs/exec-server/src/client/accepted.rs :: ExecServerClient::connect_accepted_websocket、Inner::accept_replacement_connection
rust
if options.resume_session_id.is_some() {
    return Err(ExecServerError::Protocol(
        "accepted exec-server initial connection cannot resume a session".to_string(),
    ));
}
let connection_source = AcceptedConnectionSource::new(options.clone());
Self::connect_with_recovery(
    JsonRpcConnection::from_axum_websocket(
        websocket,
        "accepted exec-server websocket".to_string(),
    ),
    options,
    Some(ExecServerReconnectStrategy::Accepted(connection_source)),
)
.await

replacement handoff 同时限制并发替换数量为 1,并确保旧 RPC transport 先进入 recovery 再关闭,避免两个 connection 同时消费同一 session。

4. 断连保活 ​

连接 detach 后,session 不立即删除,而是记录 detached connection 和过期时间;DETACHED_SESSION_TTL 测试环境为 200ms,生产为 30s。TTL 内 resume 可以继续使用后台 process;过期 task 从 registry 移除 entry 并 shutdown process。

源码位置:codex-rs/exec-server/src/server/session_registry.rs :: expire_if_detached

rust
tokio::time::sleep(DETACHED_SESSION_TTL).await;
let removed = {
    let mut sessions = self.sessions.lock().await;
    let Some(entry) = sessions.get(&session_id) else {
        return;
    };
    if !entry.is_detached_connection_expired(connection_id, Instant::now()) {
        return;
    }
    sessions.remove(&session_id)
};
if let Some(entry) = removed {
    entry.process.shutdown().await;
}

远程 client 的 detached 保活与 server registry 的 TTL 是两层机制:server 保留 SessionEntry 和 ProcessHandler;client 侧保存 session/process 事件与 reconnect 状态。transport 断开不等价于 process 已退出,只有恢复失败或 TTL 到期才会进入真正的资源 shutdown。

stdio transport 特别说明:它只服务一个连接,transport 完成后 processor 立即 shutdown,因此 detached session 不能通过 stdio resume。

5. 验证 ​

5.1 Transport与握手 ​

transport 测试覆盖 stdio、stdio://、合法/非法 ws listen URL 和 stdio initialize;initialize 测试断言返回 session ID,并断言 initialized 前发送 environment request 会得到协议错误。它们证明入口与状态门槛,不展开完整字段协议。

源码位置:

  • codex-rs/exec-server/src/server/transport_tests.rs
  • codex-rs/exec-server/tests/initialize.rs
text
cd codex-rs
cargo test -p codex-exec-server --test initialize -- --test-threads=1
cargo test -p codex-exec-server --lib 'server::transport::transport_tests::' -- --test-threads=1

5.2 Metadata与恢复 ​

environment_info_is_cached 覆盖 initialize 带 metadata 与 legacy 延迟 info 两条路径;retryable_startup_failure_does_not_burn_environment 先失败再恢复,断言后续调用共享 replacement client。accepted WebSocket 测试覆盖 initial/replacement socket、metadata、ready 和 process recovery。

源码位置:

  • codex-rs/exec-server/src/client.rs :: environment_info_is_cached、retryable_startup_failure_does_not_burn_environment
  • codex-rs/exec-server/tests/accepted_websocket.rs
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
cargo test -p codex-exec-server --test accepted_websocket -- --test-threads=1

5.3 Session与进程 ​

process 测试断言 detached session 在 TTL 内 resume 不会杀死后台进程;handler 测试覆盖 active resume 拒绝、重复 process ID、长 poll 在 resume 后失败、退出后 terminate 和 notification receiver 关闭后的 output/exit 保留。

源码位置:

  • codex-rs/exec-server/tests/process.rs :: exec_server_resumes_detached_session_without_killing_processes
  • codex-rs/exec-server/src/server/handler/tests.rs
text
cd codex-rs
cargo test -p codex-exec-server --test process exec_server_resumes_detached_session_without_killing_processes -- --test-threads=1
cargo test -p codex-exec-server --lib 'server::handler::tests::' -- --test-threads=1

6. 源码排查 ​

text
rg -n "parse_listen_url|run_stdio_connection|websocket_upgrade_handler" codex-rs/exec-server/src/server/transport.rs
rg -n "initialize_requested|initialized\(|require_initialized_for" codex-rs/exec-server/src/server
rg -n "resume_session_id|DETACHED_SESSION_TTL|expire_if_detached" codex-rs/exec-server/src/server/session_registry.rs

连接主线是:transport 建立字节通道,initialize attach session,initialized 打开业务请求,断连进入 detached TTL;WebSocket session 可在 TTL 内 resume,stdio 连接结束后直接 shutdown。下一篇将展开 ExecServer 的 JSON-RPC request/response/event 消息模型。