Skip to content

断连清理与资源回收

追踪连接关闭后的 RPC 排空、outgoing、FS、process、thread listener 和后台任务回收。

基于rust-v0.150.0
CodexRustAppServerConnectionCleanup

断连清理与资源回收 ​

本文承接Connection RPC Gate和OutgoingMessage写队列,面向已经理解 connection state、gate、writer queue 和 pending request 的读者。本文回答连接关闭后资源按什么顺序回收、哪些任务等待或强制终止,以及 graceful shutdown 与 forced shutdown 的差异;不展开单个 transport 的 close frame。

1. 清理主线 ​

断连不是一个单独函数,而是一条跨层路径:主 loop 移除 connection、关闭 gate、通知 outbound router,然后把 processor cleanup 放入 JoinSet。这样可以先停止新请求,再异步排空旧任务。

2. 主Loop处理 ​

源码位置:codex-rs/app-server/src/lib.rs :: TransportEvent::ConnectionClosed

rust
let Some(connection_state) = connections.remove(&connection_id) else {
    continue;
};
let stdio_closed = connection_state.origin == ConnectionOrigin::Stdio;
connection_state.session.rpc_gate.close().await;
let outbound_closed = outbound_control_tx
    .send(OutboundControlEvent::Closed { connection_id })
    .await
    .is_ok();
let processor = Arc::clone(&processor);
connection_cleanup_tasks.spawn(async move {
    processor.connection_closed(connection_id, &connection_state.session).await;
});
if !outbound_closed {
    break "outbound_router_closed";
}
if single_client_mode && stdio_closed {
    break "stdio_connection_closed";
}

先从 map 移除可防止新 incoming message 找到连接;gate close 阻止迟到 handler;cleanup task 稍后处理具体资源。outbound_closed=false 会终止主 loop,因为写路由本身已不可用。

3. RPC排空 ​

源码位置:codex-rs/app-server/src/message_processor.rs :: connection_closed

rust
if timeout(
    CONNECTION_RPC_DRAIN_TIMEOUT,
    session_state.rpc_gate.shutdown(),
).await.is_err() {
    tracing::warn!(
        ?connection_id,
        timeout_seconds = CONNECTION_RPC_DRAIN_TIMEOUT.as_secs(),
        "timed out waiting for connection RPCs to drain"
    );
}
self.outgoing.connection_closed(connection_id).await;

清理先关闭 gate 和 MCP event streams,再等待已经取得 gate token 的 handler,最长 30 秒;超时只记录 warning 并继续清理。它不会自动取消 handler future,也不保证其业务副作用回滚。

4. 资源清理 ​

MessageProcessor::connection_closed 随后调用 outgoing、FS watch、command exec、process exec 和 thread processor 的 connection cleanup。每个子系统按 connection ID 过滤自己的资源,避免误删其他连接。

源码位置:codex-rs/app-server/src/message_processor.rs :: connection_closed

rust
self.outgoing.connection_closed(connection_id).await;
self.fs_processor.connection_closed(connection_id).await;
self.command_exec_processor
    .connection_closed(connection_id)
    .await;
self.process_exec_processor
    .connection_closed(connection_id)
    .await;
self.thread_processor.connection_closed(connection_id).await;

processor 只把 connection id 传给各自 owner;outgoing 清理 request contexts,FS 清理 watches,command/process 清理会话,Thread processor 清理 listener。这里没有跨子系统共享的一把“大锁”,所以每个子系统都必须保证只删除属于该 connection 的对象。

5. Process回收 ​

源码位置:codex-rs/app-server/src/request_processors/process_exec_processor.rs :: ProcessExecManager::connection_closed

rust
let controls = {
    let mut sessions = self.sessions.lock().await;
    let handles = sessions.keys()
        .filter(|handle| handle.connection_id == connection_id)
        .cloned().collect::<Vec<_>>();
    handles.into_iter().filter_map(|handle| sessions.remove(&handle)).collect::<Vec<_>>()
};
for control in controls {
    let _ = control.control_tx.send(ProcessControlRequest {
        control: ProcessControl::Kill,
        response_tx: None,
    }).await;
}

先从 session map 移除,再发送 Kill,避免断连期间新的 write/resize 找回旧 process。kill 没有 response receiver,因为连接已经无法接收控制结果。

6. Cleanup任务 ​

源码位置:codex-rs/app-server/src/connection_cleanup.rs :: ConnectionCleanupTasks

rust
pub(crate) async fn drain(&mut self) {
    while let Some(result) = self.tasks.join_next().await {
        log_cleanup_result(result);
    }
}

pub(crate) async fn reap_next(&mut self) {
    if self.tasks.is_empty() {
        pending::<()>().await;
    }
    if let Some(result) = self.tasks.join_next().await {
        log_cleanup_result(result);
    }
}

pub(crate) fn abort(&mut self) {
    self.tasks.abort_all();
}

正常 shutdown 使用 drain 等待连接 cleanup;forced shutdown 直接 abort JoinSet。cleanup task 的 panic/错误只记录 warning,不会重新建立连接。

7. 失败与恢复 ​

cleanup timeout、writer closed、process kill 失败和 thread listener 清理失败是不同故障。断连后客户端可重新连接并 resume thread,但无法恢复已经丢失的实时 notification;应从 rollout/history projection 重建,而不是依赖旧 connection queue。

8. 源码验证 ​

connection cleanup 测试验证 close 后 gate drain、outgoing request context 删除、FS watch/process session 清理和 thread listener detach;process tests 验证断连会向每个 connection-owned process 发送 Kill;graceful shutdown tests 验证非强制路径 drain、强制路径 abort。它们证明回收顺序,不证明 OS 一定立即终止所有子进程。

源码位置:

  • codex-rs/app-server/src/message_processor.rs :: connection_closed
  • codex-rs/app-server/src/request_processors/process_exec_processor.rs :: connection_closed
  • codex-rs/app-server/src/connection_cleanup.rs :: ConnectionCleanupTasks
  • codex-rs/app-server/src/lib.rs :: graceful shutdown
text
cd codex-rs
cargo test -p codex-app-server connection_closed
cargo test -p codex-app-server process_exec
cargo test -p codex-app-server connection_cleanup
cargo test -p codex-app-server graceful_shutdown

9. 断连排查 ​

遇到断连后资源仍占用,先查 connection map、RPC gate token 和 cleanup JoinSet,再查 process/fs/thread 子系统;遇到 resume 后实时事件缺失,确认这是旧 queue 丢失还是 rollout 可重建;遇到强制关闭,区分 abort cleanup 与业务 handler 未完成。