Skip to content

OutgoingMessage写队列

追踪 response、notification 和 server request 如何进入写队列、关联回调并在断连时清理。

基于rust-v0.150.0
CodexRustAppServerTransportQueue

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

rust
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

rust
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

rust
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

rust
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

rust
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_context
  • codex-rs/app-server/src/outgoing_message.rs :: pending_requests_for_thread_returns_thread_requests_in_request_id_order
  • codex-rs/app-server/src/transport.rs :: send_message_to_connection
text
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 transport

9. 写队列排查 ​

遇到 response 丢失,检查 request context 是否提前移除、目标 connection 是否仍在 map;遇到 server request 卡住,检查 callback map 与客户端 response ID;遇到慢连接断开,检查 writer queue 和 connection origin;遇到 shutdown 后仍有等待,检查 cancel_all_requests 是否覆盖对应 thread。