服务端请求与客户端响应
本文承接OutgoingMessage写队列和协议映射与事件转换。本文只讨论 OutgoingMessageSender 的真实实现:App Server 如何把审批、动态工具或用户输入变成反向 request,如何保存等待者,如何处理客户端 response/error,以及 thread transition、resume 和 shutdown 怎样结束 pending request。
1. 两张状态表
源码位置:codex-rs/app-server/src/outgoing_message.rs :: OutgoingMessageSender
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,
}request_id_to_callback 面向“服务端发给客户端的反向 request”:key 是服务端生成的 RequestId,value 保存 oneshot sender、可选 ThreadId 和完整 ServerRequest。request_contexts 面向“客户端发给服务端的入站 request”:key 是 (ConnectionId, RequestId),用于最终 response/error 的 tracing parent 和 in-flight 计数。两者都叫 request,但生命周期和清理入口完全不同。
2. ID与注册
源码位置: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 callbacks = self.request_id_to_callback.lock().await;
callbacks.insert(
id,
PendingCallbackEntry {
callback: tx_approve,
thread_id,
request: request.clone(),
_diagnostics_guard: PENDING_SERVER_REQUESTS.track(),
},
);
}next_server_request_id 使用 AtomicI64::fetch_add 生成进程内递增 ID。注册发生在发送之前,保证客户端快速响应时 callback 已经存在。PendingCallbackEntry 还保存原始 request,是后续按 thread 重放和记录响应类型的依据。
3. 广播与定向发送
源码位置:codex-rs/app-server/src/outgoing_message.rs :: OutgoingEnvelope
pub(crate) enum OutgoingEnvelope {
ToConnection {
connection_id: ConnectionId,
message: OutgoingMessage,
write_complete_tx: Option<oneshot::Sender<()>>,
},
Broadcast { message: OutgoingMessage },
}send_request_to_connections 的 None 分支发送 Broadcast;Some(&[ConnectionId]) 则逐个发送 ToConnection。ThreadScopedOutgoingMessageSender 正是通过后者把请求限定到当前 thread 关联的连接集合。任一发送失败都会记录 warning 并从 callback map 移除 entry,避免 Core 永远等待一个已经无法送达的请求。
4. Response与Error
源码位置:codex-rs/app-server/src/outgoing_message.rs :: notify_client_response、notify_client_error
let entry = self.take_request_callback(&id).await;
match entry {
Some((id, entry)) => {
if let Ok(response) = entry.request.response_from_result(result.clone()) {
self.analytics_events_client
.track_server_response(now_unix_timestamp_ms(), response);
}
entry.callback.send(Ok(result)).ok();
}
None => warn!("could not find callback for {id:?}"),
}response 和 error 都先调用 take_request_callback,在 map 锁内原子移除 entry,再在锁外完成日志、analytics 和 oneshot send。这样重复 response、迟到 response 或 response/error 竞争时,只有第一个事件能拿到 entry;后续事件只能走 “could not find callback” 分支。notify_client_error 不记录错误消息和 data,因为其中可能包含凭据。
5. 入站响应路径
源码位置:codex-rs/app-server/src/outgoing_message.rs :: send_response_as_inner
let request_context = self.take_request_context(&request_id).await;
let outgoing_message = OutgoingMessage::Response(OutgoingResponse {
id: request_id.request_id,
result: Box::new(response),
});
self.send_outgoing_message_to_connection(
request_context,
request_id.connection_id,
outgoing_message,
"response",
).await;这里处理的是客户端先发来的 request。ClientResponsePayload 先被装进 OutgoingResponse,再按 ConnectionRequestId 找到并移除 RequestContext。因此 response 发出时,in-flight GaugeGuard 和 tracing context 一并结束;它不会触碰 request_id_to_callback,因为那张表属于反向 request。
6. Thread重放
源码位置:codex-rs/app-server/src/outgoing_message.rs :: pending_requests_for_thread、replay_requests_to_connection_for_thread
let mut requests = callbacks
.values()
.filter_map(|entry| {
(entry.thread_id == Some(thread_id)).then_some(entry.request.clone())
})
.collect::<Vec<_>>();
requests.sort_by(|left, right| left.id().cmp(right.id()));resume 时只读取指定 ThreadId 的 pending entries,按 request ID 排序后发送到新连接。重放使用原 request 和原 ID,不重新插入 callback,也不创建第二个 waiter;客户端回复后仍由原 callback 完成 Core 的等待。
7. 取消与关闭
源码位置:codex-rs/app-server/src/outgoing_message.rs :: cancel_requests_for_thread、cancel_all_requests
let entries = {
let mut callbacks = self.request_id_to_callback.lock().await;
let ids = callbacks.iter()
.filter_map(|(id, entry)|
(entry.thread_id == Some(thread_id)).then_some(id.clone()))
.collect::<Vec<_>>();
ids.into_iter()
.filter_map(|id| callbacks.remove(&id))
.collect::<Vec<_>>()
};
for entry in entries {
if let Some(error) = error.as_ref() {
entry.callback.send(Err(error.clone())).ok();
}
}thread 取消先在锁内完成“找出并移除”,再在锁外发送错误;因此 callback map 已经不再包含这些 request,迟到响应不会重新唤醒 waiter。cancel_all_requests 使用 HashMap::drain 做同样的全局清理。cancel_request 则只移除一个 ID,主要用于单请求撤销;是否向 waiter 发送错误取决于调用方传入的路径。
8. Thread切换语义
ThreadScopedOutgoingMessageSender::abort_pending_server_requests 会调用 cancel_requests_for_thread,并传入带 TURN_TRANSITION_PENDING_REQUEST_ERROR_REASON 的内部错误。这个错误不是客户端 response,而是服务端主动结束 pending waiter 的结果;它让 Core 能区分“用户拒绝”“客户端报错”和“Turn 状态已改变”。因此 Turn transition、连接断开和 runtime shutdown 虽然都表现为 waiter 结束,来源信息仍由各自的清理入口决定。
9. 测试边界
源码位置:
codex-rs/app-server/src/outgoing_message.rs::notify_client_error_forwards_error_to_waitercodex-rs/app-server/src/outgoing_message.rs::pending_requests_for_thread_returns_thread_requests_in_request_id_ordercodex-rs/app-server/src/outgoing_message.rs::cancel_requests_for_thread_cancels_all_thread_requestscodex-rs/app-server/src/outgoing_message.rs::send_response_clears_registered_request_contextcodex-rs/app-server/src/outgoing_message.rs::connection_closed_clears_registered_request_contexts
测试输入是受控的 server request payload、thread ID、连接 ID 和错误值;断言检查错误是否到达 oneshot waiter、thread pending request 是否按 ID 排序、取消后 map 是否为空,以及入站 response 或 connection close 是否清理 request context。它们验证 callback 与 context map 的生命周期,不证明真实传输层一定成功送达,也不证明客户端 UI 如何展示错误。
cd codex-rs
cargo test -p codex-app-server --lib outgoing_message -- --nocapture --test-threads=1补充 outgoing message 的生命周期图,分别展示请求注册、响应匹配和 thread 重放。
10. 排查顺序
server request 卡住时,先确认 request ID 是否已注册,再看发送目标是 broadcast 还是定向 connection;客户端回复无效时,检查 response ID 是否仍在 callback map。resume 后重复提示时,检查是否复用了原 request ID。入站 response 未返回时,检查 ConnectionRequestId 是否匹配并从 request_contexts 移除。Turn transition 或 shutdown 后仍等待,则分别追踪 cancel_requests_for_thread 与 cancel_all_requests 是否执行。
