断连清理与资源回收
本文承接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
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
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
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
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
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_closedcodex-rs/app-server/src/request_processors/process_exec_processor.rs :: connection_closedcodex-rs/app-server/src/connection_cleanup.rs :: ConnectionCleanupTaskscodex-rs/app-server/src/lib.rs :: graceful shutdown
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_shutdown9. 断连排查
遇到断连后资源仍占用,先查 connection map、RPC gate token 和 cleanup JoinSet,再查 process/fs/thread 子系统;遇到 resume 后实时事件缺失,确认这是旧 queue 丢失还是 rollout 可重建;遇到强制关闭,区分 abort cleanup 与业务 handler 未完成。
