Skip to content

MessageProcessor读循环

追踪 TransportEvent、JSON-RPC 消息分派、请求上下文和 malformed/overload 失败路径。

基于rust-v0.150.0
CodexRustAppServerMessageProcessorTransport

MessageProcessor读循环 ​

本文承接初始化握手与ClientInfo和Connection核心状态,面向已经理解 connection state、request ID 和 outgoing queue 的读者。本文回答 transport 如何把字节变成 JSONRPCMessage、runtime loop 如何按消息种类分派,以及 malformed、未知连接和队列过载如何结束;不展开具体业务 processor。

1. 读循环结构 ​

读循环分两段:transport 负责 framing/JSON decode 并发送 TransportEvent;App Server 主 loop 负责连接 map、processor 调用和 cleanup。MessageProcessor 不是 socket reader,它消费已经解码的消息。

2. Transport事件 ​

源码位置:codex-rs/app-server-transport/src/transport/mod.rs :: TransportEvent

rust
pub enum TransportEvent {
    ConnectionOpened {
        connection_id: ConnectionId,
        origin: ConnectionOrigin,
        writer: mpsc::Sender<QueuedOutgoingMessage>,
        disconnect_sender: Option<CancellationToken>,
    },
    ConnectionClosed { connection_id: ConnectionId },
    IncomingMessage {
        connection_id: ConnectionId,
        message: JSONRPCMessage,
    },
}

connection open 携带 writer 和 disconnect capability;incoming message 已经绑定 connection ID;close 只携带 ID,主 loop 再根据 map 找到 session 和 outbound state。

3. JSON解码 ​

源码位置:codex-rs/app-server-transport/src/transport/mod.rs :: forward_incoming_message、enqueue_incoming_message

rust
match serde_json::from_str::<JSONRPCMessage>(payload) {
    Ok(message) => {
        enqueue_incoming_message(transport_event_tx, writer, connection_id, message).await
    }
    Err(err) => {
        error!("Failed to deserialize JSONRPCMessage: {err}");
        true
    }
}

JSON 解码失败停留在 transport 层,不会进入 MessageProcessor。返回值由具体 transport 用来决定是否继续读取;错误日志不等于向客户端发送了业务 response。

源码位置:codex-rs/app-server-transport/src/transport/mod.rs :: enqueue_incoming_message

rust
let event = TransportEvent::IncomingMessage {
    connection_id,
    message,
};
match transport_event_tx.try_send(event) {
    Ok(()) => true,
    Err(mpsc::error::TrySendError::Closed(_)) => false,
    Err(mpsc::error::TrySendError::Full(TransportEvent::IncomingMessage {
        connection_id,
        message: JSONRPCMessage::Request(request),
    })) => {
        let overload_error = OutgoingMessage::Error(OutgoingError {
            id: request.id,
            error: JSONRPCErrorError {
                code: OVERLOADED_ERROR_CODE,
                message: "Server overloaded; retry later.".to_string(),
                data: None,
            },
        });
        writer.try_send(QueuedOutgoingMessage::new(overload_error)).is_ok()
    }
    Err(mpsc::error::TrySendError::Full(event)) => {
        transport_event_tx.send(event).await.is_ok()
    }
}

入站队列满时,带 request id 的请求尝试直接向该连接写回 overload error;notification 没有 response id,则按 transport 队列策略等待或丢弃。这个分支发生在业务 processor 之前。

4. 主循环分派 ​

源码位置:codex-rs/app-server/src/lib.rs :: TransportEvent::IncomingMessage

rust
match message {
    JSONRPCMessage::Request(request) => {
        let Some(connection_state) = connections.get_mut(&connection_id) else {
            warn!("dropping request from unknown connection: {connection_id:?}");
            continue;
        };
        processor.process_request(
            connection_id,
            request,
            &transport,
            Arc::clone(&connection_state.session),
        ).await;
    }
    JSONRPCMessage::Response(response) => processor.process_response(response).await,
    JSONRPCMessage::Notification(notification) => processor.process_notification(notification).await,
    JSONRPCMessage::Error(err) => processor.process_error(err).await,
}

Request 需要找到 connection session;response/error 主要用于完成 server request callback;JSON-RPC notification 当前由 process_notification 记录,typed in-process notification 则走独立的 process_client_notification。未知连接的消息被丢弃,不会创建临时 session。

5. Request上下文 ​

源码位置:codex-rs/app-server/src/message_processor.rs :: process_request、run_request_with_context

rust
let request_id = ConnectionRequestId {
    connection_id,
    request_id: request.id.clone(),
};
let request_span =
    crate::app_server_tracing::request_span(&request, transport, connection_id, &session);
let request_trace = request.trace.as_ref().map(|trace| W3cTraceContext {
    traceparent: trace.traceparent.clone(),
    tracestate: trace.tracestate.clone(),
});
let request_context = RequestContext::new(request_id.clone(), request_span, request_trace);
Self::run_request_with_context(
    Arc::clone(&self.outgoing),
    request_context.clone(),
    async {
        let codex_request = deserialize_client_request(request);
        let result = match codex_request {
            Ok(codex_request) => {
                self.handle_client_request(
                    request_id.clone(),
                    codex_request,
                    Arc::clone(&session),
                    None,
                    request_context.clone(),
                )
                .await
            }
            Err(error) => Err(error),
        };
        if let Err(error) = result {
            self.outgoing.send_error(request_id.clone(), error).await;
        }
    },
).await;

async fn run_request_with_context<F>(
    outgoing: Arc<OutgoingMessageSender>,
    request_context: RequestContext,
    request_fut: F,
) where F: Future<Output = ()> {
    outgoing.register_request_context(request_context.clone()).await;
    request_fut.instrument(request_context.span()).await;
}

每个 request 先注册 context,再做 typed decode 和 handler 分派。deserialize_client_request 还会先拒绝已删除的旧字段;context 同时承载 connection/request ID、trace span 和后续 response/error 路由。

6. 过载保护 ​

源码位置:codex-rs/app-server-transport/src/transport/mod.rs :: enqueue_incoming_message

rust
Err(mpsc::error::TrySendError::Full(TransportEvent::IncomingMessage {
    connection_id,
    message: JSONRPCMessage::Request(request),
})) => {
    let overload_error = OutgoingMessage::Error(OutgoingError {
        id: request.id,
        error: JSONRPCErrorError {
            code: OVERLOADED_ERROR_CODE,
            message: "Server overloaded; retry later.".to_string(),
            data: None,
        },
    });
    match writer.try_send(QueuedOutgoingMessage::new(overload_error)) {
        Ok(()) => true,
        Err(mpsc::error::TrySendError::Closed(_)) => false,
        Err(mpsc::error::TrySendError::Full(_)) => {
            warn!(
                "dropping overload response for connection {:?}: outbound queue is full",
                connection_id
            );
            true
        }
    }
}

只有 request 可以得到 overload error,因为它有 ID;notification 没有 response 目标。过载错误由 transport writer 直接发送,不会进入业务 processor。

7. 关闭与清理 ​

收到 ConnectionClosed 后主 loop 移除 connection state、关闭 RPC gate、通知 outbound router,并把 processor.connection_closed 放入 cleanup task。读循环停止接收该 connection 的后续消息,但 Core thread 是否继续由其 owner 决定。

8. 源码验证 ​

transport 测试验证 malformed JSON、incoming queue 满时 request overload response、notification 丢弃和 channel close;message processor 测试验证 request context、response/error callback 和 unknown connection;stdio/UDS/WebSocket 测试验证 framing 后产生同样的 TransportEvent。它们证明读循环边界,不证明业务 handler 全部成功。

源码位置:

  • codex-rs/app-server-transport/src/transport/mod.rs :: enqueue_incoming_message
  • codex-rs/app-server/src/lib.rs :: TransportEvent::IncomingMessage
  • codex-rs/app-server/src/message_processor.rs :: run_request_with_context
text
cd codex-rs
cargo test -p codex-app-server-transport incoming_message
cargo test -p codex-app-server message_processor
cargo test -p codex-app-server-transport stdio

9. 读循环排查 ​

遇到 malformed 日志,先检查 transport payload;遇到 overload,区分 request 是否收到错误 response、notification 是否被丢弃;遇到 response 无法完成 callback,检查 connection ID、request ID 和 outbound context;遇到断连后仍有 Core 活动,回到 thread owner 而不是继续查 socket reader。