Skip to content

请求串行化与顺序

解释 V2 请求如何按资源 key 排队、共享读取、跨资源并行,以及队列如何与连接 gate 协作。

基于rust-v0.150.0
CodexRustAppServerConcurrencyOrdering

请求串行化与顺序 ​

本文承接V2请求分派总表和Connection RPC Gate。目标不是给出抽象的“请求会排队”结论,而是沿着 ClientRequestSerializationScope、RequestSerializationQueueKey、RequestSerializationQueues::enqueue 和 drain 的真实调用边界,解释一个请求怎样获得资源顺序,以及为什么只读请求可以并发。

1. 两层顺序 ​

请求顺序由资源队列和连接 gate 两层共同决定。资源队列回答“同一个资源什么时候轮到这个请求”;gate 回答“轮到后,这条连接是否仍允许执行 handler”。两者职责不同:不同 key 的队列可以并行,同 key 仍受 FIFO 或 shared-read 批次约束,而关闭的 gate 会在 future 真正执行前跳过它。

2. Scope到Key ​

源码位置:codex-rs/app-server/src/request_serialization.rs :: RequestSerializationQueueKey::from_scope

rust
match scope {
    ClientRequestSerializationScope::Global(name) =>
        (Self::Global(name), RequestSerializationAccess::Exclusive),
    ClientRequestSerializationScope::GlobalSharedRead(name) =>
        (Self::Global(name), RequestSerializationAccess::SharedRead),
    ClientRequestSerializationScope::Thread { thread_id } =>
        (Self::Thread { thread_id }, RequestSerializationAccess::Exclusive),
    ClientRequestSerializationScope::ThreadPath { path } =>
        (Self::ThreadPath { path }, RequestSerializationAccess::Exclusive),
    ClientRequestSerializationScope::CommandExecProcess { process_id } =>
        (Self::CommandExecProcess { connection_id, process_id },
         RequestSerializationAccess::Exclusive),
    ClientRequestSerializationScope::Process { process_handle } =>
        (Self::Process { connection_id, process_handle },
         RequestSerializationAccess::Exclusive),
    ClientRequestSerializationScope::FuzzyFileSearchSession { session_id } =>
        (Self::FuzzyFileSearchSession { session_id },
         RequestSerializationAccess::Exclusive),
    ClientRequestSerializationScope::FsWatch { watch_id } =>
        (Self::FsWatch { connection_id, watch_id },
         RequestSerializationAccess::Exclusive),
    ClientRequestSerializationScope::McpOauth { server_name } =>
        (Self::McpOauth { server_name }, RequestSerializationAccess::Exclusive),
}

GlobalSharedRead 与 Global 共享同一个 Global(name) key,区别只在访问模式;因此读和写会相遇在同一条队列上。Process、CommandExecProcess 和 FsWatch 把 ConnectionId 放进 key,防止两个连接恰好使用相同的句柄或 watch ID 时互相阻塞。ThreadPath、模糊搜索 session 和 MCP OAuth 则按各自的稳定身份串行化。

3. 请求对象 ​

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

rust
pub(crate) struct QueuedInitializedRequest {
    gate: Option<Arc<ConnectionRpcGate>>,
    future: BoxFutureUnit,
}

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

RPC 请求携带 Some(gate),所以资源顺序和连接生命周期都得到处理;enqueue_background 创建 gate: None 的请求,让 app-owned work 与 RPC 共用同一资源队列,但不被某个客户端连接关闭所跳过。这是把“资源互斥”与“客户端是否还活着”分开的关键设计。

4. 入队与唤醒 ​

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

rust
struct RequestSerializationQueue {
    requests: VecDeque<QueuedSerializedRequest>,
    changed: Arc<Notify>,
}

let request = QueuedSerializedRequest {
    access,
    request,
    _diagnostics_guard: QUEUED_REQUESTS.track(),
};
let should_spawn = {
    let mut queues = self.inner.lock().await;
    match queues.get_mut(&key) {
        Some(queue) => {
            queue.requests.push_back(request);
            queue.changed.notify_one();
            false
        }
        None => {
            let mut requests = VecDeque::new();
            requests.push_back(request);
            queues.insert(key.clone(), RequestSerializationQueue {
                requests,
                changed: Arc::new(Notify::new()),
            });
            true
        }
    }
};

每个 key 只有第一次入队会启动一个带 tracing span 的 drain task;后续请求只追加到同一 VecDeque 并通知 Notify。GaugeGuard 的生命周期覆盖排队阶段,因此诊断指标表示仍在队列中的请求,而不是已经被 drain 取走的请求。不同 key 在 map 中拥有独立队列,天然允许并行。

5. Exclusive执行 ​

源码位置: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
            }) {
                requests.push(queue.requests.pop_front().expect("front exists"));
            }
        }
        (requests, Arc::clone(&queue.changed))
    }
    None => {
        queues.remove(&key);
        return;
    }
}

当队首是 Exclusive 时,drain 逐个调用 request.request.run().await。因此同 key 的写请求严格保持入队顺序;写请求后面的读不能越过它。当前实现还会在队列空时删除 key,避免已完成资源永久占据 map。

6. SharedRead批次 ​

队首是 SharedRead 时,先取出当前连续的 shared-read,并放入 FuturesUnordered 并发执行。但批次不是静态快照:只要其中仍有读 future 运行,tokio::select! 会同时等待某个读完成或 changed.notified();后者触发时,drain 再次从队首吸收连续的 shared-read。

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

rust
let mut running_reads = requests
    .into_iter()
    .map(|request| request.request.run())
    .collect::<FuturesUnordered<_>>();

loop {
    tokio::select! {
        Some(()) = running_reads.next() => {
            if running_reads.is_empty() { break; }
        }
        () = changed.notified() => {
            let requests = /* 从同一 key 队首继续取连续 SharedRead */;
            for request in requests {
                running_reads.push(request.request.run());
            }
        }
    }
}

这解释了 later_shared_read_joins_running_shared_read 的行为:第一个读已经执行时,后来到达的读仍可加入当前批次。相反,如果队列中已经出现 Exclusive,吸收循环会停止;其后的读必须等待这个写完成。这样既保留写入顺序,也避免读请求因为批次建立时机而无谓延迟。

7. Gate边界 ​

ConnectionRpcGate::run 只在 future 被取出后检查连接状态。排队中的请求不会提前占用 gate token;请求轮到执行时,已关闭的 gate 会跳过其 future,之后同 key 的请求仍能继续。已经取得 token 的 handler 则由 gate 的 token 生命周期负责完成收尾,不能把“已开始执行”误解为“断连后立即取消”。

队列保证资源顺序,gate 保证连接级准入和 shutdown 等待;数据库事务、业务锁以及 handler 内部异步操作的顺序仍需在各自模块验证。

8. 测试边界 ​

源码位置:

  • codex-rs/app-server/src/request_serialization.rs :: same_key_requests_run_fifo
  • codex-rs/app-server/src/request_serialization.rs :: different_keys_run_concurrently
  • codex-rs/app-server/src/request_serialization.rs :: closed_gate_request_is_skipped_and_following_requests_continue
  • codex-rs/app-server/src/request_serialization.rs :: shutdown_of_live_gate_skips_already_queued_requests
  • codex-rs/app-server/src/request_serialization.rs :: same_key_shared_reads_run_concurrently
  • codex-rs/app-server/src/request_serialization.rs :: later_shared_read_joins_running_shared_read
  • codex-rs/app-server/src/request_serialization.rs :: later_shared_read_waits_behind_writer_queued_during_running_read
  • codex-rs/app-server/src/request_serialization.rs :: startup_config_reads_run_concurrently_and_exclude_writes
  • codex-rs/app-server/src/request_serialization.rs :: later_shared_reads_do_not_jump_ahead_of_queued_write

这些测试的输入是带有不同 access、key 和 gate 状态的受控 future;断言分别检查 FIFO 值序列、跨 key 是否能先完成、读 future 是否同时启动、写入是否阻挡后续读,以及关闭或 shutdown 后排队请求是否被跳过。测试结果证明的是队列调度合同,不证明真实文件系统、网络服务或业务 handler 的内部并发安全。

客观边界:队列只约束由同一个 RequestSerializationQueueKey 标识的资源;不同 key 之间没有全局顺序,handler 内部的锁、事务和外部服务并不自动继承这里的 FIFO 语义。

text
cd codex-rs
cargo test -p codex-app-server --lib request_serialization -- --nocapture --test-threads=1

补充队列键、批次和 gate 的关系图。

9. 排查顺序 ​

遇到顺序异常,先记录 serialization_scope 转换出的 key 和 access,再确认请求是否真的落在同一个队列。若同 key 乱序,检查是否把写误标成 SharedRead;若读越过写,检查写请求是否在队列中先出现。若断连后仍有副作用,区分请求是在 gate 关闭前已取得 token,还是仍处于排队状态。最后再检查 handler 自身的锁和事务,因为它们不属于这套资源队列的保证范围。