OutgoingMessage写队列
本文承接MessageProcessor读循环和Connection核心状态,面向已经理解 inbound request、connection ID 和 bounded writer queue 的读者。本文回答 OutgoingMessageSender 如何管理 request context、server request callback、定向发送和广播,以及队列满和连接关闭时发生什么;不展开具体 transport 的字节 framing。
1. 写路径
业务 processor 不直接写 socket,而是提交 OutgoingEnvelope。sender 负责 response/error 和 callback,transport router 再按 connection 状态过滤并写入每个连接的 writer。
2. Sender状态
源码位置:codex-rs/app-server/src/outgoing_message.rs :: OutgoingMessageSender、ThreadScopedOutgoingMessageSender
pub(crate) struct OutgoingMessageSender {
next_server_request_id: AtomicI64,
sender: mpsc::Sender<OutgoingEnvelope>,
request_id_to_callback: Mutex<HashMap<RequestId, PendingCallbackEntry>>,
request_contexts: Mutex<HashMap<ConnectionRequestId, RequestContext>>,
analytics_events_client: AnalyticsEventsClient,
}
pub(crate) struct ThreadScopedOutgoingMessageSender {
outgoing: Arc<OutgoingMessageSender>,
connection_ids: Arc<Vec<ConnectionId>>,
thread_id: ThreadId,
}OutgoingMessageSender 管两类关联:入站 request context 和服务端发出的 pending callback。ThreadScopedOutgoingMessageSender 再把一组 connection id 与 ThreadId 固定下来,保证 Thread notification/request 只发给订阅该 Thread 的连接,并能按 Thread 取消 pending server requests。
3. Server回调
源码位置:codex-rs/app-server/src/outgoing_message.rs :: send_request_to_connections
let id = self.next_request_id();
let request = request.request_with_id(id.clone());
let (tx_approve, rx_approve) = oneshot::channel();
{
let mut request_id_to_callback = self.request_id_to_callback.lock().await;
request_id_to_callback.insert(
id,
PendingCallbackEntry {
callback: tx_approve,
thread_id,
request: request.clone(),
_diagnostics_guard: PENDING_SERVER_REQUESTS.track(),
},
);
}
let outgoing_message = OutgoingMessage::Request(request.clone());
let send_result = match connection_ids {
None => {
self.sender
.send(OutgoingEnvelope::Broadcast {
message: outgoing_message,
})
.await
}
Some(connection_ids) => {
let mut send_error = None;
for connection_id in connection_ids {
if let Err(err) = self
.sender
.send(OutgoingEnvelope::ToConnection {
connection_id: *connection_id,
message: outgoing_message.clone(),
write_complete_tx: None,
})
.await
{
send_error = Some(err);
break;
}
}
match send_error {
Some(err) => Err(err),
None => Ok(()),
}
}
};
if send_result.is_err() {
let mut request_id_to_callback = self.request_id_to_callback.lock().await;
request_id_to_callback.remove(&id);
}callback 必须先写入 map,再发送 envelope;发送失败会回收 callback。定向发送逐连接提交同一 request 副本,广播则交给 router 再按 initialized/capability 过滤。
4. Response清理
源码位置:codex-rs/app-server/src/outgoing_message.rs :: send_response_as_inner、send_error
async fn send_response_as_inner(
&self,
request_id: ConnectionRequestId,
response: ClientResponsePayload,
thread_originator: Option<String>,
) {
let connection_id = request_id.connection_id;
let request_id_for_analytics = request_id.request_id.clone();
match thread_originator {
Some(thread_originator) => {
self.analytics_events_client
.track_response_with_thread_originator(
connection_id.0,
request_id_for_analytics,
&response,
thread_originator,
);
}
None => {
self.analytics_events_client.track_response(
connection_id.0,
request_id_for_analytics,
&response,
);
}
}
let response = Box::new(response);
let request_context = self.take_request_context(&request_id).await;
let outgoing_message = OutgoingMessage::Response(OutgoingResponse {
id: request_id.request_id,
result: response,
});
self.send_outgoing_message_to_connection(
request_context,
connection_id,
outgoing_message,
"response",
)
.await;
}最终 response/error 发送前会移除 request context,并把 context span 继续用于发送 tracing。response/error 总是定向到原始 connection,而 notification 才有广播路径。
5. 广播与过滤
Broadcast 只发送到 initialized 且未被 capability/opt-out 过滤的连接;定向 ToConnection 则先查 connection 是否存在。实验 notification 可能被跳过,server request 通常不会走 notification opt-out。
6. 背压处理
源码位置:codex-rs/app-server/src/transport.rs :: send_message_to_connection
if connection_state.can_disconnect() {
match writer.try_send(queued_message) {
Ok(()) => false,
Err(mpsc::error::TrySendError::Full(_)) => {
warn!("disconnecting slow connection after outbound queue filled");
disconnect_connection(connections, connection_id)
}
Err(mpsc::error::TrySendError::Closed(_)) => {
disconnect_connection(connections, connection_id)
}
}
} else if writer.send(queued_message).await.is_err() {
disconnect_connection(connections, connection_id)
}可主动断开的网络连接使用 try_send,队列满就断开慢连接;不可主动断开的嵌入式连接使用 await send。两种策略都会避免无限堆积,但阻塞/断连语义不同。
7. 取消与清理
源码位置:codex-rs/app-server/src/outgoing_message.rs :: cancel_all_requests、cancel_requests_for_thread
pub(crate) async fn cancel_all_requests(&self, error: Option<JSONRPCErrorError>) {
let entries = {
let mut request_id_to_callback = self.request_id_to_callback.lock().await;
request_id_to_callback
.drain()
.map(|(_, entry)| entry)
.collect::<Vec<_>>()
};
for entry in entries {
self.analytics_events_client
.track_server_request_aborted(now_unix_timestamp_ms(), entry.request.id().clone());
if let Some(error) = error.as_ref()
&& entry.callback.send(Err(error.clone())).is_err()
{
let request_id = entry.request.id();
warn!("could not notify callback for {request_id:?}: receiver dropped");
}
}
}Turn transition、connection close 或 runtime shutdown 会取消 pending server requests,向等待者发送 internal error,并清理入站 request contexts。取消只结束协议等待,不撤销已经发生的 Core/OS 副作用。
8. 源码验证
outgoing tests 验证 response/error 路由到目标 connection、发送后 request context 被移除、pending callback 按 ID 完成、thread-scoped cancellation 只影响对应请求;transport tests 验证 broadcast filter、slow connection disconnect 和 writer close。它们证明队列与关联边界,不证明底层 socket 已经把字节刷新到客户端。
源码位置:
codex-rs/app-server/src/outgoing_message.rs :: send_response_clears_registered_request_contextcodex-rs/app-server/src/outgoing_message.rs :: pending_requests_for_thread_returns_thread_requests_in_request_id_ordercodex-rs/app-server/src/transport.rs :: send_message_to_connection
cd codex-rs
cargo test -p codex-app-server send_response_clears_registered_request_context
cargo test -p codex-app-server pending_requests_for_thread
cargo test -p codex-app-server transport9. 写队列排查
遇到 response 丢失,检查 request context 是否提前移除、目标 connection 是否仍在 map;遇到 server request 卡住,检查 callback map 与客户端 response ID;遇到慢连接断开,检查 writer queue 和 connection origin;遇到 shutdown 后仍有等待,检查 cancel_all_requests 是否覆盖对应 thread。
