Skip to content

Event与EventMsg总表

从 Session 事件发布入口,跟读 EventMsg 的持久化、状态投影、历史 reducer、App Server 映射和终止语义。

基于rust-v0.150.0
CodexRustProtocolEvent

Event与EventMsg总表 ​

本文承接Submission与Op总表,面向已经理解 submission queue、准备阅读 EQ 的读者。本文回答一个 EventMsg 从 session 产生后会经过哪些消费者、哪些事件进入 rollout、何时更新 AgentStatus、为什么 RawResponseCompleted 不等于 TurnComplete,以及 channel 关闭时实时消费和持久化会怎样分离。不展开每个工具的完整业务实现,也不把 App Server notification 当成内部事件的同义词。

在 rust-v0.150.0 中,事件至少有四种用途:Core 内部状态推进、rollout/history 持久化、实时 event channel、上层 App Server/TUI 投影。一个 EventMsg 可以被多个消费者读取,但它们的筛选和生效时机不同。理解这一点,比背诵几十个 enum variant 更能帮助定位“事件丢了”“状态没更新”或“恢复后历史不一致”。

1. 事件外壳 ​

Event 只有关联 id 和 EventMsg。id 通常是产生该事件的 submission/turn id,但事件本身不携带一个“唯一业务响应”假设;一个 Turn 会产生开始、内容、工具、诊断和终止等多条事件。

源码位置:codex-rs/protocol/src/protocol.rs :: Event、EventMsg

rust
/// Event Queue Entry - events from agent
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct Event {
    /// Submission `id` that this event is correlated with.
    pub id: String,
    /// Payload
    pub msg: EventMsg,
}

/// Response event from the agent
/// NOTE: Make sure none of these values have optional types, as it will mess up the extension code-gen.
#[derive(Debug, Clone, Deserialize, Serialize, Display, JsonSchema, TS)]
#[serde(tag = "type", rename_all = "snake_case")]
#[ts(tag = "type")]
#[strum(serialize_all = "snake_case")]
pub enum EventMsg {
    Error(ErrorEvent),
    Warning(WarningEvent),
    GuardianWarning(WarningEvent),
    RealtimeConversationStarted(RealtimeConversationStartedEvent),
    RealtimeConversationRealtime(RealtimeConversationRealtimeEvent),
    RealtimeConversationClosed(RealtimeConversationClosedEvent),
    RealtimeConversationSdp(RealtimeConversationSdpEvent),
    ModelReroute(ModelRerouteEvent),
    ModelVerification(ModelVerificationEvent),
    TurnModerationMetadata(TurnModerationMetadataEvent),
    SafetyBuffering(SafetyBufferingEvent),
    ContextCompacted(ContextCompactedEvent),
    ThreadRolledBack(ThreadRolledBackEvent),
    #[serde(rename = "task_started", alias = "turn_started")]
    TurnStarted(TurnStartedEvent),
    ThreadSettingsApplied(ThreadSettingsAppliedEvent),
    #[serde(rename = "task_complete", alias = "turn_complete")]
    TurnComplete(TurnCompleteEvent),
}

这段 enum 是源码中从错误到 TurnComplete 的连续片段。后续变体还包括消息、工具、审批、TurnAborted、ShutdownComplete 和 raw response 事件。两个兼容点是:Rust variant TurnStarted 在旧 v1 wire 上序列化为 task_started,同时接受 turn_started;TurnComplete 同理使用 task_complete 和 turn_complete。EventMsg 不能被当作“每个 variant 都会进入 rollout 或 App Server”的保证。

2. 发布入口 ​

Session::send_event 是带 Turn 上下文的主要入口。它先为影响 Turn status 的 Error 写入 terminal_error,再记录 Codex turn trace 和 tool-call trace,构造 Event,进入 raw 发布;随后处理 MultiAgentV2 parent completion、realtime handoff,并可能发布 legacy projection events。

源码位置:codex-rs/core/src/session/mod.rs :: Session::send_event

rust
pub(crate) async fn send_event(&self, turn_context: &TurnContext, msg: EventMsg) {
    let legacy_source = msg.clone();
    if let EventMsg::Error(error) = &legacy_source
        && error
            .codex_error_info
            .as_ref()
            .is_some_and(CodexErrorInfo::affects_turn_status)
    {
        turn_context
            .terminal_error
            .lock()
            .await
            .replace(error.clone());
    }
    self.services
        .rollout_thread_trace
        .record_codex_turn_event(&turn_context.sub_id, &legacy_source);
    self.services
        .rollout_thread_trace
        .record_tool_call_event(turn_context.sub_id.clone(), &legacy_source);
    let event = Event {
        id: turn_context.sub_id.clone(),
        msg,
    };
    self.send_event_raw(event).await;
    self.maybe_notify_parent_of_terminal_turn(turn_context, &legacy_source)
        .await;
    self.maybe_mirror_event_text_to_realtime(&legacy_source)
        .await;
    self.maybe_clear_realtime_handoff_for_event(&legacy_source)
        .await;
}

legacy_source 的 clone 不是多余复制:raw event 发送会消费 msg,而 parent/realtime/legacy projection 仍需要读取同一个逻辑事件。Terminal error 也必须在 raw delivery 前写入,后续 parent completion 才能把错误投影为 AgentStatus::Errored。

3. 持久化与投递 ​

raw 发布分成两个阶段:先由 MCP runtime observer 观察,再按 persist 选择 append rollout,然后记录 protocol trace,最后投递到 tx_event。deliver_event_raw 同时把影响状态的事件投影到 watch channel;event channel 关闭只会丢实时消息,不会撤销之前完成的持久化和 trace。

源码位置: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) {
    self.services.mcp_runtime.observe_event(&event.msg);
    // Persist the event into rollout storage; the store applies its persistence policy.
    if persist {
        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) {
    // Record the last known agent status.
    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 {
        debug!("dropping event because channel is closed: {e}");
    }
}

这里的 persist_rollout_items 本身还会检查 live thread 和 append 错误;append 失败记录 error,但不会阻塞 event channel 的投递。实时消费者因此可能看到一条事件,而 rollout 没有同一条记录,排查时必须分别检查两条路径。

4. 延迟物化 ​

某些线程尚未 materialize rollout,却需要发送 settings、诊断或连接事件。send_event_raw_without_materializing_rollout 先探测当前 rollout path 是否已经存在;没有现有文件时才允许不创建本地 rollout,探测失败则选择 persist = true,避免把 IO 不确定性解释成“无需持久化”。

源码位置:codex-rs/core/src/session/mod.rs :: send_event_raw_without_materializing_rollout

rust
pub(crate) async fn send_event_raw_without_materializing_rollout(&self, event: Event) {
    let persist = match self.current_rollout_path().await {
        Ok(Some(path)) => codex_rollout::existing_rollout_path(&path).await.is_some(),
        Ok(None) => true,
        Err(err) => {
            warn!("failed to check whether thread persistence is materialized: {err}");
            true
        }
    };
    self.send_event_raw_with_persistence(event, persist).await;
}

这个入口不是“无副作用 send”。如果 path 检查返回 None,persist = true 可能触发 live thread append;只有已经存在的 rollout 文件才会沿用现有持久化状态。设置事件的调用方必须知道自己是否允许首次物化。

5. 状态投影 ​

agent_status_from_event 只关心少数事件:TurnStarted 进入 Running,TurnComplete 根据 error 进入 Completed 或 Errored,TurnAborted 根据 reason 区分 Interrupted/Errored,ShutdownComplete 进入 Shutdown。普通 message、tool delta、Warning 和 TokenCount 不改变这个 watch 状态。

源码位置:codex-rs/core/src/agent/status.rs :: agent_status_from_event

rust
pub(crate) fn agent_status_from_event(msg: &EventMsg) -> Option<AgentStatus> {
    match msg {
        EventMsg::TurnStarted(_) => Some(AgentStatus::Running),
        EventMsg::TurnComplete(ev) => Some(match &ev.error {
            Some(error) => AgentStatus::Errored(error.message.clone()),
            None => AgentStatus::Completed(ev.last_agent_message.clone()),
        }),
        EventMsg::TurnAborted(ev) => match ev.reason {
            TurnAbortReason::Interrupted | TurnAbortReason::BudgetLimited => {
                Some(AgentStatus::Interrupted)
            }
            _ => Some(AgentStatus::Errored(format!("{:?}", ev.reason))),
        },
        EventMsg::Error(ev) => Some(AgentStatus::Errored(ev.message.clone())),
        EventMsg::ShutdownComplete => Some(AgentStatus::Shutdown),
        _ => None,
    }
}

因此 Error 可以先让 status 变成 Errored,随后同一 Turn 的 TurnComplete { error: None } 又可能投影为 Completed;是否影响 Turn 的最终状态还取决于 ErrorEvent.codex_error_info.affects_turn_status 和 Turn 收尾逻辑,不能只看某一条事件。

6. 终止事件 ​

6.1 Turn与Step ​

ItemCompleted、ExecCommandEnd、RawResponseCompleted 都是局部完成信号;TurnComplete 和 TurnAborted 才是整个 agent Turn 的终端事件。RawResponseCompleted 记录一次 upstream Responses response 的精确 id/usage,之后 Core 还要 flush assistant segments、更新预算、判断 follow-up,才可能产生 TurnComplete。

源码位置:codex-rs/protocol/src/protocol.rs :: RawResponseCompletedEvent、TurnCompleteEvent、TurnAbortedEvent

rust
/// Exact usage reported by one upstream Responses API completion.
#[derive(Debug, Clone, Deserialize, Serialize, TS, JsonSchema)]
pub struct RawResponseCompletedEvent {
    pub response_id: String,
    pub token_usage: Option<TokenUsage>,
}

#[derive(Debug, Clone, Deserialize, Serialize, JsonSchema, TS)]
pub struct TurnCompleteEvent {
    pub turn_id: String,
    pub last_agent_message: Option<String>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub error: Option<ErrorEvent>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub started_at: Option<i64>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub completed_at: Option<i64>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub duration_ms: Option<i64>,
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub time_to_first_token_ms: Option<i64>,
}

#[derive(Debug, Clone, Deserialize, Serialize, JsonSchema, TS)]
pub struct TurnAbortedEvent {
    pub turn_id: Option<String>,
    pub started_at: Option<i64>,
    pub reason: TurnAbortReason,
    pub completed_at: Option<i64>,
    pub duration_ms: Option<i64>,
}

6.2 Status machine ​

状态图表达的是 Core 消费关系,不是 EventMsg enum 的自动状态机。比如 ExecCommandEnd 可以回到 Running,TurnAborted 直接终止,RawResponseCompleted 还可能触发下一轮 sampling。

7. 历史Reducer ​

ThreadHistoryBuilder 是有状态 reducer,不是 EventMsg -> JSON 的纯映射。它维护 turns、current turn、item index 和 change set,按事件类型更新 Thread/Turn/ThreadItem;无法进入 rollout history 的 WorldState、RealtimeItem 和 SecurityRiskScore 在 handle_rollout_item 中跳过。

源码位置:codex-rs/app-server-protocol/src/protocol/thread_history.rs :: ThreadHistoryBuilder::handle_event、handle_event_with_changes

rust
pub fn handle_event(&mut self, event: &EventMsg) {
    match event {
        EventMsg::UserMessage(payload) => self.handle_user_message(payload),
        EventMsg::AgentMessage(payload) => self.handle_agent_message(payload),
        EventMsg::AgentReasoning(payload) => self.handle_agent_reasoning(payload),
        EventMsg::ExecCommandBegin(payload) => self.handle_exec_command_begin(payload),
        EventMsg::ExecCommandEnd(payload) => self.handle_exec_command_end(payload),
        EventMsg::GuardianAssessment(payload) => self.handle_guardian_assessment(payload),
        EventMsg::DynamicToolCallRequest(payload) => {
            self.handle_dynamic_tool_call_request(payload)
        }
        EventMsg::DynamicToolCallResponse(payload) => {
            self.handle_dynamic_tool_call_response(payload)
        }
        EventMsg::ItemStarted(payload) => self.handle_item_started(payload),
        EventMsg::ItemCompleted(payload) => self.handle_item_completed(payload),
        EventMsg::TurnAborted(payload) => self.handle_turn_aborted(payload),
        EventMsg::TurnStarted(payload) => self.handle_turn_started(payload),
        EventMsg::TurnComplete(payload) => self.handle_turn_complete(payload),
        _ => {}
    }
}

pub fn handle_event_with_changes(&mut self, event: &EventMsg) -> ThreadHistoryChangeSet {
    self.collect_changes(|builder| builder.handle_event(event))
}

handle_event_with_changes 的 change set 给 running thread resume/rejoin 使用;同一事件可以先改变 current item,再由 accumulator 合并成一次对外变化。若只看最终 Thread JSON,会漏掉这个增量生效时机。

8. App Server映射 ​

App Server 的 item_event_to_server_notification 只负责有一对一关系的 item 事件。EventMsg::ExecCommandOutputDelta 被投影为 CommandExecutionOutputDeltaNotification;ExecCommandEnd 通过 item builder 生成 ItemCompleted。Warning、TurnComplete 和复杂状态事件由其他 adapter 处理。

源码位置:codex-rs/app-server-protocol/src/protocol/event_mapping.rs :: item_event_to_server_notification

rust
EventMsg::ItemStarted(item_started_event) => {
    ServerNotification::ItemStarted(ItemStartedNotification {
        thread_id,
        turn_id,
        item: item_started_event.item.into(),
        started_at_ms: item_started_event.started_at_ms,
    })
}
EventMsg::ItemCompleted(item_completed_event) => {
    ServerNotification::ItemCompleted(ItemCompletedNotification {
        thread_id,
        turn_id,
        item: item_completed_event.item.into(),
        completed_at_ms: item_completed_event.completed_at_ms,
    })
}
EventMsg::ExecCommandOutputDelta(event) => {
    ServerNotification::CommandExecutionOutputDelta(
        CommandExecutionOutputDeltaNotification {
            thread_id,
            turn_id,
            item_id: event.call_id,
            delta: String::from_utf8_lossy(&event.chunk).to_string(),
        },
    )
}

这个 helper 的 _ => unreachable!("unsupported item event") 是刻意的边界:调用方必须先筛选 item event,不能把完整 EQ 当作 App Server notification。App Server notification 还会补 thread/turn 字段,和内部 Event.id 不是同一个关联字段。

9. 跨线程消费者 ​

9.1 子Agent完成 ​

MultiAgentV2 子线程的 TurnComplete/TurnAborted 会被 parent session 观察。若 terminal_error 存在,parent 收到 AgentStatus::Errored;否则从 terminal event 计算 status,再通过 agent control 投影给父线程。普通 child event 不会全部广播给 parent。

源码位置:codex-rs/core/src/session/mod.rs :: maybe_notify_parent_of_terminal_turn

rust
if turn_context.multi_agent_version != MultiAgentVersion::V2 {
    return;
}

if !matches!(msg, EventMsg::TurnComplete(_) | EventMsg::TurnAborted(_)) {
    return;
}

let status = match turn_context.terminal_error.lock().await.take() {
    Some(error) => {
        let status = AgentStatus::Errored(error.message);
        self.agent_status.send_replace(status.clone());
        status
    }
    None => {
        let Some(status) = agent_status_from_event(msg) else {
            return;
        };
        status
    }
};

9.2 Realtime handoff ​

Realtime handoff 只在 TurnComplete 时清理 active handoff;message mirror 和 close 由不同 helper 处理。因而同一 EventMsg 的实时文本镜像和 Turn 收尾是两个独立消费者,不能按 event channel 到达顺序推断全部 handoff 已完成。

源码位置:codex-rs/core/src/session/mod.rs :: maybe_clear_realtime_handoff_for_event

rust
async fn maybe_clear_realtime_handoff_for_event(&self, msg: &EventMsg) {
    if !matches!(msg, EventMsg::TurnComplete(_)) {
        return;
    }
    if let Err(err) = self.conversation.handoff_complete().await {
        debug!("failed to finalize realtime handoff output: {err}");
    }
    self.conversation.clear_active_handoff().await;
}

10. 测试与验证 ​

第一组测试验证 status projection:输入 TurnStarted、带 error/不带 error 的 TurnComplete、Interrupted TurnAborted 和 ShutdownComplete,断言分别得到 Running、Errored/Completed、Interrupted、Shutdown。这些测试只覆盖 agent_status_from_event,不证明事件已经持久化。

源码位置:codex-rs/core/src/agent/control_tests.rs :: on_event_updates_status_from_task_started、on_event_updates_status_from_task_complete、on_event_updates_status_from_turn_aborted、on_event_updates_status_from_shutdown_complete

rust
let status = agent_status_from_event(&EventMsg::TurnStarted(TurnStartedEvent {
    turn_id: "turn-1".to_string(),
    trace_id: None,
    started_at: None,
    model_context_window: None,
    collaboration_mode_kind: ModeKind::Default,
}));
assert_eq!(status, Some(AgentStatus::Running));

第二组测试验证 channel-close teardown:关闭 submission channel 后,session 仍调用 thread-stop contributor 和 live thread shutdown;活动 Turn 场景额外断言 abort lifecycle 先于 thread-stop。它反向证明“事件/状态收尾”和“实时 channel 是否仍有 receiver”是两件事。

源码位置:codex-rs/core/src/session/tests.rs :: submission_loop_channel_close_runs_full_thread_teardown、submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle

text
cd codex-rs
cargo test -p codex-core --lib on_event_updates_status_from_task_started -- --nocapture --test-threads=1
cargo test -p codex-core --lib on_event_updates_status_from_task_complete -- --nocapture --test-threads=1
cargo test -p codex-core --lib on_event_updates_status_from_turn_aborted -- --nocapture --test-threads=1
cargo test -p codex-core --lib on_event_updates_status_from_shutdown_complete -- --nocapture --test-threads=1
cargo test -p codex-core --lib submission_loop_channel_close_runs_full_thread_teardown -- --nocapture --test-threads=1
cargo test -p codex-core --lib submission_loop_channel_close_aborts_active_turn_before_thread_stop_lifecycle -- --nocapture --test-threads=1

这些测试证明状态映射和关闭顺序;不能证明每个 EventMsg variant 都进入 history、App Server mapper 都保留字段,或远端客户端一定在线接收。

11. 源码定位 ​

遇到“事件没有显示”,先沿四个问题定位:

  1. handler 是否调用 send_event,还是只写了 persist_rollout_items;
  2. send_event_raw_with_persistence 的 persist 是否为 true,rollout append 是否报错;
  3. agent_status_from_event 是否会更新 watch status;
  4. event channel、history reducer 或 App Server mapper 是否过滤/投影了该 variant。

遇到“命令结束但 Turn 仍运行”,检查 ExecCommandEnd、RawResponseCompleted 和 TurnComplete 的层级;遇到“恢复后历史出现但实时 UI 没有”,分别检查 rollout append 和 tx_event.send。完成这条路径后,可以继续阅读协议映射与事件转换深入研究非一对一的 App Server projection。