Skip to content

ExecServer消息模型

从 JSON-RPC 信封、有界 JSON 解码到双向 pending 表,追踪 exec-server 消息如何被分类、路由和完成。

基于rust-v0.150.0
CodexRustExecutionExecServerJSON-RPC

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

rust
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

rust
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

rust
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 :: 方法常量

rust
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

rust
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

rust
#[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

rust
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

rust
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!(&params, 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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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

rust
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 decodeError(-32602) 或连接级 malformed event
方法没有注册router 查找Error(-32601)
handler 找不到资源业务 handlerError(-32004)
对端明确返回 errorresponse readerRpcCallError::Server
发送后连接关闭transport / readerRpcCallError::Closed
等待超过调用 deadlinecall_inner / call_with_timeoutRpcCallError::TimedOut,并移除 pending
并发槽位用尽semaphorePendingRequestLimitExceeded

这张表揭示一个边界: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

text
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

text
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 和输出通知的字段转换与生命周期。