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
/// 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
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
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
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
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
/// 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
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
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
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
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
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
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. 源码定位
遇到“事件没有显示”,先沿四个问题定位:
- handler 是否调用
send_event,还是只写了persist_rollout_items; send_event_raw_with_persistence的persist是否为 true,rollout append 是否报错;agent_status_from_event是否会更新 watch status;- event channel、history reducer 或 App Server mapper 是否过滤/投影了该 variant。
遇到“命令结束但 Turn 仍运行”,检查 ExecCommandEnd、RawResponseCompleted 和 TurnComplete 的层级;遇到“恢复后历史出现但实时 UI 没有”,分别检查 rollout append 和 tx_event.send。完成这条路径后,可以继续阅读协议映射与事件转换深入研究非一对一的 App Server projection。
