请求串行化与顺序
本文承接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
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
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
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
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
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_fifocodex-rs/app-server/src/request_serialization.rs::different_keys_run_concurrentlycodex-rs/app-server/src/request_serialization.rs::closed_gate_request_is_skipped_and_following_requests_continuecodex-rs/app-server/src/request_serialization.rs::shutdown_of_live_gate_skips_already_queued_requestscodex-rs/app-server/src/request_serialization.rs::same_key_shared_reads_run_concurrentlycodex-rs/app-server/src/request_serialization.rs::later_shared_read_joins_running_shared_readcodex-rs/app-server/src/request_serialization.rs::later_shared_read_waits_behind_writer_queued_during_running_readcodex-rs/app-server/src/request_serialization.rs::startup_config_reads_run_concurrently_and_exclude_writescodex-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 语义。
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 自身的锁和事务,因为它们不属于这套资源队列的保证范围。
