Skip to content

Bespoke事件处理

从 Core EventMsg 到 App Server 通知、反向请求和 Turn 终态,解析不能直接一一映射的事件如何补状态、聚合和抑制重复。

基于rust-v0.150.0
CodexRustAppServerEventProtocol

Bespoke事件处理 ​

Core 的 EventMsg 是运行时事件流,App Server 的 ServerNotification 却是稳定的对外协议;两者并非一一对应。一个事件可能需要读取 ThreadState 才能补齐字段,另一个事件只为兼容旧消费者而被丢弃,ItemStarted 还可能同时触发一个发给客户端的反向 request。本文面向了解 Rust enum、异步 channel 和锁的读者,承接AppServer架构总览、V2请求分派总表和服务端请求与客户端响应。

本文围绕 apply_bespoke_event_handling 的真实分派,解释事件的所有者、状态变化、协议消费者和失败路径;不重复讲每个 processor 的请求入口。

1. 映射边界 ​

EventMsg 的 id 是 Core 事件所在 turn,外层 conversation_id 是 Thread。映射函数必须显式传入两者,否则实时事件或历史回放会把 turn 关联错。

事件处理器同时读写 Thread 状态并驱动出站协议,不能当作无副作用的 enum 转换函数。

同一个 Core 事件可能只更新状态、生成通知,或启动需要客户端回答的反向请求。

2. 入口分派 ​

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

相关函数:apply_bespoke_event_handling。

rust
pub(crate) async fn apply_bespoke_event_handling(
    event: Event,
    conversation_id: ThreadId,
    conversation: Arc<CodexThread>,
    thread_manager: Arc<ThreadManager>,
    outgoing: ThreadScopedOutgoingMessageSender,
    thread_state: Arc<tokio::sync::Mutex<ThreadState>>,
    thread_watch_manager: ThreadWatchManager,
    thread_list_state_permit: Arc<tokio::sync::Semaphore>,
    fallback_model_provider: String,
) {
    let Event { id: event_turn_id, msg } = event;
    match msg {
        EventMsg::TurnStarted(payload) => { /* 建立当前 turn 快照并通知 */ }
        EventMsg::TurnComplete(payload) => { /* 聚合 summary 后通知终态 */ }
        EventMsg::ItemStarted(event) => { /* 可能发送动态工具请求 */ }
        EventMsg::ItemCompleted(event) => { /* canonical item 的副作用 */ }
        EventMsg::Error(error) => { /* 区分回滚错误和 turn 错误 */ }
        _ => { /* 每个分支定义自己的协议边界 */ }
    }
}

这里的参数列表本身就是设计证据:conversation 用于读取运行时扩展数据,thread_state 拥有当前 turn 的可变摘要,outgoing 是唯一对外发送者,thread_watch_manager 维护 Thread 级观察状态。函数不是纯转换器。

3. Turn终态 ​

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

相关函数:handle_turn_complete、handle_turn_interrupted、find_and_remove_turn_summary。

rust
async fn handle_turn_complete(
    conversation_id: ThreadId,
    event_turn_id: String,
    turn_complete_event: TurnCompleteEvent,
    outgoing: &ThreadScopedOutgoingMessageSender,
    thread_state: &Arc<Mutex<ThreadState>>,
) {
    let turn_summary = find_and_remove_turn_summary(conversation_id, thread_state).await;
    let (status, error, last_agent_message) = match turn_summary.last_error {
        Some(error) => (TurnStatus::Failed, Some(error), None),
        None => (TurnStatus::Completed, None, turn_summary.last_agent_message),
    };
    emit_turn_completed_with_status(
        conversation_id, event_turn_id,
        TurnCompletionMetadata {
            status, error, last_agent_message,
            started_at: turn_summary.started_at,
            completed_at: turn_complete_event.completed_at,
            duration_ms: turn_complete_event.duration_ms,
        }, outgoing,
    ).await;
}

TurnCompleteEvent 自身没有最终错误摘要;错误先由 EventMsg::Error 写入 turn_summary.last_error,终态处理再取出并清空摘要。std::mem::take 让同一个 Thread 的下一轮从空摘要开始,避免上一轮错误泄漏。中断路径固定输出 Interrupted,不会把取消伪装成失败。

StreamError 回到 Running,普通 Error 才进入终态摘要;这是客户端判断是否等待重试的关键。

4. 直接通知 ​

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

rust
EventMsg::TokenCount(token_count_event) => {
    handle_token_count_event(
        conversation_id, event_turn_id, token_count_event, &outgoing
    ).await;
}

async fn handle_token_count_event(
    conversation_id: ThreadId,
    turn_id: String,
    token_count_event: TokenCountEvent,
    outgoing: &ThreadScopedOutgoingMessageSender,
) {
    let TokenCountEvent { info, rate_limits } = token_count_event;
    if let Some(token_usage) = info.map(ThreadTokenUsage::from) {
        outgoing.send_server_notification(
            ServerNotification::ThreadTokenUsageUpdated(
                ThreadTokenUsageUpdatedNotification {
                    thread_id: conversation_id.to_string(), turn_id, token_usage,
                }
            )
        ).await;
    }
    if let Some(rate_limits) = rate_limits {
        outgoing.send_server_notification(
            ServerNotification::AccountRateLimitsUpdated(
                AccountRateLimitsUpdatedNotification { rate_limits: rate_limits.into() }
            )
        ).await;
    }
}

一个 Core 事件可以产生两条不同范围的通知:token usage 带 Thread/Turn,rate limit 没有 Thread ID。消费者不能通过“同一事件”假设两条通知具有相同关联键。

5. 错误与重试 ​

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

rust
EventMsg::Error(ev) => {
    let codex_error_info = ev.codex_error_info.clone();
    if matches!(codex_error_info,
        Some(CoreCodexErrorInfo::ThreadRollbackFailed)) {
        return handle_thread_rollback_failed(
            conversation_id, ev.message, &thread_state, &outgoing
        ).await;
    }
    if !ev.affects_turn_status() { return; }
    handle_error_notification(
        conversation_id, &event_turn_id,
        TurnError { message: ev.message,
            codex_error_info: ev.codex_error_info.map(Into::into),
            additional_details: None },
        &outgoing, &thread_state,
    ).await;
}

EventMsg::StreamError(ev) => {
    outgoing.send_server_notification(ServerNotification::Error(
        ErrorNotification { error: ev.into(), will_retry: true,
            thread_id: conversation_id.to_string(), turn_id: event_turn_id }
    )).await;
}

回滚错误只完成 pending rollback request,不发送普通 turn error;不影响 turn 状态的错误直接结束分支;stream error 是中间重试状态,通知客户端但不写入终态摘要。三条路径对应三种不同消费者,不能统一成“收到 Error 就失败”。

6. Item与反向请求 ​

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

rust
EventMsg::ItemStarted(event) => {
    let should_emit = match &event.item {
        CoreTurnItem::CommandExecution(item) => thread_state.lock().await
            .turn_summary.command_execution_started.insert(item.id.clone()),
        _ => true,
    };
    if should_emit {
        outgoing.send_server_notification(item_event_to_server_notification(
            EventMsg::ItemStarted(event), &conversation_id.to_string(), &event_turn_id
        )).await;
    }
    if let Some(params) = dynamic_tool_call_params {
        let call_id = params.call_id.clone();
        let (_id, rx) = outgoing.send_request(
            ServerRequestPayload::DynamicToolCall(params)
        ).await;
        tokio::spawn(async move {
            crate::dynamic_tools::on_call_response(call_id, rx, conversation).await;
        });
    }
}

审批/Guardian 流程可能先发命令开始事件,canonical item 随后才到;command_execution_started 集合因此承担去重。动态工具则是“通知 + 反向 request”两条并行路径,response 由独立 task 消费,不能阻塞事件循环等待客户端。

7. 有意忽略 ​

源码文件:codex-rs/app-server/src/bespoke_event_handling.rs

rust
EventMsg::ContextCompacted(..) => {
    // Core 仍为 raw-event/rollout 兼容消费者广播;
    // v2 客户端接收 canonical ContextCompaction item。
}
EventMsg::McpToolCallBegin(_) | EventMsg::McpToolCallEnd(_) => {
    // v2 使用 TurnItem::McpToolCall 生命周期,此处保留旧消费者兼容。
}

空分支不是遗漏,而是协议版本选择:旧事件仍进入 raw-event 或 rollout,v2 客户端改读 canonical item。把这些事件再次发送会造成重复时间线;这也是本文明确不覆盖旧 raw-event 消费者内部实现的边界。

兼容分支保留旧消费者所需的数据,同时避免 v2 时间线重复。

8. 复核路径 ​

源码测试覆盖三类证明:test_handle_turn_complete_emits_completed_without_error 检查无错误摘要时的 Completed 负载;test_handle_turn_complete_emits_failed_with_error 先注入错误再断言终态携带错误;test_handle_turn_interrupted_emits_interrupted_without_error 验证中断不被映射为 Failed。实时历史测试还检查 V1/V2/V3 事件版本的增量与关闭通知。

bash
rg -n "apply_bespoke_event_handling|handle_turn_complete|EventMsg::" codex-rs/app-server/src/bespoke_event_handling.rs
cargo test -p codex-app-server bespoke_event_handling

遇到客户端重复 item,先查去重集合而非协议解码;遇到 Turn 状态错误,沿 EventMsg::Error → turn_summary → handle_turn_complete 追踪;遇到动态工具卡住,检查 pending request callback 和 spawned response task。这样才能把 Core 事件、Thread 状态和对外通知放回各自的生命周期。

去重和动态工具 callback 是两条独立的生命周期;一个被抑制不代表另一个被取消。