CodexThread背压
“Codex 会发事件”这句话隐藏了三套完全不同的传输语义:Core 用无界 event queue 保证 Turn 不被 UI 反压;App Server 用单 listener 消费 Core event,再按 Thread subscriber 做 fan-out;App Server client 面对有界消费队列时,才把通知分成 lossless 与 best-effort。
因此,背压不是一条从 UI 直接传回模型 task 的连续链。不同层选择了不同失败策略:Core 慢消费者会 积累内存;App Server outgoing queue 会 await;in-process client 的本地事件队列保持无界,使未读取的 lossless notification 不阻塞 request response;remote client 仍可能收到 Lagged marker。不能把某一层 的队列容量外推到整个链路。
本文面向已经了解 CodexThread公共API 中 next_event() 单消费者语义的 读者。本文只讨论Event从Core到App Server再到client facade的队列、持久化和背压,不展开每种 EventMsg 的业务字段,也不把Realtime/exec自己的广播队列混入同一条链。
读完后,应能判断事件卡在哪一层、哪类通知允许丢弃、为什么event receiver关闭不会取消模型Task,以及 为什么看到 Lagged { skipped } 后必须把它当传输健康信号而不是业务事件。
1. Event队列链
下面的类图只表达 owner 与队列类型,不展开事件变体。CodexThread 没有 fan-out 容器;它只持有 SessionIo 的单一 receiver。
agent_status 是旁路 watch,不经过 App Server event reducer;即使某个中间 event 没被客户端消费, status getter 仍可能已经看到更新后的最新状态。
2. Core生产者
Session 创建容量 512 的 submission queue,却为 event 使用 unbounded channel。输入生产者会在队列满时 等待,事件生产者则不会因为 UI 暂停消费而阻塞模型/工具循环。
源码位置:codex-rs/core/src/session/mod.rs :: SessionIo channels, SessionIo::next_event
pub(crate) const SUBMISSION_CHANNEL_CAPACITY: usize = 512;
let (tx_sub, rx_sub) = async_channel::bounded(SUBMISSION_CHANNEL_CAPACITY);
// event 无界是有意的隔离点:慢 UI 不直接阻塞 Turn,但积压会转化为内存压力。
let (tx_event, rx_event) = async_channel::unbounded();
pub(crate) struct SessionIo {
pub(crate) tx_sub: Sender<Submission>,
pub(crate) rx_event: Receiver<Event>,
pub(crate) agent_status: watch::Receiver<AgentStatus>,
pub(crate) session_loop_termination: SessionLoopTermination,
}
pub(crate) async fn next_event(&self) -> CodexResult<Event> {
let event = self
.rx_event
.recv()
.await
// sender 全部释放且队列排空后,消费者得到 InternalAgentDied,而不是空 Event。
.map_err(|_| CodexErr::InternalAgentDied)?;
Ok(event)
}unbounded 不等于“不会丢”。进程崩溃会丢内存队列;receiver 关闭后,新 event 会被 Session 记录为 debug drop。它只表示正常运行时没有容量上限和 producer backpressure。
并发调用 CodexThread::next_event() 也不会复制事件。多个 future 在同一 receiver 上竞争,每个 Event 只交给一个消费者。产品需要多订阅方时,必须先建立单 listener,再在更高层 fan-out。
3. 投递顺序
send_event() 先更新 terminal error/trace,再把 Turn ID 包装进 Event。真正的 raw pipeline 会先尝试 写 rollout,然后记录 protocol trace,最后更新 status watch 并发送 event。
源码位置:codex-rs/core/src/session/mod.rs :: send_event_raw_with_persistence, deliver_event_raw
async fn send_event_raw_with_persistence(&self, event: Event, persist: bool) {
if persist {
// 先把 EventMsg 交给 rollout policy;不是所有 event 都会成为 durable item。
let rollout_items = vec![RolloutItem::EventMsg(event.msg.clone())];
self.persist_rollout_items(&rollout_items).await;
}
self.services
.rollout_thread_trace
.record_protocol_event(&event.msg);
self.deliver_event_raw(event).await;
}
async fn deliver_event_raw(&self, event: Event) {
// status 在 event queue 发送前更新;event receiver 关闭不回滚最新状态。
if let Some(status) = agent_status_from_event(&event.msg) {
self.agent_status.send_replace(status);
}
if let Err(e) = self.tx_event.send(event).await {
// 没有消费者不会让 Session 失败,也不会把错误传播回 Turn task。
debug!("dropping event because channel is closed: {e}");
}
}
pub(crate) async fn persist_rollout_items(&self, items: &[RolloutItem]) {
if let Some(live_thread) = self.live_thread()
&& let Err(e) = live_thread.append_items(items).await
{
// persistence failure 被记录但不阻止后续实时事件投递。
error!("failed to record rollout items: {e:#}");
}
}这里保证的是调用顺序,不是“客户端收到事件时一定 durable”。persist_rollout_items() 吞掉 store error 并记日志,随后仍会投递;writer 的 flush/durability 又是更强的独立边界。实时事件和恢复历史应尽量 一致,但 I/O 故障时系统优先保留 live progress。
send_event_raw_without_materializing_rollout() 还会先检查本地 rollout 是否已经存在。对于尚未 materialize 的 Thread,它可以只投递 event,避免一条启动期 warning 意外创建持久化文件。
4. Rollout
持久化策略按 ThreadHistoryMode 和 event 语义筛选。下面的流程图展示主要分类,不替代源码中的完整 枚举匹配。
源码位置:codex-rs/rollout/src/policy.rs :: should_persist_event_msg(节选)
pub fn should_persist_event_msg(ev: &EventMsg, history_mode: ThreadHistoryMode) -> bool {
match ev {
EventMsg::ItemCompleted(event) => {
// Paginated 以 TurnItem 为主;Legacy 主要依赖 raw ResponseItem/旧事件表示。
matches!(history_mode, ThreadHistoryMode::Paginated)
|| matches!(
event.item,
TurnItem::Plan(_) | TurnItem::Extension(ExtensionItem::Sleep(_))
)
}
EventMsg::TokenCount(_)
| EventMsg::ThreadGoalUpdated(_)
| EventMsg::ThreadRolledBack(_)
| EventMsg::TurnAborted(_)
| EventMsg::TurnStarted(_)
| EventMsg::TurnComplete(_)
| EventMsg::ThreadSettingsApplied(_) => true,
EventMsg::UserMessage(_)
| EventMsg::AgentMessage(_)
| EventMsg::AgentReasoning(_)
| EventMsg::PatchApplyEnd(_)
| EventMsg::ContextCompacted(_)
| EventMsg::McpToolCallEnd(_)
| EventMsg::WebSearchEnd(_)
| EventMsg::ImageGenerationEnd(_) => matches!(history_mode, ThreadHistoryMode::Legacy),
EventMsg::SubAgentActivity(event) => {
// Paginated history stores completed sub-agent activity as TurnItem;
// non-completed activity remains a legacy-only compatibility record.
matches!(history_mode, ThreadHistoryMode::Legacy)
&& event.kind != SubAgentActivityKind::Completed
}
// Error、Warning、实时delta、审批请求和启动进度等属于 transient live stream。
EventMsg::Error(_)
| EventMsg::Warning(_)
| EventMsg::StreamError(_)
| EventMsg::SessionConfigured(_)
| EventMsg::ItemStarted(_)
| EventMsg::AgentMessageContentDelta(_)
| EventMsg::PlanDelta(_)
| EventMsg::ExecCommandOutputDelta(_)
| EventMsg::RequestUserInput(_)
| EventMsg::McpStartupUpdate(_)
| EventMsg::ShutdownComplete => false,
// ... 其余 transient 变体同样返回 false。
}
}源码中的完整 transient 分支比上面节选更多。关键结论是:SessionConfigured、delta、approval/request、 warning/error 不靠 EventMsg 直接恢复;durable history 由 lifecycle、ResponseItem、TurnItem 和其他 rollout record 共同组成。
5. App边界
每个 loaded Thread 的 App Server listener 独占 next_event()。它先更新 thread-local reducer,再根据 raw-event opt-in 过滤,随后查询当前订阅连接并创建 ThreadScopedOutgoingMessageSender。
源码位置:codex-rs/app-server/src/request_processors/thread_lifecycle.rs :: thread listener event branch
event = conversation.next_event() => {
let event = match event {
Ok(event) => event,
Err(err) => {
// Core event channel 终止后 listener 退出,不无限重试 closed receiver。
tracing::warn!("thread.next_event() failed with: {err}");
break;
}
};
let raw_events_enabled = {
let mut thread_state = thread_state.lock().await;
// 即使 raw event 不向客户端发送,也先更新 thread-local reducer 状态。
thread_state.track_current_turn_event(&event.id, &event.msg);
thread_state.experimental_raw_events
};
if matches!(
&event.msg,
EventMsg::RawResponseItem(_) | EventMsg::RawResponseCompleted(_)
) && !raw_events_enabled
{
continue;
}
let subscribed_connection_ids = thread_state_manager
.subscribed_connection_ids(conversation_id)
.await;
let thread_outgoing = ThreadScopedOutgoingMessageSender::new(
outgoing_for_task.clone(),
subscribed_connection_ids,
conversation_id,
);
apply_bespoke_event_handling(
event.clone(),
conversation_id,
conversation.clone(),
thread_manager.clone(),
thread_outgoing,
thread_state.clone(),
thread_watch_manager.clone(),
thread_list_state_permit.clone(),
fallback_model_provider.clone(),
)
.await;
}apply_bespoke_event_handling() 把 Core event 投影成一个或多个 v2 notification、server request 或状态 更新。listener 只有一个,所以 event ordering 在进入 reducer 前保持;fan-out 后不同 transport 的实际 写完成时间可以不同。
App Server 的 OutgoingMessageSender 内部使用容量 128 的 Tokio mpsc。对每个 connection 的 enqueue 调用会 send().await,所以从这里开始,满队列可以反压 Thread listener;WebSocket writer 另有更大的 32K outbound queue,但那是 transport 层缓冲,不改变 Core event queue 的无界性质。
6. Event投影
同一事实经过 Core、App Server 和客户端 facade 时使用不同 envelope。下面的 ER 图关注关联键和一对多 投影,不表示它们都被持久化。
Core Event.id 通常是 submission/Turn ID;ServerNotification payload 会再携带 thread/turn/item ID; transport envelope 增加 emitted timestamp。不要依赖 JSON-RPC request ID 去关联通知,它只属于有响应的 request/response 交换。
7. 丢弃策略位置
当前 in-process client 使用无界 event consumer channel;worker 从 runtime 事件流按顺序转发通知、server request 和 Lagged marker。command channel 与底层 runtime 仍受容量约束,但调用方暂时不读 event 不会 阻塞后续 request response。
源码位置:codex-rs/app-server-client/src/lib.rs :: InProcessAppServerClient::start
let (event_tx, event_rx) = mpsc::unbounded_channel::<InProcessServerEvent>();
// 本地 event queue 无界,避免未读取通知阻塞 request response。
if event_tx.send(event).is_err() {
event_stream_enabled = false;
}remote WebSocket 客户端则由 transport/consumer 的具体缓冲策略决定是否出现 Lagged;因此不能把旧版 “所有 client 都有界并按 notification 变体丢弃”的描述继续套用到当前 in-process 实现。
有 dropped event 后,facade 会在容量允许时先发送 Lagged { skipped }。如果下一条是 lossless event, lag marker 本身也使用 blocking send,保证客户端先知道中间发生过跳过,再收到不可丢终态。
CommandExecutionOutputDelta 等高频进度即使跳过,最终 ItemCompleted 仍能给出权威结果;assistant text delta 则不能丢,否则 Markdown transcript 会永久缺字。仓库用 capacity=1 的 remote/in-process 测试验证 lossless 通知和 lag marker 在强背压下仍可观察。
8. 关闭传播
| 关闭点 | 上游看到什么 | 下游动作 |
|---|---|---|
| CodexThread event receiver 被 drop | Session send 失败只记 debug | Core task 继续,不自动 shutdown |
| Session 所有 event sender drop | next_event() 返回 InternalAgentDied | App Server listener 记录 warning 并退出 |
| App Server subscriber 为空 | thread-scoped sender 直接返回 | reducer仍可更新,连接无通知 |
| Outgoing mpsc 关闭 | send 记录 warning | 当前 notification 丢失 |
| client event channel 关闭 | DisableStream | worker 停止轮询 server event,但仍处理 command 至 shutdown |
| transport 断开 | remote client 产生 Disconnected | 调用方决定重连或结束 |
status watch 与 event channel 的关闭也不对称。agent_status() 读取最后一个 watch value;订阅者等待 changed() 才会观察 sender 关闭。不能用一次 status getter 成功证明 event listener 仍存活。
App Server 的 pending server requests 还有独立 oneshot 生命周期。Turn 状态变化或 unload 时, abort_pending_server_requests() 会按 ThreadId 取消 waiter;若一个 ServerRequest 因 client queue 满被 丢弃,facade 必须立即回送 overload error,否则 server 会永久等待不存在的响应。
9. 三个故障演练
9.1 UI 停止读取
首先区分停在哪一层。若 CodexThread receiver 仍存在但无人 poll,无界 Core queue 持续增长;若 receiver 已 drop,Core 记录 event drop,history/status 仍可能继续。若只是 client facade queue 满,则看到 Lagged 且 best-effort 事件减少,lossless transcript 应保持完整。
9.2 TurnComplete恢复
检查 rollout append error 日志和 should_persist_event_msg()。Core 的顺序是先尝试持久化再投递,但 persistence error 不阻止 live event;因此实时成功不能单独证明 durable。还要确认 history mode、writer flush 与对应 TurnItem/ResponseItem 是否进入 rollout。
9.3 一个连接缺通知
先查该连接是否仍订阅 Thread,再查 outgoing/transport queue。Core listener 是共享的:如果 reducer 已 看到 event 且其他连接正常,问题位于 per-connection fan-out 之后;若所有连接都缺同一 event,才回到 raw-event opt-in、bespoke mapping 或 Core event channel。
10. 事件代码审阅
- 这条消息在哪一层产生:Core Event、ServerNotification、ServerRequest 还是 client marker?
- 它是否必须 durable;Legacy 与 Paginated 的策略是否相同?
- queue 是 bounded、unbounded、watch 还是 oneshot;满载与关闭分别怎样处理?
- 谁是唯一消费者,fan-out 在消费前还是消费后发生?
- 丢失这条消息会破坏 transcript/终态,还是只降低进度可见性?
只写“事件异步发送”无法回答任何一个问题。Codex 的可靠性来自每一层明确选择不同保证,而不是某个 全局 event bus 自动提供 exactly-once、durable 和 broadcast。
11. 背压边界测试
11.1 AppServer
Outgoing测试创建容量2的envelope channel,向两个ConnectionId发送一条ConfigWarning,然后接收两份 envelope。除连接ID不同外,两份通知必须共享同一个 emitted_at_ms,证明fan-out复用一次逻辑事件时间, 而不是为每个连接重建业务通知。
源码位置:codex-rs/app-server/src/outgoing_message.rs
// :: send_server_notification_to_connections_reuses_timestamp(断言节选)
outgoing
.send_server_notification_to_connections(
&[ConnectionId(1), ConnectionId(2)],
ServerNotification::ConfigWarning(ConfigWarningNotification {
summary: "test".to_string(),
details: None,
path: None,
range: None,
}),
)
.await;
let timestamps = [
rx.recv().await.expect("first connection should receive notification"),
rx.recv().await.expect("second connection should receive notification"),
]
.map(|envelope| match envelope {
OutgoingEnvelope::ToConnection {
message: OutgoingMessage::AppServerNotification(envelope),
..
} => envelope.emitted_at_ms,
_ => panic!("expected targeted server notification"),
});
// 两个连接收到独立envelope,但观察同一逻辑事件时间。
assert_eq!(timestamps[0], timestamps[1]);它证明目标连接fan-out,不证明某个慢连接不会拖累具体transport;等待写完成的策略由outbound router和 transport实现负责。
11.2 Best-effort丢弃
In-process client测试先把 stdout-1 填满容量1队列,再尝试发送 stdout-2。第二条属于可丢弃的 CommandExecutionOutputDelta,因此 skipped_events 变为1。消费者启动后,producer发送AgentMessage delta、ItemCompleted和TurnCompleted;这些lossless通知必须等待容量,并在它们之前插入Lagged。
源码位置:codex-rs/app-server-client/src/lib.rs
// :: forward_in_process_event_preserves_transcript_notifications_under_backpressure(节选)
let result = forward_in_process_event(
&event_tx,
&mut skipped_events,
InProcessServerEvent::ServerNotification(Box::new(
command_execution_output_delta_notification("stdout-2"),
)),
|_| {},
)
.await;
assert_eq!(result, ForwardEventResult::Continue);
// 满队列丢弃best-effort输出,并累计一个待报告skip。
assert_eq!(skipped_events, 1);
for notification in [
agent_message_delta_notification("hello"),
item_completed_notification("hello"),
turn_completed_notification(),
] {
let result = forward_in_process_event(
&event_tx,
&mut skipped_events,
InProcessServerEvent::ServerNotification(Box::new(notification)),
|_| {},
)
.await;
assert_eq!(result, ForwardEventResult::Continue);
}
assert_eq!(skipped_events, 0);
let events = receive_task.await.expect("receiver task should join successfully");
assert!(matches!(
&events[1],
// Lagged先于后续lossless transcript事件抵达。
InProcessServerEvent::Lagged { skipped: 1 }
));测试后续逐项断言AgentMessageDelta文本、ItemCompleted内容和TurnCompleted状态,证明三者没有因容量1丢失。 它不保证best-effort stdout-2 可恢复;Lagged只报告数量,不保存被丢payload。
11.3 Remote传输策略
remote_backpressure_preserves_transcript_notifications 让测试WebSocket连续发送两条stdout delta和三条 transcript/terminal通知,client channel容量仍为1。最终事件允许包含 Lagged { skipped: 1 },但收集到 的lossless名称必须严格为 agent_message_delta、item_completed、turn_completed。
这证明remote和in-process都委托同一个 server_notification_requires_delivery() 分类,不证明网络断开后 的重连重放;Remote client断线只产生Disconnected事件,历史恢复需要重新读取Thread/rollout。
11.4 Turn迁移
App Server可能正在等待客户端回答RequestUserInput。测试注册pending request后调用 abort_pending_server_requests();回调必须解析为错误,message说明Turn state改变,data中的reason固定为 turnTransition,避免旧Turn请求在新Turn中永久等待。
源码位置:codex-rs/app-server/src/request_processors/thread_processor_tests.rs
// :: aborting_pending_request_clears_pending_state(断言节选)
let (request_id, client_request_rx) = thread_outgoing
.send_request(ServerRequestPayload::ToolRequestUserInput(
ToolRequestUserInputParams {
thread_id: thread_id.to_string(),
turn_id: "turn-1".to_string(),
item_id: "call-1".to_string(),
questions: vec![],
is_blocking: true,
auto_resolution_ms: None,
},
))
.await;
thread_outgoing.abort_pending_server_requests().await;
let response = client_request_rx.await.expect("callback should be resolved");
let error = response.expect_err("request should be aborted during cleanup");
// waiter以结构化Turn迁移错误结束,而不是靠sender drop得到无上下文取消。
assert_eq!(
error.message,
"client request resolved because the turn state was changed"
);
assert_eq!(error.data, Some(json!({ "reason": "turnTransition" })));原测试还断言发出的request ID与目标connection一致。本文省略该部分,因为这里关注的是close/transition 如何解除waiter,而不是request routing。
Core最内层 Session::deliver_event_raw() 在 tx_event.send() 失败时只debug并继续,但当前没有专门测试drop CodexThread event receiver后模型Task仍完成且rollout仍写入。该结论由实现和队列owner推导,保留为测试 缺口;不能用client facade的Lagged测试替代Core receiver-close测试。
12. 背压与关闭
- 从
Session::send_event_raw_with_persistence()开始,依次说出rollout、Core event queue、App Server listener、connection fan-out和client facade;指出每层的容量与失败策略。 - 解释为什么AgentMessageDelta必须lossless,而CommandExecutionOutputDelta可以best-effort;再说明Lagged 能恢复什么、不能恢复什么。
- 给定“UI缺一段stdout但最终回答完整”和“UI永远等不到TurnComplete”两种现象,分别定位最可能违反的 分类或队列边界。
- 使用只读命令核对分类、转发和pending request清理:
rg -n "server_notification_requires_delivery|forward_in_process_event|Lagged" \
codex-rs/app-server-client/src/lib.rs
rg -n "abort_pending_server_requests|cancel_requests_for_thread|deliver_event_raw" \
codex-rs/app-server/src codex-rs/core/src/session/mod.rs继续学习 Session运行时处理,理解哪些handler直接发送raw event; 再阅读 Session关闭流程,观察event receiver关闭为何不能成为 teardown的必要条件。
