Skip to content

CodexThread背压

解析 CodexThread 事件从 Core 无界队列到 App Server fan-out 的持久化、背压、丢弃和关闭语义。

基于rust-v0.150.0
CodexRustCodexThreadBackpressure

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

rust
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

rust
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(节选)

rust
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

rust
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

rust
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 被 dropSession send 失败只记 debugCore task 继续,不自动 shutdown
Session 所有 event sender dropnext_event() 返回 InternalAgentDiedApp Server listener 记录 warning 并退出
App Server subscriber 为空thread-scoped sender 直接返回reducer仍可更新,连接无通知
Outgoing mpsc 关闭send 记录 warning当前 notification 丢失
client event channel 关闭DisableStreamworker 停止轮询 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. 事件代码审阅 ​

  1. 这条消息在哪一层产生:Core Event、ServerNotification、ServerRequest 还是 client marker?
  2. 它是否必须 durable;Legacy 与 Paginated 的策略是否相同?
  3. queue 是 bounded、unbounded、watch 还是 oneshot;满载与关闭分别怎样处理?
  4. 谁是唯一消费者,fan-out 在消费前还是消费后发生?
  5. 丢失这条消息会破坏 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

rust
// :: 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

rust
// :: 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

rust
// :: 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. 背压与关闭 ​

  1. 从 Session::send_event_raw_with_persistence() 开始,依次说出rollout、Core event queue、App Server listener、connection fan-out和client facade;指出每层的容量与失败策略。
  2. 解释为什么AgentMessageDelta必须lossless,而CommandExecutionOutputDelta可以best-effort;再说明Lagged 能恢复什么、不能恢复什么。
  3. 给定“UI缺一段stdout但最终回答完整”和“UI永远等不到TurnComplete”两种现象,分别定位最可能违反的 分类或队列边界。
  4. 使用只读命令核对分类、转发和pending request清理:
bash
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的必要条件。