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。
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。
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
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
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
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
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 事件版本的增量与关闭通知。
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 是两条独立的生命周期;一个被抑制不代表另一个被取消。
