Skip to content

CodeModeHost生命周期

追踪Code Mode process/WebSocket与gRPC host的连接复用、session generation、cell/delegate路由、失败和关闭清理。

基于rust-v0.150.0
CodexRustExecutionCodeMode

CodeModeHost生命周期 ​

Code Mode host 生命周期不等于一个 socket 的打开和关闭。process/WebSocket 路径由共享 OwnedCodeModeHost 复用 physical connection,每个 logical session 在 connection driver 中拥有 generation、 cell registry、pending operation 和 delegate tasks;host 端另有 HostPeer、session runtime registry 和 request tracker。gRPC 路径则把 session 绑定到 lease stream,用 ReconnectableSession 在 host 重启后创建 新 generation。任何一层失败都必须唤醒 caller、撤销 callback、关闭 cell 并阻止迟到消息进入新状态。

本文承接CodeMode协议与CodeMode架构总览。范围是 provider、connection、session binding、cell/delegate 路由和 shutdown;不展开 V8 cell 内部算法。读完后, 你应能区分 connection 重建与 session shutdown,解释 stale generation 如何被拒绝,并从一次断连追踪到 pending request、delegate callback 和 cell closure 的清理路径。

1. 生命周期层次 ​

process/WebSocket 与 gRPC 共享 domain CodeModeSession,但生命周期 owner 不同:

层次process/WebSocketgRPC
provider ownerOwnedCodeModeHostSharedTransport
logical sessionProcessOwnedCodeModeSessionReconnectableSession
physical bindingV1 Connection + RemoteSessionGrpcCodeModeSession + lease
generationRemoteSession.generationSessionBinding.generation
host session ownerHostState.sessionsgRPC host lease registry
callback routeV1 delegate messagessession lease + tool subscription
disconnect signaldriver cancellation/failurestopped token / stream failure

2. Provider复用 ​

2.1 Physical连接 ​

OwnedCodeModeHost 缓存一条 process 或 WebSocket Connection。connection() 先无锁语义地读取 live cache,再用单 permit 串行协调创建;获取 permit 后重新检查,避免并发调用重复 spawn host。旧 connection 死亡后不再被 live_connection() 返回,下一次使用建立新连接。

源码位置:codex-rs/code-mode/src/remote_session.rs :: OwnedCodeModeHost::connection

rust
async fn connection(&self) -> Result<Arc<Connection>, ConnectionError> {
    if let Some(connection) = self.live_connection() {
        return Ok(connection);
    }
    let _connect_permit = self.connect_permit.acquire().await.map_err(|_| {
        ConnectionError::Other("code-mode host connection coordinator closed".into())
    })?;
    if let Some(connection) = self.live_connection() {
        return Ok(connection);
    }
    let new_connection = match &self.endpoint {
        HostEndpoint::Process(host_program) => Connection::spawn(host_program).await?,
        HostEndpoint::WebSocket { websocket_url, http_client_factory } => {
            Connection::connect_websocket(websocket_url, http_client_factory).await?
        }
    };
    let new_connection = Arc::new(new_connection);
    *self.connection.lock().unwrap() = Some(Arc::clone(&new_connection));
    Ok(new_connection)
}

provider 共享 physical host,但逻辑 session 不共享 store。每次 create_host_session 创建新的 ProcessOwnedCodeModeSession,并在返回前确保 binding 已打开。

2.2 Session状态 ​

process/WebSocket session 使用 New → Opening → Open → Closing → Closed。Opening 保存 remote session identity 和 watch receiver,多个并发 operation 等待同一个 open 结果。RemoteSession 含 session ID 与 generation;driver registry 会比较完整 identity,旧 generation 即使 session ID 相同也被拒绝。

源码位置:codex-rs/code-mode/src/remote_session.rs :: SessionState、SessionInner

rust
enum SessionState {
    New,
    Opening {
        remote: RemoteSession,
        result_rx: watch::Receiver<Option<Result<SessionBinding, String>>>,
    },
    Open(SessionBinding),
    Closing,
    Closed,
}

struct SessionInner {
    host: Arc<OwnedCodeModeHost>,
    delegate: Arc<dyn CodeModeSessionDelegate>,
    limits: CodeModeSessionCellExecutionLimits,
    state: StdMutex<SessionState>,
    next_generation: AtomicU64,
    shutdown_requested: AtomicBool,
    shutdown_result: StdMutex<Option<ShutdownResultReceiver>>,
    retired_cleanups: StdMutex<Vec<SessionCleanup>>,
}

3. V1 Host连接 ​

3.1 握手与双lane ​

host run_connection 先协商 Protocol V1/capabilities。若选择 dual WebSocket,必须在 10 秒内收到 pairing connection;随后创建 control outgoing channel,以及容量与 delegate 上限相同的 bulk channel。HostPeer 选择 lane,writer supervisor 负责把任一 writer 异常转成统一 peer failure。

源码位置:codex-rs/code-mode-host/src/lib.rs :: run_connection、negotiate

rust
let negotiated = negotiate(&mut reader, &mut writer, bulk_connections.as_ref()).await?;
let bulk_connection = match negotiated {
    NegotiatedConnection::Rejected => return Ok(()),
    NegotiatedConnection::Single => None,
    NegotiatedConnection::Dual(mut registration) => {
        match tokio::time::timeout(BULK_PAIRING_TIMEOUT, registration.receive()).await {
            Ok(Ok(connection)) => Some(connection),
            Ok(Err(_)) => anyhow::bail!("code-mode host bulk websocket pairing was abandoned"),
            Err(_) => anyhow::bail!("timed out pairing code-mode host bulk websocket"),
        }
    }
};

input loop 使用 biased select,让 disconnect、control session operation 和 shutdown 在 bulk callback 高负载时 仍能推进。dual 模式会校验每条消息的 transport lane;错误 lane 直接关闭连接。

3.2 HostState ​

连接建立后,host 创建 HostState。它拥有 runtime session map、session ID tombstone/seen set、operation request registry、request tasks、closing flag 和 HostPeer。request admission 受 host limits 控制;任务 panic 由 supervisor 转成 peer failure,而不是静默丢失 response。

源码位置:codex-rs/code-mode-host/src/lib.rs :: HostState

rust
struct HostState {
    sessions: Mutex<HashMap<SessionId, Arc<InProcessCodeModeSession>>>,
    limits: Arc<HostLimits>,
    seen_session_ids: Mutex<SeenSessionIds>,
    requests: Mutex<RequestRegistry>,
    request_tasks: TaskTracker,
    closing: AtomicBool,
    peer: Arc<HostPeer>,
}

4. HostPeer路由 ​

HostPeer 是 host 到 client 的 callback owner。它维护 pending delegate map、1024 permits、每个 (SessionId, CellId) 的 route、递增 delegate ID、disconnect token 和首个 failure。cell admission 前到达的 callback 可以进入 Pending queue;start_cell 建立 Active channel 后回放。重复 active route 或队列溢出会 disconnect,避免把 callback 交给错误 cell。

源码位置:codex-rs/code-mode-host/src/peer.rs :: HostPeer、CellRoute、HostPeer::start_cell

rust
pub(super) struct HostPeer {
    outgoing_tx: mpsc::Sender<EncodedFrame>,
    bulk_tx: Option<mpsc::Sender<EncodedFrame>>,
    pending: Mutex<HashMap<DelegateRequestId, PendingDelegate>>,
    delegate_permits: Arc<Semaphore>,
    cell_routes: StdMutex<HashMap<(SessionId, CellId), CellRoute>>,
    cell_routes_changed: Notify,
    next_request_id: AtomicI64,
    disconnected: CancellationToken,
    failure: StdMutex<Option<String>>,
}

enum CellRoute {
    Pending(VecDeque<CellMessage>),
    Active(mpsc::Sender<CellMessage>),
}

call 先登记 pending 和 permit,再路由到 cell,等待 dispatched acknowledgment 后才等待 response。取消只在 pending 仍存在时发送 CancelDelegateRequest;response、caller cancellation 和 connection disconnect 三条 路径都会移除 pending,permit 随 guard 释放。

5. Client Driver ​

5.1 Event循环 ​

V1 client ConnectionDriver 同时处理 connection cancellation、host event、execute 接管确认和 API command。 接管确认使用独立 channel,因为 Execute 收到 ExecutionStarted 后,caller 必须声明已接管 StartedCell; 如果 caller 在 admission 后取消,driver 会终止未被 claim 的 remote cell。

源码位置:codex-rs/code-mode/src/remote_session/connection/driver.rs :: ConnectionDriver::run

rust
loop {
    tokio::select! {
        biased;
        _ = self.cancellation.cancelled() => {
            self.fail("code-mode host connection closed".to_string());
            return;
        }
        event = self.event_rx.recv() => { /* host messages and failures */ }
        claim = self.execute_claim_rx.recv() => {
            let Some(request_id) = claim else {
                self.fail("code-mode execute claim stream closed".to_string());
                return;
            };
            self.requests.claim_execute(request_id);
        }
        command = self.command_rx.recv() => { /* session API command */ }
    }
}

5.2 Session registry ​

client registry 保存 SessionId → RemoteSession + delegate + cleanup + phase + wire/public cell map。每次操作 先 require_ready,完整比较 generation;shutdown 把 phase 改为 Closing,禁止新 operation。cell admission 把 wire ID 映射成带 generation 的 public ID,cell close 时只通知对应 delegate。

源码位置:codex-rs/code-mode/src/remote_session/connection/driver/session_registry.rs :: SessionRegistry

rust
pub(super) fn require_ready(&self, session: &RemoteSession) -> Result<(), String> {
    let record = self.records.get(&session.id)
        .ok_or_else(|| format!("unknown code-mode session {}", session.id))?;
    if record.remote != *session {
        return Err("stale code-mode session generation".to_string());
    }
    if record.phase != SessionPhase::Ready {
        return Err("code-mode session is shutting down".to_string());
    }
    Ok(())
}

5.3 Delegate runtime ​

client delegate runtime 最多保留 1024 active callbacks,并保留最近 4096 个 completed/cancelled ID 防重放。 callback task panic 被转换为 tool error;取消先 revoke completion path,再移除 active call,迟到 future 即使完成 也无法发送 response。cell close 会取消该 cell 全部 callbacks,然后调用 delegate.cell_closed。

源码位置:codex-rs/code-mode/src/remote_session/connection/driver/delegate_runtime.rs :: DelegateRuntime

6. 统一失败清理 ​

V1 driver 的 fail 只执行一次:标记 alive=false,保存首个原因,fail 所有 request,drain session registry, 让 delegate runtime fail 所有 session cleanup,最后取消 connection token。Drop for ConnectionDriver 也调用 fail,防止 task 意外退出留下 waiter。

源码位置:codex-rs/code-mode/src/remote_session/connection/driver.rs :: ConnectionDriver::fail、Drop

rust
fn fail(&mut self, reason: String) {
    if self.failed {
        return;
    }
    self.failed = true;
    self.alive.store(false, Ordering::Release);
    let reason = {
        let mut failure = self.failure.lock().unwrap();
        failure.get_or_insert(reason).clone()
    };
    self.requests.fail_all(&reason);
    let failed_sessions = self.sessions.drain();
    self.delegates.fail_all(failed_sessions);
    self.cancellation.cancel();
}

host side 输入结束后也先 peer.disconnect(),再以 timeout 执行 HostState::disconnect,等待/中止 session 与 request tasks,最后监督 writer。socket close 因此只是清理起点,不是资源已经释放的证明。

7. gRPC Lease ​

7.1 重连Binding ​

gRPC ReconnectableSession 保存可选 binding、单 permit opening coordinator、递增 generation、shutdown token 和共享 shutdown result。operation 先复用 live binding;若 lease 已 stopped,则等待旧 binding shutdown,再 创建新 generation 和 GenerationDelegate。shutdown 期间新 binding 不会 publish。

源码位置:codex-rs/code-mode/src/grpc_session/reconnect.rs :: ReconnectInner::get_or_open_binding

rust
if let Some(binding) = self.live_binding() {
    return Ok(binding);
}
let _opening_permit = tokio::select! {
    biased;
    _ = self.shutdown_requested.cancelled() => {
        return Err(SHUTDOWN_ERROR.to_string());
    }
    permit = self.opening_permit.acquire() => permit
        .map_err(|_| "gRPC code-mode session opening coordinator closed".to_string())?,
};
let generation = self.next_generation.fetch_add(1, Ordering::Relaxed);
let delegate = Arc::new(GenerationDelegate {
    delegate: Arc::clone(&self.delegate),
    generation,
});
let session = self.provider.open_binding(delegate, self.limits.clone()).await?;

public cell ID 在第二代以后包含 generation 前缀;wait/terminate 先验证并剥离当前 generation。旧 ID 无法 操作新 host 上碰巧相同的 wire cell ID。

7.2 Lease关闭 ​

GrpcCodeModeSession 的 lease/event/tool-subscription tasks 共享 stopped token 与 TaskTracker。任一 fatal stream 失败调用 close_state:取出 live cells、取消 stopped、关闭 tasks,并逐 cell 通知 delegate。 显式 shutdown 只在 session 仍 open 时调用 CloseSession,然后无论 RPC 成功与否都 close local state 并等待 stream tasks。

源码位置:codex-rs/code-mode/src/grpc_session/mod.rs :: SessionInner::close_state、drive_shutdown

rust
fn close_state(&self, failure: Option<String>) {
    let cells = self.state.lock().unwrap().close(failure);
    self.stopped.cancel();
    self.stream_tasks.close();
    for cell_id in cells {
        self.report_closed_cell(Some(cell_id));
    }
}

7.3 Wait retirement ​

gRPC wait 有独立 wait ID。caller drop 后 client 发 CancelWait;host 取消 observer,并等待 retired token 后才 acknowledge。client 收到 ack 才允许下一 wait,防止旧 observer 抢走新 outcome。这是 gRPC 相对 V1 driver 更显式的 cancellation barrier。

源码位置:

  • codex-rs/code-mode-host/src/grpc/waits.rs :: wait registry retirement
  • codex-rs/code-mode/src/grpc_session/operations.rs :: canceled wait cleanup

8. 源码验证 ​

codex-code-mode 的 70 项测试覆盖 V1 driver 和 gRPC client lifecycle:abandoned execute、wait retirement、 wrong response ID、connection failure、delegate panic/cancel/capacity、shared host reuse、gRPC lease state、 generation mapping、deadline 与 callback ownership。

源码位置:

  • codex-rs/code-mode/src/remote_session/connection/driver_tests.rs
  • codex-rs/code-mode/src/remote_session_tests.rs
  • codex-rs/code-mode/src/grpc_session/state_tests.rs
  • codex-rs/code-mode/src/grpc_session/generation_tests.rs
text
cd codex-rs
cargo test -p codex-code-mode --lib -- --test-threads=1

codex-code-mode-protocol 的 37 项测试补充握手、frame、wrong lane、消息形状、resource limits 和 identifier 边界。

text
cargo test -p codex-code-mode-protocol --lib -- --test-threads=1

host 内部的 stdio/WebSocket/gRPC integration tests 需要构建 V8 host binary。若 rusty_v8 对当前目标没有 可用 archive,这些动态测试不会进入断言;通过 client/codec 测试不能替代 host runtime 清理验证。

下一篇V8Runtime初始化将进入 host 内部,解释 V8 platform、isolate、module loader、globals 和 runtime thread 的初始化顺序。