ExecServer消息模型
exec-server 的消息协议看起来像 JSON-RPC,但线上对象省略了 jsonrpc 字段。真正决定消息含义的是 method、id、result 和 error 的组合;真正决定一次调用是否结束的,则是连接两端各自维护的 pending 表。本文从当前源码追踪这两条主线:协议 crate 负责把 JSON 解码成四类 envelope,server crate 负责把 envelope 交给 typed handler,并用 RequestId 将乱序响应交回正确的 caller。
本文面向已经读过ExecServer架构和ExecServer连接与握手 的读者。前者解释连接、环境和 session 的边界,后者解释连接何时进入 initialized 状态;这里不重复 具体 transport 握手,也不逐一展开文件系统或进程字段,而是解释消息层如何承载这些方法。读完后, 你应能从一个 JSON 对象定位其 envelope、找到参数类型和路由函数,并判断超时、重复请求 ID、连接关闭 时 pending caller 会发生什么。
1. Wire信封
1.1 四种变体
JSONRPCMessage 是线上所有对象的入口。反序列化先把对象读成受限的 serde_json::Value,再依据字段 存在性选择变体:带 method 且带 id 是 request,带 method 但没有 id 是 notification,带 result 是成功 response,其余对象按 error 结构解码。这个顺序很重要:notification 没有响应 ID, 因此 handler 完成后不能再向对端发送 response。
源码位置:codex-rs/exec-server-protocol/src/rpc.rs :: JSONRPCMessage
pub enum JSONRPCMessage {
Request(JSONRPCRequest),
Notification(JSONRPCNotification),
Response(JSONRPCResponse),
Error(JSONRPCError),
}
if object.contains_key("method") {
if object.contains_key("id") {
JSONRPCRequest::deserialize(value).map(Self::Request)
} else {
JSONRPCNotification::deserialize(value).map(Self::Notification)
}
} else if object.contains_key("result") {
JSONRPCResponse::deserialize(value).map(Self::Response)
} else {
JSONRPCError::deserialize(value).map(Self::Error)
}线上因此使用如下四种形状。params 可省略;response 必须有 id 和 result;error 的 id 仍然 用于关联原请求。
源码位置:codex-rs/exec-server-protocol/src/rpc.rs :: JSONRPCRequest、JSONRPCNotification、JSONRPCResponse、JSONRPCError
pub struct JSONRPCRequest {
pub id: RequestId,
pub method: String,
pub params: Option<serde_json::Value>,
pub trace: Option<W3cTraceContext>,
}
pub struct JSONRPCNotification {
pub method: String,
pub params: Option<serde_json::Value>,
}
pub struct JSONRPCResponse {
pub id: RequestId,
pub result: Result,
}
pub struct JSONRPCError {
pub error: JSONRPCErrorError,
pub id: RequestId,
}1.2 省略版本字段
协议文件的模块注释明确指出,exec-server 使用 Codex 自己的 JSON-RPC 方言,wire 上不发送 "jsonrpc": "2.0"。JSONRPC_VERSION 只是内部常量,不能据此要求每个输入对象包含版本字段。 这也是为什么分类逻辑直接检查 method、id 和 result,而不是先验证通用 JSON-RPC 的版本成员。
源码位置:codex-rs/exec-server-protocol/src/rpc.rs :: JSONRPC_VERSION、RequestId
pub const JSONRPC_VERSION: &str = "2.0";
#[serde(untagged)]
pub enum RequestId {
String(String),
Integer(i64),
}
impl fmt::Display for RequestId {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
match self {
Self::String(value) => f.write_str(value),
Self::Integer(value) => write!(f, "{value}"),
}
}
}RequestId 同时允许字符串和整数,并实现 Hash、Eq、排序和显示。它不是操作系统 PID,也不是 session ID;它只在一条 RPC 连接的 pending map 中承担关联键。调用方若把整数 ID 和字符串形式的 相同数字混用,两个 key 仍然不同。
2. 参数类型
2.1 方法常量
方法名集中定义在 protocol crate,server 和 client 共享同一字符串常量,避免一端发送 process/read、另一端注册 processRead 这种静默不匹配。
源码位置:codex-rs/exec-server-protocol/src/protocol.rs :: 方法常量
pub const INITIALIZE_METHOD: &str = "initialize";
pub const EXEC_METHOD: &str = "process/start";
pub const EXEC_READ_METHOD: &str = "process/read";
pub const EXEC_WRITE_METHOD: &str = "process/write";
pub const EXEC_SIGNAL_METHOD: &str = "process/signal";
pub const EXEC_TERMINATE_METHOD: &str = "process/terminate";
pub const EXEC_OUTPUT_DELTA_METHOD: &str = "process/output";
pub const EXEC_EXITED_METHOD: &str = "process/exited";
pub const EXEC_CLOSED_METHOD: &str = "process/closed";
pub const FS_READ_FILE_METHOD: &str = "fs/readFile";
pub const HTTP_REQUEST_METHOD: &str = "http/request";这些常量描述的是协议路由名,不等于 Rust 函数名。比如 process/output 是通知名,消费者可能在 进程 watcher 中注册 notification route,而 process/read 是带 response 的 request。
2.2 typed payload
信封中的 params 初始是 Option<Value>,只有命中路由后才解码成具体类型。以读取进程输出为例, ReadParams 将进程句柄、序号、字节预算和等待窗口放在一个稳定的 camelCase 对象中;返回值再把 每个带序号的 chunk 与退出、关闭和失败状态一起交给 caller。
源码位置:codex-rs/exec-server-protocol/src/protocol.rs :: ReadParams、ReadResponse、ProcessOutputChunk
pub struct ReadParams {
pub process_id: ProcessId,
pub after_seq: Option<u64>,
pub max_bytes: Option<usize>,
pub wait_ms: Option<u64>,
}
pub struct ProcessOutputChunk {
pub seq: u64,
pub stream: ExecOutputStream,
pub chunk: ByteChunk,
}
pub struct ReadResponse {
pub chunks: Vec<ProcessOutputChunk>,
pub next_seq: u64,
pub exited: bool,
pub exit_code: Option<i32>,
pub closed: bool,
pub failure: Option<String>,
pub sandbox_denied: bool,
}ByteChunk 使用 base64 作为 JSON 的透明表示。它保证 JSON 层不需要猜测二进制编码,同时保留 Vec<u8> 的所有权转换。
源码位置:codex-rs/exec-server-protocol/src/protocol.rs :: ByteChunk
#[serde(transparent)]
pub struct ByteChunk(#[serde(with = "base64_bytes")] pub Vec<u8>);
impl ByteChunk {
pub fn into_inner(self) -> Vec<u8> {
self.0
}
}3. 路由解码
3.1 请求和通知
RpcRouter::request 注册的是带返回值的 handler。它先保存 request 的 id,再把 params 解码为 泛型 P;参数解码失败时直接构造 invalid-params error,不会调用业务 handler。handler 成功后,返回 值被序列化为 Response;业务错误则变成 Error。
源码位置:codex-rs/exec-server/src/rpc.rs :: RpcRouter::request
pub(crate) fn request<P, R, F, Fut>(&mut self, method: &'static str, handler: F)
where
P: DeserializeOwned + Send + 'static,
R: Serialize + Send + 'static,
F: Fn(Arc<S>, P) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<R, JSONRPCErrorError>> + Send + 'static,
{
self.request_routes.insert(method, Box::new(move |state, request| {
let request_id = request.id;
let params = request.params;
let response = decode_request_params::<P>(params)
.map(|params| handler(state, params));
Box::pin(async move {
let response = match response {
Ok(response) => response.await,
Err(error) => return Some(RpcServerOutboundMessage::Error { request_id, error }),
};
Some(match response {
Ok(result) => match serde_json::to_value(result) {
Ok(result) => RpcServerOutboundMessage::Response { request_id, result },
Err(err) => RpcServerOutboundMessage::Error {
request_id,
error: internal_error(err.to_string()),
},
},
Err(error) => RpcServerOutboundMessage::Error { request_id, error },
})
})
}));
}通知使用另一张 route map,返回 Result<(), String>,因此没有 response envelope。这个区分让 process/output 之类的异步事件不会占用 caller 的 pending 表。
源码位置:codex-rs/exec-server/src/rpc.rs :: RpcRouter::notification、decode_params
pub(crate) fn notification<P, F, Fut>(&mut self, method: &'static str, handler: F)
where
P: DeserializeOwned + Send + 'static,
F: Fn(Arc<S>, P) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<(), String>> + Send + 'static,
{
self.notification_routes.insert(method, Box::new(move |state, notification| {
let params = decode_notification_params::<P>(notification.params)
.map(|params| handler(state, params));
Box::pin(async move {
let handler = params?;
handler.await
})
}));
}
fn decode_params<P>(params: Option<Value>) -> Result<P, serde_json::Error>
where
P: DeserializeOwned,
{
let params = params.unwrap_or(Value::Null);
let retry_as_null = matches!(¶ms, Value::Object(map) if map.is_empty());
match serde_json::from_value(params) {
Ok(params) => Ok(params),
Err(err) if retry_as_null => serde_json::from_value(Value::Null).map_err(|_| err),
Err(err) => Err(err),
}
}3.2 入站请求门槛
双向连接意味着 client 也可能收到对端 request。RpcClient::admit_inbound_request 先限制 request ID 必须是非负整数,再用 inbound_request_ids 拒绝同一 ID 的并发重复请求,最后从 semaphore 取得执行 槽位。任一条件失败,都不会留下半注册的 ID。
源码位置:codex-rs/exec-server/src/rpc.rs :: RpcClient::admit_inbound_request、RpcInboundRequestGuard
pub(crate) fn admit_inbound_request(
&self,
request_id: &RequestId,
call_slots: &Arc<Semaphore>,
) -> Result<RpcInboundRequestGuard, RpcInboundRequestAdmissionError> {
let request_id = match request_id {
RequestId::Integer(request_id) if *request_id >= 0 => request_id,
RequestId::Integer(_) | RequestId::String(_) => {
return Err(RpcInboundRequestAdmissionError::InvalidRequestId);
}
};
let request_id = RequestId::Integer(*request_id);
if !self.inbound_request_ids.lock().unwrap().insert(request_id.clone()) {
return Err(RpcInboundRequestAdmissionError::DuplicateRequestId);
}
let call_slot = Arc::clone(call_slots).try_acquire_owned().map_err(|_| {
self.inbound_request_ids.lock().unwrap().remove(&request_id);
RpcInboundRequestAdmissionError::AtCapacity
})?;
Ok(RpcInboundRequestGuard {
request_id,
request_ids: Arc::clone(&self.inbound_request_ids),
_call_slot: call_slot,
})
}guard 的 Drop 会删除 ID,并释放 semaphore permit。也就是说,重复请求的保护周期覆盖整个 handler 生命周期,而不是只覆盖“收到 JSON”这一瞬间。
4. 有界解码
4.1 节点预算
协议层没有直接调用 serde_json::from_str::<Value>,而是用 BoundedValueSeed 和 BoundedValueVisitor 递归解码。每个 JSON 节点先扣减 MAX_JSONRPC_VALUE_NODES;当前上限是 256 * 1024。这个限制针对的是数组、对象和嵌套值的数量,而不是字符串长度,所以一个很长但单一 的 scalar 仍然可以通过,紧凑的大数组则会在分配大量 Value 之前失败。
源码位置:codex-rs/exec-server-protocol/src/rpc.rs :: MAX_JSONRPC_VALUE_NODES、BoundedValueSeed::deserialize
const MAX_JSONRPC_VALUE_NODES: usize = 256 * 1024;
let Some(remaining) = self.remaining.checked_sub(1) else {
return Err(de::Error::custom(format!(
"JSON-RPC message exceeds the limit of {MAX_JSONRPC_VALUE_NODES} JSON values"
)));
};
*self.remaining = remaining;
deserializer.deserialize_any(BoundedValueVisitor {
remaining: self.remaining,
})visitor 还处理 serde_json arbitrary-precision 模式产生的 number/raw-value 包装,并在 raw wrapper 内部继续使用同一预算;包装层不能绕过限制。
4.2 重复键
对象 visitor 在插入 map 前检查 key 是否已经存在。重复 key 会产生明确的反序列化错误,而不是接受 “最后一个值覆盖前一个值”的解析器差异。这一行为对请求 envelope 尤其重要:重复 method 或 id 不应让不同语言的实现得到不同路由结果。
源码位置:codex-rs/exec-server-protocol/src/rpc.rs :: BoundedValueVisitor::visit_map
let mut values = Map::new();
let first_value = object.next_value_seed(BoundedValueSeed {
remaining: &mut *self.remaining,
})?;
values.insert(first_key, first_value);
while let Some(key) = object.next_key::<String>()? {
if values.contains_key(&key) {
return Err(de::Error::custom(format!(
"duplicate JSON object key `{key}`"
)));
}
let value = object.next_value_seed(BoundedValueSeed {
remaining: &mut *self.remaining,
})?;
values.insert(key, value);
}5. Pending生命周期
5.1 client发起调用
RpcClient::call_inner 的顺序决定了关闭竞态是否安全:先生成整数 RequestId,在持有 async mutex 时同时检查 closed 状态并插入 oneshot sender,再把 request 放入写队列,最后等待 receiver。 如果序列化、发送或等待超时失败,都会删除对应 map 项。
源码位置:codex-rs/exec-server/src/rpc.rs :: RpcClient::call_inner
let request_id = RequestId::Integer(
self.next_request_id.fetch_add(1, Ordering::SeqCst),
);
let (response_tx, response_rx) = oneshot::channel();
{
let mut pending = self.pending.lock().await;
if self.closed.load(Ordering::Acquire) || *self.disconnected_rx.borrow() {
return Err(RpcCallError::Closed);
}
pending.retain(|_, sender| !sender.is_closed());
pending.insert(request_id.clone(), response_tx);
}
self.write_tx.send(JSONRPCMessage::Request(JSONRPCRequest {
id: request_id.clone(),
method: method.to_string(),
params: Some(params),
trace: codex_otel::current_span_w3c_trace_context(),
})).await.map_err(|_| RpcCallError::Closed)?;普通调用共享 MAX_IN_FLIGHT_REGULAR_CALLS = 1024 个槽位;清理调用在普通槽位耗尽时可以使用保留的 cleanup 槽位。这样连接故障时仍能尝试发送 terminate 等收尾请求,而普通流量不会无限积压。
源码位置:codex-rs/exec-server/src/rpc.rs :: acquire_regular_call_slot、call_for_cleanup
const MAX_IN_FLIGHT_REGULAR_CALLS: usize = 1024;
const RESERVED_CLEANUP_CALLS: usize = 1;
fn acquire_regular_call_slot(&self) -> Result<SemaphorePermit<'_>, RpcCallError> {
self.shared_call_slots.try_acquire().map_err(|_| {
RpcCallError::PendingRequestLimitExceeded {
limit: MAX_IN_FLIGHT_REGULAR_CALLS,
}
})
}5.2 响应和关闭
reader task 收到 Response 或 Error 时按 RequestId 从 pending map 移除 sender,再唤醒 caller。 连接事件在同一有序输入队列中先于 Disconnected 到达,因此已经读到的 response 会优先完成;reader 随后调用 drain_pending,将真正未完成的调用全部变成 RpcCallError::Closed。
源码位置:codex-rs/exec-server/src/rpc.rs :: handle_server_message、drain_pending
match message {
JSONRPCMessage::Response(JSONRPCResponse { id, result }) => {
if let Some(pending) = pending.lock().await.remove(&id) {
let _ = pending.send(Ok(result));
}
}
JSONRPCMessage::Error(JSONRPCError { id, error }) => {
if let Some(pending) = pending.lock().await.remove(&id) {
let _ = pending.send(Err(RpcCallError::Server(error)));
}
}
JSONRPCMessage::Request(request) => { /* 转给 RpcClientEvent */ }
JSONRPCMessage::Notification(notification) => { /* 转给 RpcClientEvent */ }
}
async fn drain_pending(pending: &Mutex<HashMap<RequestId, PendingRequest>>) {
let pending = {
let mut pending = pending.lock().await;
pending.drain().map(|(_, sender)| sender).collect::<Vec<_>>()
};
for sender in pending {
let _ = sender.send(Err(RpcCallError::Closed));
}
}5.3 server反向调用
server 处理 client request 时也可能需要反向调用 client,例如询问 executor 能力。这个方向由 RpcServerRequestSender 管理另一张 pending map,容量是 256;每次调用注册 RequestId::Integer、 发送 RpcServerOutboundMessage::Request,再由 complete 按 ID 完成。它与 RpcClient 的 map 相互 独立,因此两端可以同时等待对方的请求而不会覆盖 key。
源码位置:codex-rs/exec-server/src/rpc_server_requests.rs :: RpcServerRequestSender::call_with_timeout、complete、close
pub const MAX_IN_FLIGHT_SERVER_CALLS: usize = 256;
let request_id = RequestId::Integer(
self.inner.next_request_id.fetch_add(1, Ordering::SeqCst),
);
let (response_tx, response_rx) = oneshot::channel();
self.inner.pending.lock().unwrap().insert(request_id.clone(), response_tx);
self.inner.outgoing_tx.send(RpcServerOutboundMessage::Request(
JSONRPCRequest {
id: request_id,
method: method.to_string(),
params: Some(params),
trace: codex_otel::current_span_w3c_trace_context(),
},
)).await.map_err(|_| RpcCallError::Closed)?;6. 错误语义
6.1 稳定错误码
server helper 将协议层错误编码集中在 rpc.rs:invalid request 是 -32600,method not found 是 -32601,invalid params 是 -32602,internal error 是 -32603;找不到 process 等资源时使用 -32004,session 已被其他连接占用时使用专用的 -32010。错误消息可以提供上下文,但不应把 Rust debug 格式当作客户端可依赖的协议字段。
源码位置:codex-rs/exec-server/src/rpc.rs :: invalid_request、method_not_found、invalid_params、not_found、internal_error
pub(crate) fn invalid_params(message: String) -> JSONRPCErrorError {
JSONRPCErrorError { code: -32602, data: None, message }
}
pub(crate) fn not_found(message: String) -> JSONRPCErrorError {
JSONRPCErrorError { code: -32004, data: None, message }
}
pub(crate) fn internal_error(message: String) -> JSONRPCErrorError {
JSONRPCErrorError { code: -32603, data: None, message }
}6.2 失败边界
消息层能区分几类失败,但它们的消费者不同:
| 现象 | 产生位置 | caller 看到的结果 |
|---|---|---|
| JSON 结构或参数不匹配 | bounded decode / typed decode | Error(-32602) 或连接级 malformed event |
| 方法没有注册 | router 查找 | Error(-32601) |
| handler 找不到资源 | 业务 handler | Error(-32004) |
| 对端明确返回 error | response reader | RpcCallError::Server |
| 发送后连接关闭 | transport / reader | RpcCallError::Closed |
| 等待超过调用 deadline | call_inner / call_with_timeout | RpcCallError::TimedOut,并移除 pending |
| 并发槽位用尽 | semaphore | PendingRequestLimitExceeded |
这张表揭示一个边界:Error 是协议消息,Closed、TimedOut 和容量错误是本地调用状态。不能把 连接断开伪装成某个远端 JSON-RPC code,否则 caller 会误以为请求已经被 server 处理。
7. 源码验证
先验证 envelope 的四种变体,再验证有界解码的反向边界。协议测试不仅检查 round-trip,还检查任意 精度数字、raw wrapper 不能绕过节点预算、重复 key 被拒绝,以及紧凑数组会触发上限。
源码位置:codex-rs/exec-server-protocol/src/rpc_tests.rs :: round_trips_every_jsonrpc_message_variant、rejects_duplicate_object_keys、rejects_compact_array_heap_amplification
cd codex-rs
cargo test -p codex-exec-server-protocol --lib rpc::tests -- --test-threads=1这些测试证明的是 wire 解码边界,不证明具体业务方法的权限或 process 生命周期。
然后验证 server 反向调用。测试同时启动 slow 和 fast 两个请求,以不同顺序调用 complete,断言 每个 caller 收到自己的值;timeout 测试推进虚拟时间并断言 pending 已删除;close 测试填满 256 个 槽位后确认溢出调用被拒绝,并且关闭会唤醒全部等待者。
源码位置:codex-rs/exec-server/src/rpc_server_requests_tests.rs :: rpc_server_sender_matches_out_of_order_responses_by_request_id、rpc_server_sender_timeout_removes_pending_request、rpc_server_sender_bounds_and_drains_pending_requests_on_close
cargo test -p codex-exec-server --lib rpc::tests::rpc_client_matches_out_of_order_responses_by_request_id -- --test-threads=1
cargo test -p codex-exec-server --lib rpc::tests::rpc_client_timeout_removes_pending_request -- --test-threads=1
cargo test -p codex-exec-server --lib rpc_server_requests::tests -- --test-threads=1阅读源码时可以沿着这条路径复述一次调用:RpcClient::call 获取 semaphore → call_inner 登记 RequestId → 写队列输出 JSONRPCRequest → reader 的 handle_server_message 按 ID 移除 sender → caller 反序列化 typed response;若连接先结束,则 drain_pending 将剩余调用标记为 Closed。要定位 一个“请求一直不返回”的问题,应依次检查 request 是否进入写队列、对端是否回相同类型的 ID、pending 是否被 timeout 删除,以及 reader 是否已经观察到 Disconnected。
下一篇ExecServer进程RPC会在这套 envelope 和 pending 机制之上,继续 追踪 process/start、process/read、process/write 和输出通知的字段转换与生命周期。
