Skip to content

Connection RPC Gate

解释连接关闭时 RPC gate 如何拦截排队请求、等待在途 handler,并与资源串行队列协作。

基于rust-v0.150.0
CodexRustAppServerRPCConcurrency

Connection RPC Gate ​

本文承接Connection核心状态和MessageProcessor读循环,面向已经理解 connection close、async task 和 request serialization 的读者。本文聚焦一个并发问题:断连发生时,哪些 RPC 允许继续、哪些排队请求必须丢弃,以及 gate 与资源级 FIFO 队列如何组合;不重复介绍整个 connection state。

1. 两层并发控制 ​

RequestSerializationQueues 管同一资源的请求顺序,ConnectionRpcGate 管某个连接是否仍接受 handler。请求可以已经进入资源队列,却还没有取得 gate token;断连后这类请求不会执行。

2. Gate结构 ​

源码位置:codex-rs/app-server/src/connection_rpc_gate.rs :: ConnectionRpcGate

rust
pub(crate) struct ConnectionRpcGate {
    accepting: Mutex<bool>,
    tasks: TaskTracker,
}

pub(crate) async fn run<F>(&self, future: F)
where
    F: Future<Output = ()>,
{
    let token = {
        let accepting = self.accepting.lock().await;
        if !*accepting {
            return;
        }
        self.tasks.token()
    };
    future.await;
    drop(token);
}

accepting 与 token 创建在同一临界区,避免 close 已发生但新 handler 仍登记为在途。future 在 gate 关闭时甚至不会被 poll。

3. Close与Shutdown ​

源码位置:codex-rs/app-server/src/connection_rpc_gate.rs :: close、shutdown

rust
pub(crate) async fn close(&self) {
    let mut accepting = self.accepting.lock().await;
    *accepting = false;
    self.tasks.close();
}

pub(crate) async fn shutdown(&self) {
    self.close().await;
    self.tasks.wait().await;
}

close 立即返回,不等待已开始 handler;shutdown 才等待全部 token 释放。两阶段设计让主 loop 能迅速停止接收,同时将耗时排空放入 cleanup。

4. Processor接入 ​

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

rust
let rpc_gate = Arc::clone(&session.rpc_gate);
let request = QueuedInitializedRequest::new(
    rpc_gate,
    async move {
        let result = processor
            .handle_initialized_client_request(
                connection_request_id,
                codex_request,
                request_context,
                session,
                event_stream_ready,
            )
            .await;
        if let Err(error) = result {
            processor.outgoing.send_error(error_request_id, error).await;
        }
    },
);

所有已初始化 RPC 都被包装为 QueuedInitializedRequest。Initialize 自身在这条路径之前处理,因此 gate 保护的是初始化后的业务 handler。

5. 资源队列 ​

源码位置:codex-rs/app-server/src/request_serialization.rs :: RequestSerializationQueueKey、QueuedInitializedRequest、RequestSerializationQueues

rust
if let Some(scope) = serialization_scope {
    let (key, access) = RequestSerializationQueueKey::from_scope(connection_id, scope);
    self.request_serialization_queues.enqueue(key, access, request).await;
} else {
    tokio::spawn(async move { request.run().await; });
}

源码位置:codex-rs/app-server/src/request_serialization.rs :: QueuedInitializedRequest::run

rust
pub(crate) async fn run(self) {
    let Self { gate, future } = self;
    match gate {
        Some(gate) => gate.run(future).await,
        None => future.await,
    }
}

有 scope 的请求先进入 global/thread/process/fs-watch 等资源队列;无 scope 的请求直接 spawn,但两者最终都调用 request.run() 并经过 gate。

6. 排队竞态 ​

源码位置:codex-rs/app-server/src/request_serialization.rs :: RequestSerializationQueues::drain

rust
match queue.requests.pop_front() {
    Some(request) => {
        let access = request.access;
        let mut requests = vec![request];
        if access == RequestSerializationAccess::SharedRead {
            while queue.requests.front().is_some_and(|request| {
                request.access == RequestSerializationAccess::SharedRead
            }) {
                let Some(request) = queue.requests.pop_front() else {
                    break;
                };
                requests.push(request);
            }
        }
        (requests, Arc::clone(&queue.changed))
    }
    None => {
        queues.remove(&key);
        return;
    }
}

同一 key 的 exclusive 请求逐个执行;连续 shared-read 请求被成批取出,并通过 FuturesUnordered 并发运行。队列空时删除 key,因此不会为已经没有请求的资源保留后台 drain task。

同 key 的 exclusive 请求 FIFO 执行;连续 shared-read 可以并发。若第一个请求运行期间连接关闭,后续排队请求出队时发现 gate closed,future 不被 poll。资源队列仍会继续 drain 并最终删除 key。

7. 断连排空 ​

源码位置: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!("timed out waiting for connection RPCs to drain");
}
self.outgoing.connection_closed(connection_id).await;
self.fs_processor.connection_closed(connection_id).await;
self.process_exec_processor.connection_closed(connection_id).await;
self.thread_processor.connection_closed(connection_id).await;

先等待在途 RPC,再清理 connection-scoped request contexts、FS/process/thread 资源。超时只记录 warning 并继续清理;它不会强制取消仍在执行的 handler future。

8. 失败边界 ​

gate 丢弃迟到 future 时不会自动发送 JSON-RPC error,因为连接已经关闭;在途 handler 完成后产生的 response 也可能没有可写连接。资源队列顺序正确不代表客户端一定观察到结果。

9. 源码验证 ​

gate 测试验证 close 后 future 不被 poll、close 不等待在途 handler、shutdown 等待 token、shutdown 期间 late run 被丢弃;serialization tests 验证同 key FIFO、shared-read 并发和不同 key 并行。它们证明调度与排空,不证明业务 handler 支持取消。

源码位置:

  • codex-rs/app-server/src/connection_rpc_gate.rs :: shutdown_drops_late_runs_while_waiting_for_inflight_work
  • codex-rs/app-server/src/request_serialization.rs :: same_key_requests_run_fifo
  • codex-rs/app-server/src/request_serialization.rs :: shared read tests
text
cd codex-rs
cargo test -p codex-app-server connection_rpc_gate
cargo test -p codex-app-server request_serialization

10. 竞态排查 ​

断连后 handler 仍执行时,确认它是否已取得 gate token;排队请求没有任何副作用时,检查是否在出队前 gate 已关闭;请求顺序异常时检查 serialization key/access,而不是 gate;shutdown 超时时定位仍持有 TaskTracker token 的 handler。