Skip to content

协议映射与事件转换

追踪 Core EventMsg 到 V2 notification、item 和 server request 的单向映射与信息边界。

基于rust-v0.150.0
CodexRustProtocolMappingEvents

协议映射与事件转换 ​

本文承接Thread History Projection和V2 Turn与Item协议,面向已经理解 Core EventMsg 与 App Server notification 的读者。本文回答事件如何被映射、哪些旧事件被忽略、哪些事件会触发反向 request,以及为什么映射不是无损转换;不展开每个 notification 的完整字段。

1. 映射总线 ​

Core 发布 EventMsg,App Server 根据事件种类选择 V2 notification、server request 或兼容性转发。映射函数同时接收 conversation/thread ID 和 event turn ID,因此外部 ID 不完全来自事件 payload 本身。

2. 生命周期映射 ​

源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: EventMsg::ItemStarted、EventMsg::ItemCompleted

rust
EventMsg::ItemStarted(event) => {
    let notification = item_event_to_server_notification(
        EventMsg::ItemStarted(event),
        &conversation_id.to_string(),
        &event_turn_id,
    );
    outgoing.send_server_notification(notification).await;
}
EventMsg::ItemCompleted(event) => {
    apply_canonical_item_completed_side_effects(
        &thread_manager,
        &thread_watch_manager,
        &thread_state,
        &event.item,
    ).await;
    let notification = item_event_to_server_notification(
        EventMsg::ItemCompleted(event),
        &conversation_id.to_string(),
        &event_turn_id,
    );
    outgoing.send_server_notification(notification).await;
}

ItemCompleted 先应用 thread state side effect,再发送 notification;完成事件的外部观察可能依赖这些状态已经更新。映射器不是简单的 serde_json rename。

3. 流式事件 ​

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

rust
msg @ (EventMsg::AgentMessageContentDelta(_)
    | EventMsg::PlanDelta(_)
    | EventMsg::ReasoningContentDelta(_)
    | EventMsg::ReasoningRawContentDelta(_)) => {
    let notification = item_event_to_server_notification(
        msg,
        &conversation_id.to_string(),
        &event_turn_id,
    );
    outgoing.send_server_notification(notification).await;
}

文本、计划和 reasoning 增量共享 mapper,但输出 notification 类型不同。它们携带同一 thread/turn 关联,却不代表同一种 item。

4. 一对一事件映射 ​

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

rust
pub fn item_event_to_server_notification(
    msg: EventMsg,
    thread_id: &str,
    turn_id: &str,
) -> ServerNotification {
    let thread_id = thread_id.to_string();
    let turn_id = turn_id.to_string();
    match msg {
        EventMsg::AgentMessageContentDelta(event) => {
            let codex_protocol::protocol::AgentMessageContentDeltaEvent {
                item_id, delta, ..
            } = event;
            ServerNotification::AgentMessageDelta(AgentMessageDeltaNotification {
                thread_id,
                turn_id,
                item_id,
                delta,
            })
        }
        EventMsg::PlanDelta(event) => ServerNotification::PlanDelta(PlanDeltaNotification {
            thread_id,
            turn_id,
            item_id: event.item_id,
            delta: event.delta,
        }),
        EventMsg::ReasoningContentDelta(event) => {
            ServerNotification::ReasoningSummaryTextDelta(
                ReasoningSummaryTextDeltaNotification {
                    thread_id,
                    turn_id,
                    item_id: event.item_id,
                    delta: event.delta,
                    summary_index: event.summary_index,
                },
            )
        }
        EventMsg::ItemStarted(event) => {
            ServerNotification::ItemStarted(ItemStartedNotification {
                thread_id,
                turn_id,
                item: event.item.into(),
                started_at_ms: event.started_at_ms,
            })
        }
        EventMsg::ItemCompleted(event) => {
            ServerNotification::ItemCompleted(ItemCompletedNotification {
                thread_id,
                turn_id,
                item: event.item.into(),
                completed_at_ms: 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(),
                },
            )
        }
        _ => unreachable!("unsupported item event"),
    }
}

mapper 只接受能够由单个事件直接构造的 item/stream notification;它不负责去重、权限 guard、pending request 或 Turn summary。比如 command output 的 bytes 在 mapper 中转换为 lossy UTF-8 文本,完整的 cap 和原始字节语义由其他协议路径承担。

5. 反向请求 ​

动态工具、permissions、MCP elicitation 等事件不会只发送 notification。App Server 会把 Core 事件转换成带 callback 的 server request,客户端响应再回到 Core 等待者。

源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: EventMsg::RequestPermissions、EventMsg::DynamicToolCallRequest

rust
let (pending_request_id, rx) = outgoing
    .send_request(ServerRequestPayload::PermissionsRequestApproval(params))
    .await;
tokio::spawn(async move {
    on_request_permissions_response(
        pending_request,
        conversation,
        thread_state,
    ).await;
});

这里的 pending request ID 属于 App Server connection callback,Core event 的 call_id/item_id 属于业务关联。两套 ID 必须分别保存。

6. 旧事件处理 ​

源码位置:codex-rs/app-server/src/bespoke_event_handling.rs :: deprecated event arms

rust
EventMsg::ExecCommandBegin(_)
| EventMsg::ExecCommandEnd(_)
| EventMsg::EnteredReviewMode(_)
| EventMsg::ExitedReviewMode(_) => {
    // V2 clients receive TurnItem lifecycle instead.
}
EventMsg::McpToolCallBegin(_) | EventMsg::McpToolCallEnd(_) => {
    // V2 receives canonical TurnItem::McpToolCall lifecycle.
}

旧事件被忽略并不代表 Core 没有产生它们,而是 V2 选择 canonical item 视图;raw-event 和 rollout 兼容消费者仍可能需要这些事件。

源码位置:codex-rs/app-server-protocol/src/protocol/item_builders.rs :: build_command_execution_begin_item、build_command_execution_end_item

rust
pub fn build_command_execution_begin_item(payload: &ExecCommandBeginEvent) -> ThreadItem {
    let presentation =
        CommandExecutionPresentation::from_raw(&payload.command, &payload.parsed_cmd, &payload.cwd);
    ThreadItem::CommandExecution {
        id: payload.call_id.clone(),
        plugin_id: payload.plugin_id.clone(),
        script_path: payload.script_path.clone(),
        command: presentation.command,
        cwd: payload.cwd.clone().into(),
        process_id: payload.process_id.clone(),
        source: payload.source.into(),
        status: CommandExecutionStatus::InProgress,
        command_actions: presentation.command_actions,
        aggregated_output: None,
        exit_code: None,
        duration_ms: None,
    }
}

pub fn build_command_execution_end_item(payload: &ExecCommandEndEvent) -> ThreadItem {
    let aggregated_output = if payload.aggregated_output.is_empty() {
        None
    } else {
        Some(payload.aggregated_output.clone())
    };
    let duration_ms = i64::try_from(payload.duration.as_millis()).unwrap_or(i64::MAX);
    let presentation =
        CommandExecutionPresentation::from_raw(&payload.command, &payload.parsed_cmd, &payload.cwd);
    ThreadItem::CommandExecution {
        id: payload.call_id.clone(),
        plugin_id: payload.plugin_id.clone(),
        script_path: payload.script_path.clone(),
        command: presentation.command,
        cwd: payload.cwd.clone().into(),
        process_id: payload.process_id.clone(),
        source: payload.source.into(),
        status: (&payload.status).into(),
        command_actions: presentation.command_actions,
        aggregated_output,
        exit_code: Some(payload.exit_code),
        duration_ms: Some(duration_ms),
    }
}

旧的 ExecCommandBegin/End 仍可由兼容消费者发出,但 V2 通过 builder 把它们投影为同一 ThreadItem::CommandExecution 形状。开始项没有 exit code 和 aggregated output,结束项才补齐终态字段;这就是“事件名称变化”之外的语义转换。

源码位置:codex-rs/app-server-protocol/src/protocol/common.rs :: ClientRequestSerializationScope

rust
pub enum ClientRequestSerializationScope {
    Global(&'static str),
    GlobalSharedRead(&'static str),
    Thread { thread_id: String },
    ThreadPath { path: PathBuf },
    CommandExecProcess { process_id: String },
    Process { process_handle: String },
    FuzzyFileSearchSession { session_id: String },
    FsWatch { watch_id: String },
    McpOauth { server_name: String },
}

请求映射还决定并发串行化范围:同一 thread、process handle、command process、FS watch 或 MCP OAuth server 的请求会进入不同 scope。scope 不会出现在 JSON wire payload 中,却影响 App Server 是否并行处理两个请求,因此不能只检查 method 名称。

7. 信息损失 ​

映射可能改变名称、聚合粒度、字段可见性和错误形状。ContextCompacted 在 V2 被 canonical item 取代;stream error 会带 will_retry,普通 error 则可能更新 thread watch 状态。不能根据一个 V2 notification 反推出完整原始 EventMsg。

8. 源码验证 ​

event mapping 测试验证 Core ItemStarted/Completed 到 V2 item notification 的字段和 ID;dynamic tool/permissions 测试验证反向 server request 的 callback 与业务 call ID 分离;deprecated event 测试验证旧事件不会重复发送 canonical V2 lifecycle。它们证明映射方向和兼容策略,不证明所有字段都能逆向恢复。

源码位置:

  • codex-rs/app-server/src/bespoke_event_handling.rs :: item mapping tests
  • codex-rs/app-server/src/bespoke_event_handling.rs :: DynamicToolCallRequest
  • codex-rs/app-server/src/bespoke_event_handling.rs :: RequestPermissions
text
cd codex-rs
cargo test -p codex-app-server event_mapping
cargo test -p codex-app-server dynamic_tool
cargo test -p codex-app-server permissions

9. 映射排查 ​

遇到 Core 有事件但 V2 没通知,先确认该事件是否被 canonical item 替代或仅供兼容消费者;遇到客户端响应无法唤醒 Core,分别核对 connection callback ID 与业务 call/item ID;遇到字段缺失,回到 mapper 判断这是设计中的信息损失还是错误路径。