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
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
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
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
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
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
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_messagecodex-rs/app-server/src/lib.rs :: TransportEvent::IncomingMessagecodex-rs/app-server/src/message_processor.rs :: run_request_with_context
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 stdio9. 读循环排查
遇到 malformed 日志,先检查 transport payload;遇到 overload,区分 request 是否收到错误 response、notification 是否被丢弃;遇到 response 无法完成 callback,检查 connection ID、request ID 和 outbound context;遇到断连后仍有 Core 活动,回到 thread owner 而不是继续查 socket reader。
