Skip to content

BEM事件与增量状态

沿委派、消息增量和客户端时间线追踪 BEM 输出,解释通道前缀、定时刷新、文本截断、UTF-8 分片、转录续段与产物去重。

基于rust-v0.150.0
CodexRealtimeRust流式处理

BEM事件与增量状态 ​

实时模型把“检查项目的构建错误”交给 Codex 后,普通 Agent 先输出 [COMMENTARY]正在查看日志,再输出 [FINAL]错误来自缺失的依赖。这两段文本怎样回到实时模型?如果标签被拆成 [COM 和 MENTARY],第一段会不会被提前发送?客户端看到的语音转录、后台回答和可视化产物,又怎样成为一条有稳定 ID 的时间线?

本文围绕这些问题阅读事件与增量状态。先读 实时会话架构,可以理解普通 Turn 与实时会话怎样共存;流式推理与消息Delta 解释普通模型的文本增量怎样进入 Core。这里假定读者了解 Rust 的 Option、Arc、Mutex 与 async/await,重点展开前缀识别、分片、完成与去重,不重复连接建立、媒体采集和 sideband 重连。

阅读时要同时关注两个消费者:实时模型接收用于后续生成的文本,App Server 则把语音和普通 Agent 的工作整理为客户端事件。它们对同一段文字保留的状态不同。读完后,应能从一个服务端委派走到实际出站 JSON,也能解释为何“转录最终值变了”和“时间线已保存的文字变了”不是同一件事。下文代码中的中文注释用于说明关键条件,原实现的英文注释保留。

1. 三套事件 ​

BEM 在这里指承担后台工作的 Codex 模型及其文本信封。实现中没有统一的 BemEvent 枚举:BEM 回答沿普通 Agent 的 item/delta 事件到达,通道信息可以编码在回答开头;实时服务另有入站事件,时间线又有独立的数据结构。

对象生产者与消费者表达的事实
RealtimeEventcodex-api 解析服务端帧,Core 消费转录、音频、委派、会话或响应事件
EventMsg::ItemStarted / AgentMessageContentDelta / ItemCompleted普通 Turn 产生,Session 镜像给实时子系统一个 Agent 输出项的开始、文本追加和完成
RealtimeOutboundCore 镜像与刷新逻辑产生,输入任务消费准备送给实时模型的追加、完整结果或独立文本
RealtimeItemApp Server 的 RealtimeHistoryState 产生可以存入时间线的语音片段、会话边界或普通 item 的展示引用

handoff 是实时服务把工作交给普通 Turn 的一次交接;V3 在线上称为 delegation。phase 是 Core 对普通回答阶段的分类;channel 是 V3 上下文追加请求的路由字段。后面会看到,phase 与 channel 有显式转换关系。

下图沿“谁生产、谁消费”区分这三套事件。委派沿左侧进入普通工作,回答回到实时模型;右侧观察事件并建立客户端时间线。

  • RealtimeOutbound 回到模型连接,ServerNotification 发往订阅 Thread 的应用客户端。
  • RealtimeItem 的内容可以只是“展示某个已有 item”,无需复制该 item 的全部正文。
  • 持久化时间线有单独的启用条件,图中的历史支路并非所有 Thread 都会执行。

先看入站模型的完整定义与 payload。RealtimeEvent 聚合了三种适配器可产生的事件,因此枚举里有一个变体,不代表每种线协议都支持它。

源码文件:codex-rs/protocol/src/protocol.rs

相关函数/类型:RealtimeEvent、RealtimeAudioFrame、RealtimeHandoffRequested(L360–L444,摘录)

rust
// 这是适配后的内部事件;可选 ID 和只剩文本的 payload 决定下游能否关联原事件。
#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeAudioFrame {
    pub data: String,
    pub sample_rate: u32,
    pub num_channels: u16,
    #[serde(skip_serializing_if = "Option::is_none")]
    pub samples_per_channel: Option<u32>,
    #[serde(skip_serializing_if = "Option::is_none")]
    pub item_id: Option<String>,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
// 转录增量不保留上游 item ID 或时间偏移。
pub struct RealtimeTranscriptDelta {
    pub delta: String,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeTranscriptDone {
    pub text: String,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeTranscriptEntry {
    pub role: String,
    pub text: String,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeHandoffRequested {
    pub handoff_id: String,
    pub item_id: String,
    pub input_transcript: String,
    pub active_transcript: Vec<RealtimeTranscriptEntry>,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeNoopRequested {
    pub call_id: String,
    pub item_id: String,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeInputAudioSpeechStarted {
    pub item_id: Option<String>,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeResponseCancelled {
    pub response_id: Option<String>,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeResponseCreated {
    pub response_id: Option<String>,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
pub struct RealtimeResponseDone {
    pub response_id: Option<String>,
}

#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
// 变体覆盖多种适配器,不表示每个版本都能产生所有变体。
pub enum RealtimeEvent {
    SessionUpdated {
        realtime_session_id: String,
        instructions: Option<String>,
    },
    InputAudioSpeechStarted(RealtimeInputAudioSpeechStarted),
    InputTranscriptDelta(RealtimeTranscriptDelta),
    InputTranscriptDone(RealtimeTranscriptDone),
    OutputTranscriptDelta(RealtimeTranscriptDelta),
    OutputTranscriptDone(RealtimeTranscriptDone),
    AudioOut(RealtimeAudioFrame),
    ResponseCreated(RealtimeResponseCreated),
    ResponseCancelled(RealtimeResponseCancelled),
    ResponseDone(RealtimeResponseDone),
    ConversationItemAdded(Value),
    ConversationItemDone {
        item_id: String,
    },
    HandoffRequested(RealtimeHandoffRequested),
    NoopRequested(RealtimeNoopRequested),
    Error(String),
}

这些字段已经体现了几种信息取舍。转录 delta 只有 delta 字符串,没有上游 item ID、时间戳或偏移;RealtimeAudioFrame.item_id 可为空;ConversationItemAdded 则保留一个任意 JSON 值。消费者不能假定所有事件都能按同一种 ID 关联。

变体组关键字段Core 与 App Server 的后续用途
SessionUpdated会话 ID、可选 instructions更新连接侧信息;不直接等于新增时间线边界
输入/输出 TranscriptDelta、TranscriptDone增量 delta 或终值 text分角色维护委派上下文;另一路产生转录通知与时间线片段
AudioOut编码数据、采样信息、可选 item ID向客户端转发音频;适用时用于打断进度
InputAudioSpeechStarted可选 item ID开始新的用户转录;V2 还参与音频截断
ResponseCreated / Cancelled / Done可选 response IDV2 响应状态协调;取消映射为客户端 itemAdded
ConversationItemAdded / Done原始 item 或 item IDAdded 可通知客户端;Done 不结束普通 Turn 的交接
HandoffRequestedhandoff ID、item ID、输入和转录启动或 steer 普通 Turn
NoopRequestedcall ID、item IDV2 为静默工具调用回填空结果
Error字符串转发错误并结束当前实时输入循环

2. 委派事件入口 ​

2.1 V3 入站映射 ​

公开入口是实验 API thread/realtime/start。客户端需要在初始化时接受实验 API,Thread 还需要启用默认关闭的 realtime_conversation feature。BEM 通道路由使用 version: "v3",它选择 RealtimeEventParser::FramelessBidi。这些入口检查和连接过程见 实时会话架构。

这里直接从 V3 连接收到一条文本帧开始。适配器先解析 JSON 的 type,再决定哪些字段进入内部事件。

源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi.rs

相关函数/类型:parse_frameless_bidi_event(L15–L37,摘录)

rust
// 适配器先按 type 分类;返回 None 的文本帧由外层 next_event 跳过。
pub(super) fn parse_frameless_bidi_event(payload: &str) -> Option<RealtimeEvent> {
    let (parsed, message_type) = parse_realtime_payload(payload, "frameless bidi")?;
    match message_type.as_str() {
        "session.started" | "session.updated" => parse_session_updated_event(&parsed),
        "output_audio.delta" => parse_output_audio_delta(&parsed),
        "input_transcript.added" => {
            parse_transcript_item(&parsed).map(RealtimeEvent::InputTranscriptDelta)
        }
        "output_transcript.added" => {
            parse_transcript_item(&parsed).map(RealtimeEvent::OutputTranscriptDelta)
        }
        "turn.done" => parse_turn_done(&parsed),
        "delegation.created" => parse_delegation_created(&parsed),
        "error" => parse_error_event(&parsed),
        _ => {
            debug!(
                "received unsupported frameless bidi event type: {message_type}, data: {payload}"
            );
            None
        }
    }
}

这段 match 覆盖 V3 当前支持的所有事件名。session.started 与 session.updated 归为同一个内部变体;turn.done 被解释为某个说话角色的转录完成,不能拿它结束 Codex 的普通 Turn。

V3 的 type必要字段内部结果与舍弃的信息
session.started / session.updatedsession.id 字符串SessionUpdated;instructions 可选
input_transcript.addeditem.text 字符串InputTranscriptDelta;不保留 item ID
output_transcript.addeditem.text 字符串OutputTranscriptDelta;不保留 item ID
turn.doneturn.role、turn.transcriptuser/assistant 分别对应两类 Done;其他 role 不产生事件
output_audio.deltaaudio 字符串AudioOut;适配器填 24,000 Hz、单声道,不保留 start/end 时间
delegation.createdclient delegation 的 id 与 content 数组HandoffRequested;只拼接 input_text 内容
error可提取的错误信息Error;缺少信息时仍可能被忽略
其他类型—返回 None,不生成一个“未知事件”供上层处理

与 BEM 交接直接相关的是 delegation.created。它不是任意 function call 的别名,而是经过目标、类型和字段检查的事件。

源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/protocol_frameless_bidi.rs

相关函数/类型:parse_delegation_created(L73–L96,摘录)

rust
// 只接收 client delegation;多个 input_text 直接拼接,不插入分隔符。
fn parse_delegation_created(parsed: &Value) -> Option<RealtimeEvent> {
    let item = parsed.get("item")?.as_object()?;
    if item.get("type").and_then(Value::as_str) != Some("delegation")
        || item.get("target").and_then(Value::as_str) != Some("client")
    {
        return None;
    }
    let item_id = item.get("id").and_then(Value::as_str)?.to_string();
    let input_transcript = item
        .get("content")
        .and_then(Value::as_array)?
        .iter()
        .filter(|content| content.get("type").and_then(Value::as_str) == Some("input_text"))
        .filter_map(|content| content.get("text").and_then(Value::as_str))
        .collect::<String>();

    Some(RealtimeEvent::HandoffRequested(RealtimeHandoffRequested {
        handoff_id: item_id.clone(),
        item_id,
        input_transcript,
        // 近期转录稍后由 update_active_transcript 附加。
        active_transcript: Vec::new(),
    }))
}

target 必须等于 client。content 中非 input_text 项和没有字符串 text 的项会被过滤,剩余文本直接连接,不会自动插入换行。例如两个输入项分别为 "检查" 和 "构建",结果是 "检查构建"。空数组可以得到空输入,是否还能提交普通 Turn,要继续看后续积累的转录。

归一化时 handoff_id 和 item_id 同时取 delegation 的 id;active_transcript 暂时为空。这只是解析器的结果,next_event 返回前还有一次有状态处理。

2.2 转录与提交 ​

RealtimeWebsocketEvents::next_event 调用 update_active_transcript。API 层的 RealtimeTranscriptState 用一个共享锁保护近期转录,并在委派到达时把这段增量上下文移入事件。

源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs

相关函数/类型:RealtimeTranscriptState、ActiveTranscriptState、update_active_transcript(L233–L244,L639–L644,摘录)

rust
// API 层缓存用于后续委派,与 App Server 的稳定时间线片段分开。
#[derive(Clone, Default)]
pub struct RealtimeTranscriptState {
    active: Arc<Mutex<ActiveTranscriptState>>,
}

#[derive(Default)]
struct ActiveTranscriptState {
    entries: Vec<RealtimeTranscriptEntry>,
    new_input_entry: bool,
    new_output_entry: bool,
}

// ...

RealtimeEvent::HandoffRequested(handoff) => {
    append_handoff_input(&mut active_transcript.entries, &handoff.input_transcript);
    // 交出本批转录并清空缓存,下一次委派从新的增量开始。
    handoff.active_transcript = std::mem::take(&mut active_transcript.entries);
    active_transcript.new_input_entry = true;
    active_transcript.new_output_entry = true;
}

mem::take 同时完成交付与清空,后续委派不会再次得到这一批缓存。两个 new_*_entry 标志阻止新一轮 delta 接到前一轮残留角色上。这里的转录服务于后续模型上下文;第 8 节的时间线状态有另外的 ID 与完成规则。

Core 在向 fanout 转发委派前,先更新自己的交接状态:

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:handle_realtime_server_event(L2369–L2373,摘录)

rust
// 这里只摘录 HandoffRequested 内部 match 的 V1 kind 分支;V3 也映射到该 kind。
RealtimeSessionKind::V1 => {
    let mut stream = handoff_state.stream.lock().await;
    stream.items.clear();
    stream.active_handoff = Some(handoff.handoff_id.clone());
}

V3 映射到内部的 RealtimeSessionKind::V1,因此执行上面的 V1 分支:新的委派替换 active_handoff,同时清掉旧的流式 item 状态。这个字段只保存一个当前 ID。V2 的活动交接处理走另一条 steering acknowledgement 路径,不能从名称中的 V1/V2 推断是否支持 BEM。

fanout 对 HandoffRequested 提取并包装输入,调用 Session::route_realtime_text_input;后者走 TurnInputMode::StartOrSteer。如果已有普通任务,输入用于 steer;否则开始新 Turn。没有直接调用某个工具处理器。完成路由后,fanout 再把原事件封装为 EventMsg::RealtimeConversationRealtime 对外发送。

下图把入站与回传连起来;图中“接受输入”不意味着普通 Turn 已经执行完。

  • API 的转录缓存先附加到委派,Core 随后才把它变成普通输入。
  • 输入任务通过事件队列交给独立的 fanout task,后者负责路由普通 Turn。
  • 图中输出按一种正常顺序展开;任务启动后,两路事件生产者可以交错。普通输出先进入 Session 的常规事件路径,再执行实时镜像,不同消费者的完成时刻没有一个总屏障。
  • ItemCompleted 完成一个输出项,TurnComplete 结束普通回合,职责不同。

畸形帧也需要按层判断。parse_realtime_payload 对无效 JSON 或缺少字符串 type 返回 None;next_event 的循环继续读取下一帧。与此不同,一个成功解析的 RealtimeEvent::Error 会先被送到输出队列,再使输入循环退出。不能把“忽略一帧”和“连接失败”统称为解析异常。

3. BEM 路由契约 ​

3.1 参数生效范围 ​

codexResponseHandoffMode 是启动实时会话时的参数,枚举本身声明了默认值与线上拼写:

源码文件:codex-rs/protocol/src/protocol.rs

相关函数/类型:CodexResponseHandoffMode(L1647–L1656,摘录)

rust
// 启动参数省略 mode 时使用 Thinking;线上写法由 serde 和 TS 一起确定。
#[derive(Debug, Clone, Copy, Default, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
#[serde(rename_all = "camelCase")]
#[ts(rename_all = "camelCase")]
pub enum CodexResponseHandoffMode {
    #[default]
    Thinking,
    Commentary,
    BemTags,
}

这里的 Thinking 是“出站消息省略 channel 字段”的本地选择。服务端如何利用这种上下文,需要由服务端契约决定,不能从枚举名推断一定静默或一定播报。BemTags 才启用文本前缀识别。

App Server 将可选字段转成 Core 的 ConversationStartParams,随后 start_inner 把它们保存到新的 RealtimeHandoffState。

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

相关函数/类型:thread_realtime_start_inner(L1145–L1154,摘录)

rust
// 这里只摘录 ConversationStartParams 构造中与回传相关的字段;其余启动字段省略。
client_managed_handoffs: params.client_managed_handoffs.unwrap_or(false),
delegation_ack_filler: params.delegation_ack_filler,
flush_transcript_tail_on_session_end: params
    .flush_transcript_tail_on_session_end
    .unwrap_or(false),
codex_responses_as_items: params.codex_responses_as_items.unwrap_or(false),
codex_response_item_prefix: params.codex_response_item_prefix,
codex_response_handoff_mode: params.codex_response_handoff_mode.unwrap_or_default(),
codex_response_handoff_channel_prefixes: params
    .codex_response_handoff_channel_prefixes,

因此这些设置在本次实时会话构造时固定,不是每个普通 Turn 从全局 Config 重新读取。重新启动会话会构造新状态;本路径没有前缀表的热更新操作。

参数或条件默认/取值对本篇路径的影响
version可用 v1/v2/v3;普通配置默认 v2只有 v3 选择 Frameless Bidi 的通道路由
codexResponseHandoffModethinking / commentary / bemTags省略 channel、固定 commentary、按标签分类
codexResponseHandoffChannelPrefixes可选 map仅覆盖配置了的 analysis/commentary/final 通道
clientManagedHandoffsfalsetrue 时跳过自动回答回传;不阻止入站委派提交普通 Turn
codexResponsesAsItemsfalsetrue 时不走流式 handoff append,完成文本走 conversation item 适配
codexResponseItemPrefix无只用于上一行的路径;BEM 判定发生在添加这个前缀之前
outputModality请求显式提供本篇 V3 使用 audio;text 模式要求 V2

例如,已启用实时 feature 的 Thread 可提交以下请求。threadId 应使用 thread/start 返回的 ID;若还需要后面的持久化时间线,创建 Thread 时选择 historyMode: "paginated"。

示例代码:

json
{
  "id": 2,
  "method": "thread/realtime/start",
  "params": {
    "threadId": "00000000-0000-4000-8000-000000000001",
    "version": "v3",
    "outputModality": "audio",
    "codexResponseHandoffMode": "bemTags",
    "codexResponseHandoffChannelPrefixes": {
      "commentary": ["[PROGRESS]", "[UPDATE]"],
      "final": ["[DONE]"]
    }
  }
}

这个请求让 commentary 接受两种新前缀,final 接受 [DONE],analysis 继续使用默认 [ANALYSIS]。它不会把用户看到的原始 Agent 文本改写成另一种协议,也不会改变普通模型的生成参数。

3.2 通道选择 ​

最终 channel 的选择集中在 v3_output_writer。它克隆 writer,只在这个克隆上设置 channel,底层 socket 仍然共享。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:v3_output_writer(L2066–L2085,摘录)

rust
// 在 writer 克隆上选择本次输出的 channel,不修改共享连接的全局路由。
fn v3_output_writer(
    writer: &RealtimeWebsocketWriter,
    phase: Option<&MessagePhase>,
    handoff_mode: CodexResponseHandoffMode,
) -> RealtimeWebsocketWriter {
    let channel = match handoff_mode {
        CodexResponseHandoffMode::Thinking => None,
        CodexResponseHandoffMode::Commentary => Some(RealtimeContextAppendChannel::Commentary),
        // BEM 的 analysis 与 commentary 都归为 MessagePhase::Commentary。
        CodexResponseHandoffMode::BemTags => match phase {
            Some(MessagePhase::FinalAnswer) => Some(RealtimeContextAppendChannel::Speakable),
            Some(MessagePhase::Commentary) => Some(RealtimeContextAppendChannel::Commentary),
            None => Some(RealtimeContextAppendChannel::Speakable),
        },
    };
    match channel {
        Some(channel) => writer.clone().with_context_append_channel(channel),
        None => writer.clone(),
    }
}

BemTags 把 Commentary 映射为 commentary,把 FinalAnswer 以及缺省 phase 映射为 speakable。文本本身保留原标签,所以虽然 [ANALYSIS] 和 [COMMENTARY] 使用同一 API channel,实时模型仍能看到两者的区别。

mode[ANALYSIS]…[COMMENTARY]…[FINAL]…无法识别的完整文本
thinking无 channel无 channel无 channel无 channel
commentarycommentarycommentarycommentarycommentary
bemTagscommentarycommentaryspeakablespeakable

thread/realtime/appendSpeech 则产生 StandaloneSpeech,V3 明确选择 speakable。它不沿用自动回答的 mode。V1/V2 不使用这套选择;“给 V1 请求加上 bemTags”不会把它升级为 V3。

4. 标签的增量识别 ​

4.1 前缀判定 ​

message_phase 是 BEM 分类的全部规则。这里没有正则、大小写归一化或通用标签解析器。

源码文件:codex-rs/core/src/realtime_conversation/bem.rs

相关函数/类型:message_phase(L5–L28,摘录)

rust
// 依次检查三个通道;配置存在时替换该通道默认前缀。
pub(super) fn message_phase(
    text: &str,
    channel_prefixes: &BTreeMap<String, Vec<String>>,
) -> Option<MessagePhase> {
    // 固定顺序决定重叠前缀的优先级,不按最长前缀排序。
    for (channel, default_prefix, phase) in [
        ("analysis", "[ANALYSIS]", MessagePhase::Commentary),
        ("commentary", "[COMMENTARY]", MessagePhase::Commentary),
        ("final", "[FINAL]", MessagePhase::FinalAnswer),
    ] {
        let matches = channel_prefixes.get(channel).map_or_else(
            || text.starts_with(default_prefix),
            |prefixes| {
                prefixes
                    .iter()
                    // 空串会匹配任意文本,必须显式排除。
                    .any(|prefix| !prefix.is_empty() && text.starts_with(prefix))
            },
        );
        if matches {
            return Some(phase);
        }
    }

    None
}

BTreeMap 负责按名称查找配置,真正的匹配优先级由数组顺序决定:analysis、commentary、final。存在某个 key 时,只尝试该 key 的自定义列表;列表为空或仅含空串,也不会回退到默认前缀。过滤空串是必要的,否则 Rust 的 starts_with("") 会匹配所有文本。

下面是直接按判定式推演的边界,而非额外的协议承诺:

配置与输入判定原因
默认;[FINAL]完成FinalAnswer起始位置、大小写精确匹配
默认; [FINAL]完成None前导空格不被 trim
默认;[final]完成None不做大小写折叠
commentary 配 ["[UPDATE]"];[COMMENTARY]运行中None自定义列表替换这一通道的默认前缀
analysis 配 [];[ANALYSIS]文本None有 key 但无可匹配前缀
analysis 配 ["["];[FINAL]完成Commentaryanalysis 先匹配;不是最长前缀优先
默认;[FINAL][COMMENTARY]文本FinalAnswer只认起始前缀;后面的标签是文本

因此自定义前缀最好互不重叠。否则同一条回答落到哪个 channel 由匹配顺序决定,增加前缀的排列长度并不能得到“最具体匹配”。

4.2 半标签缓冲 ​

网络上的模型文本可以在任何字符边界分块。ChannelParser 为一个输出项保留尚未分类的文本,并且一旦识别,就锁定该项的 phase。

源码文件:codex-rs/core/src/realtime_conversation/bem.rs

相关函数/类型:ChannelParser(L30–L67,摘录)

rust
// 每个输出项独占一份 parser;只读前缀表用 Arc 共享。
/// Buffers a streamed BEM message until its channel header is complete.
///
/// Once the channel is known, the original envelope is released unchanged so
/// the frontend model can distinguish BEM `analysis` from `commentary`.
#[derive(Debug, Default)]
pub(super) struct ChannelParser {
    channel_prefixes: Arc<BTreeMap<String, Vec<String>>>,
    buffered_text: String,
    phase: Option<MessagePhase>,
}

impl ChannelParser {
    pub(super) fn new(channel_prefixes: Arc<BTreeMap<String, Vec<String>>>) -> Self {
        Self {
            channel_prefixes,
            ..Self::default()
        }
    }

    pub(super) fn push(&mut self, text: &str) -> Option<String> {
        // 一旦分类就不再解释后续正文中的标签。
        if self.phase.is_some() {
            return Some(text.to_string());
        }

        self.buffered_text.push_str(text);
        self.phase = message_phase(&self.buffered_text, &self.channel_prefixes);
        // 尚未识别时保留缓冲;命中后交出包含前缀的原始文本。
        self.phase.as_ref()?;
        Some(std::mem::take(&mut self.buffered_text))
    }

    pub(super) fn phase(&self) -> Option<MessagePhase> {
        self.phase.clone()
    }

    pub(super) fn finish(&mut self) -> String {
        std::mem::take(&mut self.buffered_text)
    }
}

三个字段各有职责:channel_prefixes 是共享的只读配置;buffered_text 是尚未释放的原始信封;phase 是识别结果。phase.as_ref()? 使未识别时立即返回 None,但文本已经保留。识别成功后 mem::take 交出整段缓冲,包括原前缀;后续 push 直接返回新 chunk,不再重新分类。

调用返回值保留的状态
push("[COM")Nonebuffer 为 [COM,phase 为空
push("MENTARY]查日志")[COMMENTARY]查日志buffer 已取走,phase 为 Commentary
push(",继续定位"),继续定位phase 不变
finish()空串仅排空 buffer;phase 仍为 Commentary

finish 不重置 phase。如果把这个 parser 复用给下一条 item,下一条 final 会继承前一条 commentary 的判定。实际调用方为每个流式 item 新建 parser,并在 item 完成时移除整项。

下图用字段条件表示生命周期。Removed 表示外层 owner 移除了 item,并不是 ChannelParser 自带的枚举状态。

  • 未命中可能只是 chunk 太短,也可能是整个消息都不符合约定,parser 暂时不区分。
  • 第一次命中是单向转换;消息正文里的新标签不会改变 phase。
  • 完成后的降级由 RealtimeStreamedItem::finish_input 负责,parser 只返回剩余文本。

4.3 无法识别的输出 ​

外层 item 在 push 时调用 parser,在结束时把未识别的非空文本降为 FinalAnswer。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:RealtimeStreamedItem::push_text、RealtimeStreamedItem::finish_input(L218–L247,摘录)

rust
// push 只接受已分类文本;item 完成时才把非空的未知信封降级为 FinalAnswer。
fn push_text(&mut self, text: &str) {
    if text.is_empty() {
        return;
    }

    let text = if let Some(parser) = self.bem_channel_parser.as_mut() {
        let Some(text) = parser.push(text) else {
            return;
        };
        self.phase = parser.phase();
        text
    } else {
        text.to_string()
    };
    self.push_output_text(&text);
}

fn finish_input(&mut self) {
    let Some(parser) = self.bem_channel_parser.as_mut() else {
        return;
    };
    self.phase = parser.phase();
    // finish 只交出剩余缓冲;降级状态由外层 item 决定。
    let text = parser.finish();
    if self.phase.is_none() && !text.is_empty() {
        warn!("BEM output ended before a recognized channel header was received");
        self.phase = Some(MessagePhase::FinalAnswer);
    }
    self.push_output_text(&text);
}

这解释了一个很具体的现象:普通界面已经显示了一段回答,实时模型却暂时没有收到。若回答没有匹配前缀,它会一直留在 parser 中,直到 item 完成,再按 speakable 回传。修改 200 ms 刷新间隔不能解决这一问题,因为尚无已分类文本可刷。

还要区分“输出预算”与“识别缓冲”。这段 parser 对未识别文本只有 push_str,没有独立的长度上限;后面的截断发生在 push_output_text。因此 4,000 字节的下游输出预算,不能证明未知前缀的暂存内存也不超过 4,000 字节。对持续不结束的未知文本,只能从该实现得出“缓冲会继续增长”,不能宣称会按预算提前丢弃。

上游测试把拆分标签和不能识别的结束分别固定下来:

源码文件:codex-rs/core/src/realtime_conversation/bem_tests.rs

相关函数/类型:buffers_streamed_text_until_the_bem_channel_is_complete、preserves_unrecognized_output_when_the_stream_finishes(L67–L79,L96–L103,摘录)

rust
// 测试分别约束半标签释放时机和未知文本的 finish 行为。
#[test]
fn buffers_streamed_text_until_the_bem_channel_is_complete() {
    let mut parser = ChannelParser::default();

    assert_eq!(parser.push("[COM"), None);
    assert_eq!(
        parser.push("MENTARY]progress"),
        Some("[COMMENTARY]progress".to_string())
    );
    assert_eq!(parser.phase(), Some(MessagePhase::Commentary));
    assert_eq!(parser.push(" update"), Some(" update".to_string()));
}

// ...

#[test]
fn preserves_unrecognized_output_when_the_stream_finishes() {
    let mut parser = ChannelParser::default();

    assert_eq!(parser.push("plain output"), None);
    assert_eq!(parser.finish(), "plain output");
    assert_eq!(parser.phase(), None);
}

第一个断言要求半标签不产出;第二次 push 必须返回完整原信封,后续 delta 原样通过。另一个测试要求 finish 取回未知文本且 phase 仍为空,说明 FinalAnswer 的兜底不在 parser 内部。这些测试验证字符串状态转换,不涉及 WebSocket 和真实语音生成。

5. 输出项的所有权 ​

5.1 共享与独占 ​

当前交接状态保存在 RealtimeHandoffState。其中只有轻量句柄与配置会被克隆,活动 item 集合由共享锁保护。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:RealtimeHandoffState、RealtimeHandoffOutput、RealtimeHandoffStreamState(L145–L170,摘录)

rust
// Clone 复制发送端和 Arc 句柄;活动交接与 item map 仍由共享 Mutex 保护。
#[derive(Clone, Debug)]
struct RealtimeHandoffState {
    output_tx: Sender<RealtimeOutbound>,
    last_output: Arc<Mutex<Option<RealtimeHandoffOutput>>>,
    stream: Arc<Mutex<RealtimeHandoffStreamState>>,
    client_managed_handoffs: bool,
    codex_responses_as_items: bool,
    codex_response_item_prefix: Option<String>,
    codex_response_handoff_mode: CodexResponseHandoffMode,
    codex_response_handoff_channel_prefixes: Arc<BTreeMap<String, Vec<String>>>,
    session_kind: RealtimeSessionKind,
    event_parser: RealtimeEventParser,
}

#[derive(Clone, Debug)]
struct RealtimeHandoffOutput {
    text: String,
    phase: Option<MessagePhase>,
}

#[derive(Debug, Default)]
struct RealtimeHandoffStreamState {
    active_handoff: Option<String>,
    items: HashMap<String, RealtimeStreamedItem>,
}
字段来源、初值与责任
output_txstart_inner 创建的 64 容量队列;把回传请求交给输入任务
last_output初始 None;保存完整回传的末次文本与 phase,供适用的完成协议使用
stream初始无活动交接、空 items;多个事件与延迟刷新任务共享
client_managed_handoffs / codex_responses_as_items来自本次启动;分别控制自动回传与回传形态
codex_response_item_prefix启动参数;只用于 conversation item 路径
codex_response_handoff_mode启动时固定的三种路由之一
codex_response_handoff_channel_prefixes缺省为空 map,以 Arc 共享,不逐 item 复制整张表
session_kind / event_parser由实时版本派生;一个服务 Core 分支,一个选择线协议
stream.active_handoff当前交接 ID;由入站委派设置,由普通 Turn 完成清理
stream.items普通 Agent 的 item ID → 对应流式状态;一个交接可产生多个输出项

RealtimeStreamedItem 才拥有某个输出项的缓存、累计计数和刷新标记。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:RealtimeStreamedItem(L171–L184,摘录)

rust
// 这是单个普通 Agent item 的增量状态;sent_bytes 记录本地排出量。
#[derive(Debug)]
struct RealtimeStreamedItem {
    handoff_id: String,
    phase: Option<MessagePhase>,
    // 只有标签路由才安装 parser,半标签不会跨 item 共享。
    bem_channel_parser: Option<BemChannelParser>,
    prefix_final_message: bool,
    sent_bytes: usize,
    buffered_text: String,
    tail_text: String,
    truncated: bool,
    last_flush_at: Instant,
    // 同一 item 在等待期间只安排一份刷新任务。
    flush_scheduled: bool,
}
字段为什么需要它
handoff_id注册时绑定交接;后续发送使用这一 ID
phase非 BEM 模式沿用普通 item;BEM 模式等待解析
bem_channel_parser仅 BEM 标签路由时存在,独占本 item 的半标签
prefix_final_message控制旧式标签;当前 V3 注册路径为 false
sent_bytes已从缓冲排出的字节数,包含额外前缀;不是远端确认计数
buffered_text识别后等待发送的文本
tail_text超出预算后持续保留的最新尾部
truncated是否已经改为“开头 + 截断标记 + 尾部”输出
last_flush_at上次定时刷新取出文本的本地时刻
flush_scheduled同一 item 是否已有等待中的刷新任务,防止每个 delta 都 spawn

图中实线表示拥有,虚线表示共享或发送。长寿命对象是 handoff 状态,parser 跟着单个 item 存亡。

  • 每个 item 独立识别标签;不能用一份 parser 顺序处理整个 Turn。
  • 延迟任务克隆的是 handoff 句柄,再按 item ID 回到共享 map 查找。
  • writer 的 channel 是一次输出选择;共享连接不意味着全局 phase 会随之改变。

5.2 注册与镜像 ​

Session::send_event 在常规事件投递后调用 maybe_mirror_event_text_to_realtime。镜像函数明确区分 started、delta、completed 三种事件。

源码文件:codex-rs/core/src/session/mod.rs

相关函数/类型:maybe_mirror_event_text_to_realtime(L2142–L2185,摘录)

rust
// 常规 Session 事件已投递后才镜像;三种 item 事件承担不同阶段。
async fn maybe_mirror_event_text_to_realtime(&self, msg: &EventMsg) {
    if self.conversation.running_state().await.is_none() {
        return;
    }
    match msg {
        EventMsg::ItemStarted(event) => {
            if let TurnItem::AgentMessage(item) = &event.item {
                self.conversation
                    .register_handoff_stream_item(
                        item.id.clone(),
                        item.phase.clone(),
                        agent_message_text(item),
                    )
                    .await;
            }
            return;
        }
        EventMsg::AgentMessageContentDelta(event) => {
            if let Err(err) = self
                .conversation
                .stream_handoff_delta(&event.item_id, event.delta.clone())
                .await
            {
                debug!("failed to stream event text to realtime conversation: {err}");
            }
            return;
        }
        EventMsg::ItemCompleted(event) => {
            if let TurnItem::AgentMessage(item) = &event.item
                // 已按流式路径处理时,避免在 completed 后重复发送完整消息。
                && self.conversation.finish_handoff_stream_item(&item.id).await
            {
                return;
            }
        }
        _ => {}
    }
    let Some((text, phase)) = realtime_text_for_event(msg) else {
        return;
    };
    if let Err(err) = self.conversation.handoff_out(text, phase).await {
        debug!("failed to mirror event text to realtime conversation: {err}");
    }
}

started 只注册;delta 只追加;completed 先尝试排空已注册的流式项。如果这一项已通过流式路径排出文本,直接 return,避免再把完整回答发送一次。未注册、没有排出文本或不满足流式条件时,才走 realtime_text_for_event → handoff_out 的完整输出路径。

自动流式回传的开关比“实时会话正在运行”更严格:

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:streams_handoff_append、routes_handoff_by_bem(L444–L456,摘录)

rust
// 是否流式交付与是否按标签分类是两个独立判断。
impl RealtimeHandoffState {
    fn streams_handoff_append(&self) -> bool {
        self.event_parser == RealtimeEventParser::FramelessBidi
            && !self.client_managed_handoffs
            && !self.codex_responses_as_items
    }

    fn routes_handoff_by_bem(&self) -> bool {
        self.event_parser == RealtimeEventParser::FramelessBidi
            && self.codex_response_handoff_mode == CodexResponseHandoffMode::BemTags
    }
}

这两个函数必须分别理解。streams_handoff_append 决定是否进行流式回传;routes_handoff_by_bem 决定是否读取标签。V3 的 thinking 模式也可以流式回传,只是不会安装 BEM parser。

注册时还要求已有 active_handoff,并把当前 handoff ID 与普通 item ID 关联起来:

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:register_handoff_stream_item(L877–L932,摘录)

rust
// 只有活动 handoff 才注册;BEM 模式从初始文本重新判断 phase。
pub(crate) async fn register_handoff_stream_item(
    &self,
    item_id: String,
    phase: Option<MessagePhase>,
    initial_text: String,
) {
    let handoff = {
        let guard = self.state.lock().await;
        guard.as_ref().map(|state| state.handoff.clone())
    };
    let Some(handoff) = handoff else {
        return;
    };
    if !handoff.streams_handoff_append() {
        return;
    }
    let flush_delay = {
        let mut stream = handoff.stream.lock().await;
        let Some(handoff_id) = stream.active_handoff.clone() else {
            return;
        };
        let mut streamed_item = RealtimeStreamedItem {
            handoff_id,
            // 不直接相信普通 Agent item 的 phase。
            phase: if handoff.routes_handoff_by_bem() {
                None
            } else {
                phase
            },
            bem_channel_parser: handoff.routes_handoff_by_bem().then(|| {
                BemChannelParser::new(Arc::clone(
                    &handoff.codex_response_handoff_channel_prefixes,
                ))
            }),
            prefix_final_message: handoff.event_parser == RealtimeEventParser::V1,
            sent_bytes: 0,
            buffered_text: String::new(),
            tail_text: String::new(),
            truncated: false,
            last_flush_at: Instant::now(),
            flush_scheduled: false,
        };
        streamed_item.push_text(&initial_text);
        let flush_delay = if streamed_item.buffered_text.is_empty() {
            None
        } else {
            streamed_item.flush_scheduled = true;
            Some(streamed_item.next_flush_delay())
        };
        // map 的 key 是普通 item ID;value 另存它绑定的 handoff ID。
        stream.items.insert(item_id.clone(), streamed_item);
        flush_delay
    };
    if let Some(flush_delay) = flush_delay {
        schedule_streamed_handoff_flush(&handoff, item_id, flush_delay);
    }
}

BEM 模式有意忽略传入的普通 item phase,先设为 None,再按原文本判断。初始文本也会经过同一个 parser,所以标签既可能在 started 里完整出现,也可能跨 started 和后续 delta。尚无活动交接时不注册,完成事件仍有完整输出的后备路径。

5.3 两种完成 ​

finish_handoff_stream_item 先从 map 中取走 item,然后处理半标签残留并排出最后一段。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:finish_handoff_stream_item(L972–L1001,摘录)

rust
// 先移除 item 再排空,之后的延迟任务无法再次取出这份缓存。
pub(crate) async fn finish_handoff_stream_item(&self, item_id: &str) -> bool {
    let handoff = {
        let guard = self.state.lock().await;
        guard.as_ref().map(|state| state.handoff.clone())
    };
    let Some(handoff) = handoff else {
        return false;
    };
    if !handoff.streams_handoff_append() {
        return false;
    }
    let Some(mut streamed_item) = handoff.stream.lock().await.items.remove(item_id) else {
        return false;
    };
    streamed_item.finish_input();
    let chunk = streamed_item.drain_final_chunk();
    // drain 已更新 sent_bytes;这个布尔值不代表远端已确认接收。
    let sent_output = streamed_item.sent_bytes > 0;
    if let Some(text) = chunk {
        // 发送错误被忽略,不能据返回值建立交付保证。
        let _ = handoff
            .output_tx
            .send(RealtimeOutbound::HandoffAppend {
                handoff_id: streamed_item.handoff_id,
                text,
                phase: streamed_item.phase,
            })
            .await;
    }
    sent_output
}

取走 item 以后,尚未醒来的刷新任务再按 ID 查找就会返回,不会再刷新这份缓存。但 sent_output 的含义需要看计算顺序:drain_final_chunk 已增加 sent_bytes,队列 send 还没有完成;send 错误又被忽略。因此返回 true 表示本地已把这一项作为流式输出处理,不能解释为“实时服务已经接收”。

普通 Turn 完成才清理活动交接和剩余 items:

源码文件:codex-rs/core/src/session/mod.rs

相关函数/类型:maybe_clear_realtime_handoff_for_event(L2186–L2195,摘录)

rust
// 只有普通 TurnComplete 触发交接收尾,单个 item 完成不会走这里。
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;
}

V3 的 handoff_complete 因内部 session kind 为 V1 而直接返回,随后仍执行清理;它不会额外发送 V2 的 function-call acknowledgement。RealtimeEvent::ConversationItemDone 也不进入这条清理函数,所以一条普通回答完成后,同一 Turn 仍可以继续产生另一条 commentary 或 final。

当新委派替换旧委派,或者 Turn 完成清空 items,尚在 map 内的旧缓存会被丢弃。已经排出锁并放入队列的输出不属于这张 map;V3 发送分支也没有 V2 那种发送前再次比对 active handoff ID 的逻辑。因此“清空 items”不构成撤回所有已排队输出的保证。

6. 定时刷新与截断 ​

6.1 刷新调度 ​

新增 delta 先更新 item,再决定是否安排一次延迟刷新。它不是一个一直运行、每 200 ms 唤醒全部 item 的定时器。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:stream_handoff_delta(L933–L971,摘录)

rust
// 追加 delta 后复用等待中的刷新任务;没有注册的 item 不在这里补建。
pub(crate) async fn stream_handoff_delta(
    &self,
    item_id: &str,
    delta: String,
) -> CodexResult<()> {
    if delta.is_empty() {
        return Ok(());
    }
    let handoff = {
        let guard = self.state.lock().await;
        let Some(state) = guard.as_ref() else {
            return Err(CodexErr::InvalidRequest(
                "conversation is not running".to_string(),
            ));
        };
        state.handoff.clone()
    };
    if !handoff.streams_handoff_append() {
        return Ok(());
    }
    let flush_delay = {
        let mut stream = handoff.stream.lock().await;
        let Some(streamed_item) = stream.items.get_mut(item_id) else {
            return Ok(());
        };
        streamed_item.push_text(&delta);
        // 已安排刷新或头部预算耗尽时,只保留本次更新的状态。
        if streamed_item.flush_scheduled || streamed_item.streamable_text_bytes() == 0 {
            None
        } else {
            streamed_item.flush_scheduled = true;
            Some(streamed_item.next_flush_delay())
        }
    };
    if let Some(flush_delay) = flush_delay {
        schedule_streamed_handoff_flush(&handoff, item_id.to_string(), flush_delay);
    }
    Ok(())
}

flush_scheduled 为 true 时,新 delta 只合并进缓存,沿用已有任务。头部流式预算已经耗尽时,也不会安排新的定时发送,后续内容要等 item 结束。延迟由 HANDOFF_STREAM_FLUSH_INTERVAL.saturating_sub(last_flush_at.elapsed()) 计算,已过去足够时间就得到零延迟。

刷新任务醒来后再次查 map,在锁内提取要发送的字符串,离开锁后等待队列容量:

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:flush_streamed_handoff_item、schedule_streamed_handoff_flush(L2027–L2065,摘录)

rust
// 锁内取走文本,锁外等待队列;延迟任务重新按 item ID 查找。
async fn flush_streamed_handoff_item(handoff: &RealtimeHandoffState, item_id: &str) {
    let (handoff_id, text, phase) = {
        let mut stream = handoff.stream.lock().await;
        let Some(streamed_item) = stream.items.get_mut(item_id) else {
            return;
        };
        streamed_item.flush_scheduled = false;
        let Some(text) = streamed_item.drain_stream_chunk() else {
            return;
        };
        // 刷新时钟记录本地排空,不是 socket 写入完成。
        streamed_item.last_flush_at = Instant::now();
        (
            streamed_item.handoff_id.clone(),
            text,
            streamed_item.phase.clone(),
        )
    };
    // 这里也不向上返回队列 send 的错误。
    let _ = handoff
        .output_tx
        .send(RealtimeOutbound::HandoffAppend {
            handoff_id,
            text,
            phase,
        })
        .await;
}

fn schedule_streamed_handoff_flush(
    handoff: &RealtimeHandoffState,
    item_id: String,
    flush_delay: Duration,
) {
    let handoff = handoff.clone();
    let _flush_task = tokio::spawn(async move {
        tokio::time::sleep(flush_delay).await;
        flush_streamed_handoff_item(&handoff, &item_id).await;
    });
}

这里有三个时刻:等待任务到期、从缓冲取出、实际向 sideband 写入。last_flush_at 在第二个时刻更新,所以 200 ms 只是合并输出的本地调度间隔。执行器繁忙、64 容量 handoff 队列已满、socket 写入受阻,都能延后真实发送。锁没有跨越队列 send,可避免等待下游时堵住后续 item 状态访问。

6.2 头尾预算 ​

单条 BEM 回传不是无限追加。REALTIME_ASSISTANT_OUTPUT_TOKEN_BUDGET 为 1,000,approx_bytes_for_tokens 按每 token 约 4 字节估算,得到 4,000 字节预算;这不是 tokenizer 的精确 token 数。

预算分成可提前发送的头部、预留的截断标记和最后保留的尾部:

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:stream_head_byte_limit、tail_byte_limit、streamable_text_bytes(L201–L217,摘录)

rust
// 给尾部和截断标记留空间,不能把全部预算都用于提前发送。
fn stream_head_byte_limit(&self) -> usize {
    let output_byte_limit = approx_bytes_for_tokens(REALTIME_ASSISTANT_OUTPUT_TOKEN_BUDGET);
    output_byte_limit.saturating_sub(HANDOFF_STREAM_TRUNCATION_MARKER.len()) / 2
}

fn tail_byte_limit(&self) -> usize {
    approx_bytes_for_tokens(REALTIME_ASSISTANT_OUTPUT_TOKEN_BUDGET)
        .saturating_sub(self.stream_head_byte_limit())
        .saturating_sub(HANDOFF_STREAM_TRUNCATION_MARKER.len())
}

fn streamable_text_bytes(&self) -> usize {
    self.stream_head_byte_limit()
        .saturating_sub(self.sent_bytes)
        .saturating_sub(self.output_prefix().len())
}

当前截断标记 "\n…output truncated…\n" 占 24 个 UTF-8 字节,因此头、尾各预留 1,988 字节。streamable_text_bytes 还会扣除已排出的字节和适用的旧式前缀;V3 不添加 "Agent Final Message" 标签。

这是一种为未知终长准备的策略。若提前把全部 4,000 字节都发送,后续再发现输出很长,就无法撤回中段并给结论留下尾部空间。先只流出一半,才能在结束时决定保留完整剩余文本还是切换到头尾表示。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:RealtimeStreamedItem::push_output_text(L248–L277,摘录)

rust
// 第一次超过预算时决定头尾;之后只移动尾部窗口。
fn push_output_text(&mut self, text: &str) {
    if text.is_empty() {
        return;
    }
    // 进入截断模式后不再增长待发头部。
    if self.truncated {
        self.tail_text.push_str(text);
        self.tail_text =
            take_last_bytes_at_char_boundary(&self.tail_text, self.tail_byte_limit())
                .to_string();
        return;
    }

    self.buffered_text.push_str(text);
    let output_byte_limit = approx_bytes_for_tokens(REALTIME_ASSISTANT_OUTPUT_TOKEN_BUDGET);
    let remaining_text_bytes = output_byte_limit
        .saturating_sub(self.sent_bytes)
        .saturating_sub(self.output_prefix().len());
    if self.buffered_text.len() <= remaining_text_bytes {
        return;
    }

    // 在 UTF-8 字符边界截取,避免拆坏字符串。
    let head_bytes =
        take_bytes_at_char_boundary(&self.buffered_text, self.streamable_text_bytes()).len();
    self.tail_text =
        take_last_bytes_at_char_boundary(&self.buffered_text, self.tail_byte_limit())
            .to_string();
    self.buffered_text.truncate(head_bytes);
    self.truncated = true;
}

触发截断前,文本继续进入 buffered_text;第一次超过剩余预算时,从缓存截出尚可发送的头,并保留最新尾部。一旦 truncated 为 true,新文本只更新尾部,旧的中间部分不会恢复。

常规刷新和最终刷新也因此不同:

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:drain_stream_chunk、drain_final_chunk(L278–L313,摘录)

rust
// 常规刷新受头部额度约束;最终刷新才释放剩余全文或标记加尾部。
fn drain_stream_chunk(&mut self) -> Option<String> {
    let prefix = self.output_prefix();
    let available_text_bytes = self.streamable_text_bytes();
    if self.buffered_text.is_empty() || available_text_bytes == 0 {
        return None;
    }

    let requested_bytes = available_text_bytes.min(self.buffered_text.len());
    let split_at = take_bytes_at_char_boundary(&self.buffered_text, requested_bytes).len();
    if split_at == 0 {
        return None;
    }
    let text = self.buffered_text.drain(..split_at).collect::<String>();
    let text = format!("{prefix}{text}");
    // 计数在构造输出字符串时更新,尚未进行网络 I/O。
    self.sent_bytes += text.len();
    Some(text)
}

fn drain_final_chunk(&mut self) -> Option<String> {
    let prefix = self.output_prefix();
    if !self.truncated {
        if self.buffered_text.is_empty() {
            return None;
        }
        let text = self.buffered_text.drain(..).collect::<String>();
        let text = format!("{prefix}{text}");
        self.sent_bytes += text.len();
        return Some(text);
    }

    let head = self.buffered_text.drain(..).collect::<String>();
    let tail = self.tail_text.drain(..).collect::<String>();
    let text = format!("{prefix}{head}{HANDOFF_STREAM_TRUNCATION_MARKER}{tail}");
    self.sent_bytes += text.len();
    Some(text)
}

常规刷新受头部预算约束;最终刷新在未截断时排出全部剩余文本,在已截断时补上标记和尾部。截断位置都落在 UTF-8 字符边界,避免把一个中文字符或 emoji 切成非法字符串;这不保证 Markdown、JSON 或完整单词边界。

上游 streamed_handoff_preserves_a_bounded_final_tail 构造 HEAD + 5,000 个 x + TAIL,先取流式头,再取最终块,断言拼接结果不超过 4,000 字节、保留 HEAD、包含截断标记、以 TAIL 结束。它验证的是已进入输出缓冲的预算算法,不能覆盖第 4 节尚未识别的标签缓存。

6.3 跨层测试 ​

下面这组断言来自真正把普通 SSE 输出接到实时 WebSocket 的集成测试。测试先用 [PRO 和 GRESS]seed first 拆开自定义 commentary 标签,并用 oneshot gate 暂停 item 完成事件。

源码文件:codex-rs/core/tests/suite/realtime_conversation.rs

相关函数/类型:conversation_flushes_assistant_deltas_every_200ms_for_v3_handoff(L3698–L3711,L3840–L3875,L3886–L3892,摘录)

rust
// 省略 mock 初始化与事件等待;保留拆分输入、gate、出站 JSON 和防重复断言。
let initial_commentary_text = "[PRO";
let first_commentary_delta = "GRESS]seed first ";
let second_commentary_delta = "x".repeat(94);
let commentary_text =
    format!("{initial_commentary_text}{first_commentary_delta}{second_commentary_delta}");
let (gate_commentary_done_tx, gate_commentary_done_rx) = oneshot::channel();
let commentary_item_added =
    responses::ev_message_item_added("msg_commentary", initial_commentary_text);
let commentary_item_done = responses::ev_assistant_message("msg_commentary", &commentary_text);
let initial_final_text = "[DO";
let final_delta = "NE]done";
let final_text = format!("{initial_final_text}{final_delta}");
let final_item_added = responses::ev_message_item_added("msg_final", initial_final_text);
let final_item_done = responses::ev_assistant_message("msg_final", &final_text);

// ...

assert_eq!(
    wait_for_websocket_request(
        &realtime_server,
        /*connection_index*/ 0,
        /*request_index*/ 1,
    )
    .await?
    .body_json(),
    json!({
        "type": "delegation.context.append",
        "delegation_item_id": "delegation_stream",
        "channel": "commentary",
        "content": [{ "type": "input_text", "text": commentary_text }]
    })
);

// 此前已收到 commentary append,才允许普通 item 完成。
let _ = gate_commentary_done_tx.send(());
assert_eq!(
    wait_for_websocket_request(
        &realtime_server,
        /*connection_index*/ 0,
        /*request_index*/ 2,
    )
    .await?
    .body_json(),
    json!({
        "type": "delegation.context.append",
        "delegation_item_id": "delegation_stream",
        "channel": "speakable",
        "content": [{
            "type": "input_text",
            "text": final_text
        }]
    })
);

// ...

tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
    realtime_server.single_connection().len(),
    3,
    "completed assistant item must not resend the full text"
);

第一个 sideband 请求必须在 gate 放行前到达,说明输出没有等完整 item;其文本仍含 [PROGRESS],channel 为 commentary。放行后,[DO 加 NE]done 组成 final,channel 为 speakable。最后连接只有 session update 加两次 append,排除了 completed 时重发完整文本的正常路径。

测试名中的 200 ms 对应实现常量;断言重点是“完成前已收到合并文本”和“没有完整重发”,并没有把网络端到端时延测成严格的 200 ms。发送失败后的重复或丢失也不由这组 happy-path 断言保证。

7. 线协议分片 ​

7.1 内部输出类型 ​

Core 先构造 RealtimeOutbound,再由 handle_handoff_output 按适配器发送。保留这个中间层,可以让 Session 不直接构造 V1/V2/V3 的 JSON。

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:RealtimeOutbound(L324–L356,摘录)

rust
// 这是 Core 给输入任务的输出请求,不是直接序列化的线协议枚举。
#[derive(Clone, Debug, PartialEq, Eq)]
enum RealtimeOutbound {
    StandaloneHandoff {
        text: String,
        phase: Option<MessagePhase>,
    },
    StandaloneSpeech {
        text: String,
    },
    HandoffUpdate {
        handoff_id: String,
        text: String,
        phase: Option<MessagePhase>,
    },
    HandoffAppend {
        handoff_id: String,
        text: String,
        phase: Option<MessagePhase>,
    },
    CompletedHandoff {
        handoff_id: String,
        text: String,
        phase: Option<MessagePhase>,
    },
    ConversationItem {
        text: String,
        phase: Option<MessagePhase>,
    },
    HandoffCompleteAck {
        handoff_id: String,
    },
}
内部变体含义V3 最终适配
HandoffAppend流式追加delegation.context.append
HandoffUpdate / CompletedHandoff带交接 ID 的完整输出形态同样适配为 delegation context append,不是修改已发文本
StandaloneHandoff无活动交接的自动输出session.context.append
ConversationItem实验性的 item 交付方式developer 文本经适配进入 session context
StandaloneSpeech客户端显式要求送入可播报通道session context append,固定 speakable
HandoffCompleteAck不带正文的完成确认V3 无操作;相关确认语义属于 V2

这里的 Update 是 Core 内部命名。在 V3 线协议中,它仍是 append,不能让客户端据此假定能够覆盖先前送入模型的内容。

真实的流式出站分支先完成通道选择,再调用 writer:

源码文件:codex-rs/core/src/realtime_conversation.rs

相关函数/类型:handle_handoff_output(L2146–L2152,L2166–L2178,摘录)

rust
// 节选 V3 的手动 speech 与自动流式 append;原 text 不在这里去掉标签。
RealtimeOutbound::StandaloneSpeech { text } => {
    writer
        .clone()
        .with_context_append_channel(RealtimeContextAppendChannel::Speakable)
        .send_standalone_handoff(STANDALONE_HANDOFF_ID.to_string(), text)
        .await
}

// ...

RealtimeOutbound::HandoffAppend {
    handoff_id,
    text,
    phase,
} => {
    v3_output_writer(
        writer,
        phase.as_ref(),
        handoff_state.codex_response_handoff_mode,
    )
    .send_conversation_handoff_append(handoff_id, text)
    .await
}

phase 只参与选择 writer 的 channel,text 不被剥去 BEM 前缀。手动 speech 分支单独固定 speakable;完整输出、独立输出和 conversation item 也由同一 v3_output_writer 做自动模式选择。

对应的 JSON 数据类型显示:channel 是可省略字段,正文放在 content 数组里。

源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/protocol.rs

相关函数/类型:RealtimeContextAppendChannel、RealtimeOutboundMessage(L30–L36,L62–L74,摘录)

rust
// 第二段节选线协议 enum 的两个 context append 变体;None 会省略 channel 字段。
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum RealtimeContextAppendChannel {
    Speakable,
    Commentary,
}

// ...

#[serde(rename = "delegation.context.append")]
DelegationContextAppend {
    delegation_item_id: String,
    #[serde(skip_serializing_if = "Option::is_none")]
    channel: Option<RealtimeContextAppendChannel>,
    content: Vec<FramelessInputTextContent>,
},
#[serde(rename = "session.context.append")]
SessionContextAppend {
    #[serde(skip_serializing_if = "Option::is_none")]
    channel: Option<RealtimeContextAppendChannel>,
    content: Vec<FramelessInputTextContent>,
},

例如某次流式结果会变成:

示例代码:

json
{
  "type": "delegation.context.append",
  "delegation_item_id": "delegation-1",
  "channel": "commentary",
  "content": [
    { "type": "input_text", "text": "[COMMENTARY]正在查看日志" }
  ]
}

delegation_item_id 来自本 item 注册时的 handoff;channel 来自本次 writer 克隆;text 来自已分类、合并且可能截断的输出。三者在不同阶段形成,故障定位时应分别检查。

7.2 UTF-8 分片 ​

Core 的一次刷新还可能被 API writer 再分成多个文本帧。context_append_chunks 规定每段 text 不超过 500 字节。

源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods_frameless_bidi.rs

相关函数/类型:CONTEXT_APPEND_MAX_BYTES、context_append_chunks(L11–L12,L109–L126,摘录)

rust
// 500 字节只约束文本切片;每片在字符边界结束,之后才包装 JSON。
const CONTEXT_APPEND_MAX_BYTES: usize = 500;

// ...

pub(super) fn context_append_chunks(text: &str) -> Vec<String> {
    if text.len() <= CONTEXT_APPEND_MAX_BYTES {
        return vec![text.to_string()];
    }

    let mut chunks = Vec::new();
    let mut start = 0;
    while start < text.len() {
        let mut end = (start + CONTEXT_APPEND_MAX_BYTES).min(text.len());
        while end > start && !text.is_char_boundary(end) {
            end -= 1;
        }
        chunks.push(text[start..end].to_string());
        start = end;
    }
    chunks
}

算法沿字符串向前切片,若预算末尾落在多字节字符中,就回退到字符起点。每个输入字节只需被复制进一个 chunk,字符边界回退最多跨过当前字符的几个字节,整体时间与输出空间都随文本长度线性增长。返回的是一个 Vec<String>,并非不分配内存的借用迭代器。

例如 1,201 个 ASCII 字符分为 500、500、201 字节;200 个 🙂 共 800 字节,分成 500 与 300 字节的两段。上游测试同时断言每段长度上限和 chunks.concat() == text,证明字符完整且顺序可重建。

真正发送的循环如下:

源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs

相关函数/类型:RealtimeWebsocketWriter::send_json(L439–L473,摘录)

rust
// 同一次调用顺序等待每个片段;中途失败后不继续发送剩余片段。
async fn send_json(&self, message: &RealtimeOutboundMessage) -> Result<(), ApiError> {
    match message {
        RealtimeOutboundMessage::DelegationContextAppend {
            delegation_item_id,
            channel,
            content,
        } => {
            if let Some(content) = content.first() {
                // 每段继承相同的 delegation ID 与 channel。
                for chunk in context_append_chunks(&content.text) {
                    self.send_json_frame(&frameless_delegation_context_append_message(
                        delegation_item_id.clone(),
                        chunk,
                        *channel,
                    ))
                    .await?;
                }
                return Ok(());
            }
        }
        RealtimeOutboundMessage::SessionContextAppend { channel, content } => {
            if let Some(content) = content.first() {
                for chunk in context_append_chunks(&content.text) {
                    self.send_json_frame(&frameless_session_context_append_message(
                        chunk, *channel,
                    ))
                    .await?;
                }
                return Ok(());
            }
        }
        _ => {}
    }
    self.send_json_frame(message).await
}

500 字节限制只约束 content 中的文本,JSON 的字段、引号、转义和其他 envelope 开销还在其外。每帧沿用 delegation ID 与 channel;标签仅在整段文本开头,不会被重复添加到每个 500 字节片段。

await? 在某帧失败时立即结束剩余发送,已经成功写出的前缀帧不会回滚。结合第 6 节的“先排空缓冲,再发队列”,可以看到本地 item 去重、队列接受与远端完整接收之间没有事务边界。sideband 的断线重放机制见 实时会话架构,不能把这里的按序分片外推成 exactly-once 交付。

8. 转录时间线 ​

8.1 历史模式 ​

同一条 OutputTranscriptDelta 除了经传统 thread/realtime/transcript/delta 通知外,还可以产生带稳定 item ID 的时间线通知。第二条路径由 listener 构造时的历史模式快照控制:

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

相关函数/类型:ensure_listener_task_running(L249–L251,L334–L348,摘录)

rust
// listener 创建时快照 history_mode;只有 Paginated 才观察实时历史。
let config_snapshot = conversation.config_snapshot().await;
let realtime_history_enabled =
    matches!(config_snapshot.history_mode, ThreadHistoryMode::Paginated);

// ...

let (raw_events_enabled, realtime_effects) = {
    let mut thread_state = thread_state.lock().await;
    thread_state.track_current_turn_event(&event.id, &event.msg);
    let realtime_effects = if realtime_history_enabled
        && thread_state.realtime_history.should_observe(&event.msg)
    {
        let active_turn_id = thread_state.active_turn_snapshot().map(|turn| turn.id);
        thread_state
            .realtime_history
            .observe(&event.msg, active_turn_id.as_deref())
    } else {
        RealtimeEventEffects::default()
    };
    (thread_state.experimental_raw_events, realtime_effects)
};

ThreadHistoryMode 默认是 Legacy。只有 Paginated 才调用后面的历史观察器;BEM 的 bemTags 模式不控制这条支路。因此两种常见组合都合理:可以按 BEM 路由却不创建时间线,也可以在 V2 会话中创建时间线而不使用 BEM 标签路由。

完成后的时间线事实使用以下完整模型:

源码文件:codex-rs/protocol/src/realtime.rs

相关函数/类型:RealtimeItem、RealtimeItemContent、RealtimeSessionOutcome、RealtimeTranscriptRole、BemItemPresentation(L6–L59,摘录)

rust
// RealtimeItem 是独立时间线事实;BemItemPromoted 引用已有普通工作项。
/// A realtime thread item persisted in the canonical rollout.
#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize, JsonSchema, TS)]
pub struct RealtimeItem {
    pub id: String,
    pub realtime_session_id: String,
    #[serde(flatten)]
    pub content: RealtimeItemContent,
}

/// The minimum facts needed to interleave realtime speech and agent work.
#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize, JsonSchema, TS)]
#[serde(tag = "type", rename_all = "snake_case")]
#[ts(tag = "type", rename_all = "snake_case")]
pub enum RealtimeItemContent {
    RealtimeSessionStarted,
    TranscriptSegment {
        role: RealtimeTranscriptRole,
        text: String,
    },
    BemItemPromoted {
        turn_id: String,
        item_id: String,
        presentation: BemItemPresentation,
    },
    RealtimeSessionClosed {
        outcome: RealtimeSessionOutcome,
    },
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize, Serialize, JsonSchema, TS)]
#[serde(rename_all = "snake_case")]
#[ts(rename_all = "snake_case")]
pub enum RealtimeSessionOutcome {
    Ended,
    Failed,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize, Serialize, JsonSchema, TS)]
#[serde(rename_all = "snake_case")]
#[ts(rename_all = "snake_case")]
pub enum RealtimeTranscriptRole {
    User,
    Assistant,
}

/// How an existing agent item is presented in the realtime conversation.
#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize, JsonSchema, TS)]
#[serde(tag = "type", rename_all = "snake_case")]
#[ts(tag = "type", rename_all = "snake_case")]
pub enum BemItemPresentation {
    WholeItem,
    InlineMarkdown,
    InlineVisualization { index: u32 },
}

RealtimeItem.id 是本地事件 ID,realtime_session_id 关联实时会话。BemItemPromoted 另存普通 turn_id 与 item_id,表示展示已有工作项;其自身的 id 不能拿来查询普通 Agent item。

App Server 转换为 ThreadRealtimeItem 时保留这些关系,并把线上字段与 type 改为 camelCase:例如 realtimeSessionId、bemItemPromoted、inlineVisualization。这与 rollout 内的 snake_case 序列化是两层协议;生成的 TypeScript 类型也将 content flatten 成带 type 的联合。

8.2 观察器状态 ​

RealtimeHistoryState 属于 App Server 的 Thread 状态,保存的是当前处理所需的缓存和索引;历史查询由 rollout 的派生索引提供。

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

相关函数/类型:ActiveSegment、ActiveTranscriptSegments、RealtimeEventEffects、RealtimeHistoryState(L27–L41,L63–L99,摘录)

rust
// 观察器保留活动槽位、Turn 归属与去重集合;持久历史通过 rollout 索引查询。
#[derive(Debug, Clone, PartialEq, Eq)]
struct ActiveSegment {
    session_id: String,
    id: String,
    role: RealtimeTranscriptRole,
    text: String,
}

#[derive(Debug, Default)]
struct ActiveTranscriptSegments {
    user: Option<ActiveSegment>,
    assistant: Option<ActiveSegment>,
    first_active_role: Option<RealtimeTranscriptRole>,
}

// ...

#[derive(Debug)]
struct StreamingAgentMessage {
    item_id: String,
    text: String,
}

#[derive(Debug)]
pub(crate) struct RealtimeTranscriptStream {
    pub(crate) started_item: Option<RealtimeItem>,
    pub(crate) item_id: String,
    pub(crate) delta: String,
}

#[derive(Debug, Default)]
pub(crate) struct RealtimeEventEffects {
    pub(crate) items: Vec<RealtimeItem>,
    pub(crate) transcript_stream: Option<RealtimeTranscriptStream>,
}

#[derive(Clone, Copy)]
enum Continuation {
    Continue,
    Finish,
}

/// Retains only live session state; durable history is served by the rollout index.
#[derive(Debug, Default)]
pub(crate) struct RealtimeHistoryState {
    active_session_id: Option<String>,
    active_segments: ActiveTranscriptSegments,
    streaming_agent_message: Option<StreamingAgentMessage>,
    // 会话关闭后仍用于关联晚到的适用工作项。
    realtime_session_by_bem_turn: HashMap<String, String>,
    promoted_bem_presentation_keys: HashSet<String>,
    pending_handoffs: VecDeque<String>,
    failed: bool,
}
状态谁写入、谁读取
active_session_idStarted 设置、Closed 取走;决定新转录归属
active_segmentsuser/assistant 各一个活动片段,另记先开始的角色
streaming_agent_message暂存当前普通 Agent item 的文本,供展示指令识别
realtime_session_by_bem_turn关联普通 Turn 与实时会话;晚到的工作项可继续定位原会话
promoted_bem_presentation_keys已输出的 item 展示键,抑制 delta 与 completed 重复提升
pending_handoffs为后续 TurnStarted 保留会话归属
failed观察到 Realtime Error 后为 true,Closed 时决定 Failed/Ended 并重置

两类缓存的取舍不同:转录按说话角色保留两个槽位;普通 Agent 文本只有一个 streaming_agent_message,遇到新的 item ID 就清空旧文本。不能把它当作任意数量 Agent 消息可乱序交错的通用重组器。

观察器也不会在 realtime Closed 时清空 Turn 映射与所有去重键。should_observe 会继续接受已有关联 Turn 的 item 事件,使实时会话结束后才完成的适用产物仍能归到原会话;某些工具提升还另外要求当前实时会话活跃,下一节会具体区分。

会话与普通 Turn 的绑定来自下面三个事件分支,而不是从 BEM 文本反查:

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

相关函数/类型:RealtimeHistoryState::observe(L158–L177,L193–L197,L201–L211,摘录)

rust
// Started 绑定已有 Turn;后续 TurnStarted 优先消费待处理的交接会话,再回退到当前会话。
EventMsg::RealtimeConversationStarted(event) => {
    let session_id = event
        .realtime_session_id
        .clone()
        .unwrap_or_else(|| Uuid::now_v7().to_string());
    if self.active_session_id.as_deref() != Some(session_id.as_str()) {
        self.seal_segments(&mut items, Continuation::Finish);
        self.active_session_id = Some(session_id.clone());
        self.failed = false;
        items.push(RealtimeItem {
            id: Uuid::now_v7().to_string(),
            realtime_session_id: session_id.clone(),
            content: RealtimeItemContent::RealtimeSessionStarted,
        });
    }
    if let Some(turn_id) = active_turn_id {
        // 已有 Turn 的归属与当前实时连接是否仍活跃分开保存。
        self.realtime_session_by_bem_turn
            .insert(turn_id.to_string(), session_id);
    }
}

// ...

RealtimeEvent::HandoffRequested(_) => {
    if let Some(session_id) = &self.active_session_id {
        self.pending_handoffs.push_back(session_id.clone());
    }
}

// ...

EventMsg::TurnStarted(event) => {
    if let Some(session_id) = self
        .pending_handoffs
        .pop_front()
        .or_else(|| self.active_session_id.clone())
    {
        self.realtime_session_by_bem_turn
            .entry(event.turn_id.clone())
            .or_insert(session_id);
    }
}

Started 的会话 ID 缺省时生成 UUID v7;与当前 ID 不同时才生成新的会话开始项,并封口旧语音。若开始实时会话时已有普通 Turn,立即记录它的归属。后续 handoff 把会话 ID 放入 FIFO,TurnStarted 优先取出一个,再回退到当前会话;entry(...).or_insert(...) 保留已有 Turn 的归属。这套索引表达的是观察到的事件顺序,不能当作上游音频时间戳排序器。

8.3 增量与终值 ​

delta 第一次占用某个角色槽位时创建一个 UUID v7 片段,先发空内容的 started item,再发 delta;真正完成的内容尚未进入 effects.items。

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

相关函数/类型:RealtimeHistoryState::add_delta(L424–L455,摘录)

rust
// 按角色选槽位,同一片段的 started、delta、completed 共享一个 UUID。
fn add_delta(
    &mut self,
    role: RealtimeTranscriptRole,
    delta: &RealtimeTranscriptDelta,
) -> Option<RealtimeTranscriptStream> {
    let session_id = self.active_session_id.clone()?;
    self.active_segments.first_active_role.get_or_insert(role);
    let segment = self
        .active_segments
        .slot_mut(role)
        .get_or_insert_with(|| ActiveSegment {
            session_id,
            id: Uuid::now_v7().to_string(),
            role,
            text: String::new(),
        });
    // 是否发 started 取决于当前 text 为空,空 delta 因此是特殊边界。
    let started_item = segment.text.is_empty().then(|| RealtimeItem {
        id: segment.id.clone(),
        realtime_session_id: segment.session_id.clone(),
        content: RealtimeItemContent::TranscriptSegment {
            role: segment.role,
            text: String::new(),
        },
    });
    segment.text.push_str(&delta.delta);
    Some(RealtimeTranscriptStream {
        started_item,
        item_id: segment.id.clone(),
        delta: delta.delta.clone(),
    })
}

started_item 是独立可选值,因为后续 chunk 应复用同一 ID。first_active_role 保存两个角色同时活跃时的先后关系,结束或插入产物时按这个顺序封口,不靠 user/assistant 的固定排序。

实现判断“首次通知”的条件是 segment.text.is_empty(),不是另一个 started 标志。因此空 delta 不追加文字,并可能在下次 delta 再次触发同 ID 的 started;消费者与测试应注意这个条件,不能把 UUID 稳定性等同于每种空输入下都只通知一次。

Done 到达时则有两个不同分支:

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

相关函数/类型:RealtimeHistoryState::finish_segment(L456–L482,摘录)

rust
// 没有活动片段时从 done 补建;已有片段则封存累计 delta,不应用终值修订。
fn finish_segment(
    &mut self,
    items: &mut Vec<RealtimeItem>,
    role: RealtimeTranscriptRole,
    done: &RealtimeTranscriptDone,
) -> Option<RealtimeTranscriptStream> {
    let Some(segment) = self.active_segments.take(role) else {
        if done.text.is_empty() {
            return None;
        }
        let stream = self.add_delta(
            role,
            &RealtimeTranscriptDelta {
                delta: done.text.clone(),
            },
        );
        if let Some(segment) = self.active_segments.take(role) {
            self.seal_active_segment(items, segment, Continuation::Finish);
        }
        return stream;
    };
    // A split leaves an empty continuation; the upstream final may repeat
    // text that was already committed before the split.
    self.seal_active_segment(items, segment, Continuation::Finish);
    None
}

如果完全没有活动片段,非空 done 会补建一段,产生 started、delta 和完成项;如果已有片段,就保存累计的 delta 文本,忽略 done.text 的修订。这个选择避免已展示或已分段提交的内容被整体重放,但也意味着终值修订不会回写这段时间线。

下面的上游测试故意让增量与终值完全不同:

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

相关函数/类型:streams_stable_segment_items_and_ignores_upstream_transcript_revisions(L215–L250,摘录)

rust
// 用完全不同的 done 文本反证时间线不会把它当成已有片段的替换值。
#[test]
fn streams_stable_segment_items_and_ignores_upstream_transcript_revisions() {
    let mut state = started_state();
    let effects = observe_realtime(
        &mut state,
        RealtimeEvent::OutputTranscriptDelta(RealtimeTranscriptDelta {
            delta: "hello".to_string(),
        }),
    );
    assert!(effects.items.is_empty());
    let stream = effects.transcript_stream.expect("streaming delta");
    let started = stream.started_item.expect("segment start");
    assert_eq!(stream.item_id, started.id);
    assert_eq!(Uuid::parse_str(&started.id).unwrap().get_version_num(), 7);
    assert_eq!(stream.delta, "hello");

    let done = observe_realtime(
        &mut state,
        RealtimeEvent::OutputTranscriptDone(RealtimeTranscriptDone {
            text: "a substantially different upstream revision".to_string(),
        }),
    );
    assert_eq!(
        done.items,
        vec![RealtimeItem {
            id: started.id,
            realtime_session_id: "voice-1".to_string(),
            content: RealtimeItemContent::TranscriptSegment {
                role: RealtimeTranscriptRole::Assistant,
                text: "hello".to_string(),
            },
        }],
    );
    assert!(done.transcript_stream.is_none());
}

断言要求完成项仍然是 hello,ID 与第一条 started 一致。这是明确的消费者语义,不是“终值太晚到达”的竞态推测。相对地,API 层 apply_transcript_done 会在角色匹配且没有强制新项时替换其缓存的末条文本;它服务下次委派,所以同一会话的模型上下文与显示时间线可以持有不同版本的转录。

8.4 保存与通知 ​

listener 在 Thread 锁内运行 observe 得到 RealtimeEventEffects,然后释放锁,执行通知与保存。顺序是先应用历史 effects,再调用 apply_bespoke_event_handling 发传统实时通知。

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

相关函数/类型:apply_realtime_event_effects(L14–L50,摘录)

rust
// 增量通知先发送,完成项随后保存;保存失败在这一层只记 warning。
pub(crate) async fn apply_realtime_event_effects(
    conversation: &CodexThread,
    outgoing: &ThreadScopedOutgoingMessageSender,
    thread_id: ThreadId,
    effects: RealtimeEventEffects,
) {
    let thread_id = thread_id.to_string();

    if let Some(stream) = effects.transcript_stream {
        if let Some(item) = stream.started_item {
            outgoing
                .send_server_notification(ServerNotification::ThreadRealtimeItemStarted(
                    ThreadRealtimeItemStartedNotification {
                        thread_id: thread_id.clone(),
                        item: item.into(),
                    },
                ))
                .await;
        }
        outgoing
            .send_server_notification(ServerNotification::ThreadRealtimeItemTranscriptDelta(
                ThreadRealtimeItemTranscriptDeltaNotification {
                    thread_id: thread_id.clone(),
                    item_id: stream.item_id,
                    delta: stream.delta,
                },
            ))
            .await;
    }

    if let Err(error) =
        persist_realtime_items(conversation, outgoing, &thread_id, effects.items).await
    {
        warn!(thread_id, "failed to persist realtime history: {error}");
    }
}

对转录来说,客户端可以先收到 started 和 delta,稍后才等到 completed。effects 并不是持久化事务,只是本次事件计算出的“待执行结果”;锁只保护内存状态,不跨越这些异步 I/O。

完成项的保存与通知顺序如下:

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

相关函数/类型:persist_realtime_items(L51–L91,摘录)

rust
// 保存成功后才发送 completed;转录项的 started 已在增量阶段发送。
pub(crate) async fn persist_realtime_items(
    conversation: &CodexThread,
    outgoing: &ThreadScopedOutgoingMessageSender,
    thread_id: &str,
    items: Vec<RealtimeItem>,
) -> Result<(), String> {
    if items.is_empty() {
        return Ok(());
    }
    conversation
        .append_rollout_items(
            &items
                .iter()
                .cloned()
                .map(RolloutItem::RealtimeItem)
                .collect::<Vec<_>>(),
        )
        .await
        .map_err(|error| format!("failed to persist realtime history: {error}"))?;
    for item in items {
        // 会话边界与提升项没有增量阶段,在这里补齐 started。
        if !matches!(&item.content, RealtimeItemContent::TranscriptSegment { .. }) {
            outgoing
                .send_server_notification(ServerNotification::ThreadRealtimeItemStarted(
                    ThreadRealtimeItemStartedNotification {
                        thread_id: thread_id.to_string(),
                        item: item.clone().into(),
                    },
                ))
                .await;
        }
        outgoing
            .send_server_notification(ServerNotification::ThreadRealtimeItemCompleted(
                ThreadRealtimeItemCompletedNotification {
                    thread_id: thread_id.to_string(),
                    item: item.into(),
                },
            ))
            .await;
    }
    Ok(())
}

CodexThread::append_rollout_items 经 live thread 的 append_items 写入线程存储。只有它成功返回,才继续发送完成项通知。语音片段的 started 已在增量阶段发过,因此这里只给非转录项补 started,再为所有项发 completed。

保存失败会返回到 apply_realtime_event_effects,记录 warning;这里没有撤回已发 delta,也没有回滚已经更新的内存状态。完成通知证明线程存储的 append 调用已成功,不能单凭它推断掉电恢复等更强存储保证。

传统 thread/realtime/transcript/done 仍携带上游 done.text,而 thread/realtime/item/completed 使用观察器累积的片段。排查“同一时刻两种通知的文字不同”,首先应辨认订阅的是哪一条协议。

9. 产物提升 ​

9.1 展示引用 ​

promotion(产物提升) 指把普通 Turn 的某个工作项加入实时会话的展示序列。BemItemPresentation 分为整个 item、内联 Markdown、指定索引的内联可视化。它与 bemTags 的语音通道路由独立,不能因为两个名字都有 BEM 就合并解析器。

observe_item 规定哪些普通工作可以提升:

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

相关函数/类型:RealtimeHistoryState::observe_item(L281–L343,摘录)

rust
// 不同工作项有不同提升条件;MCP 与动态工具的成功判定不能混写。
fn observe_item(
    &mut self,
    items: &mut Vec<RealtimeItem>,
    turn_id: &str,
    item: &TurnItem,
    completed: bool,
) {
    match item {
        TurnItem::AgentMessage(message) => {
            let text = message
                .content
                .iter()
                .map(|content| match content {
                    AgentMessageContent::Text { text } => text.as_str(),
                })
                .collect::<String>();
            if !completed {
                self.streaming_agent_message = Some(StreamingAgentMessage {
                    item_id: message.id.clone(),
                    text: text.clone(),
                });
            }
            self.observe_assistant_message(items, turn_id, &message.id, &text);
        }
        TurnItem::ImageGeneration(image) => {
            self.add_promotion(items, turn_id, &image.id, BemItemPresentation::WholeItem);
        }
        TurnItem::Extension(extension)
            if serde_json::to_value(extension)
                .is_ok_and(|item| item["kind"] == "image_gen.generation") =>
        {
            self.add_promotion(
                items,
                turn_id,
                extension.id(),
                BemItemPresentation::WholeItem,
            );
        }
        TurnItem::SubAgentActivity(activity)
            if completed && activity.kind == SubAgentActivityKind::Started =>
        {
            self.add_promotion(items, turn_id, &activity.id, BemItemPresentation::WholeItem);
        }
        TurnItem::DynamicToolCall(call)
            if self.active_session_id.is_some()
                && completed
                && call.status == DynamicToolCallStatus::Completed
                && call.success == Some(true) =>
        {
            self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem);
        }
        TurnItem::McpToolCall(call)
            if self.active_session_id.is_some()
                && completed
                && call.server == "codex_app"
                && call.status == McpToolCallStatus::Completed =>
        {
            self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem);
        }
        _ => {}
    }
}
普通 item触发条件presentation
AgentMessage文本满足后面的展示指令InlineMarkdown 或 InlineVisualization
ImageGenerationstarted 或 completed 均可WholeItem
Extension序列化后的 kind 为 image_gen.generationWholeItem
SubAgentActivityitem completed 且 activity kind 是 StartedWholeItem
DynamicToolCall会话活跃、completed、状态 Completed、success 为 trueWholeItem
McpToolCall会话活跃、completed、server 是 codex_app、状态 CompletedWholeItem
其他 item不匹配不提升

“子 Agent 启动活动项已完成”不是“子 Agent 任务已完成”。MCP 分支也没有 DynamicToolCall 的 success == Some(true) 判断;应按各自类型的具体条件理解,不能笼统写成“所有成功工具都显示”。

9.2 文本指令 ​

Agent 消息的展示识别使用另一套规则:

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

相关函数/类型:RealtimeHistoryState::observe_assistant_message(L344–L390,摘录)

rust
// 展示指令解析独立于 BEM 通道路由;先检查内联 Markdown,再扫描可视化。
fn observe_assistant_message(
    &mut self,
    items: &mut Vec<RealtimeItem>,
    turn_id: &str,
    item_id: &str,
    text: &str,
) {
    let mut lines = text.trim_start().lines();
    let mut first = lines.next().unwrap_or_default();
    // 这里只跳过首行方括号头,并不验证默认或自定义 BEM 前缀。
    if first.starts_with('[')
        && let Some((_, content)) = first.split_once(']')
    {
        first = content.trim_start();
        if first.is_empty() {
            first = lines.next().unwrap_or_default();
        }
    }
    if first == INLINE_MARKDOWN_DIRECTIVE && text.contains('\n') {
        self.add_promotion(items, turn_id, item_id, BemItemPresentation::InlineMarkdown);
        return;
    }

    let mut in_fence = false;
    let mut visualization_index = 0;
    for line in text.lines() {
        let trimmed = line.trim_start();
        // 简单翻转围栏状态,不是完整 Markdown AST 解析。
        if trimmed.starts_with(BACKTICK_FENCE) || trimmed.starts_with(TILDE_FENCE) {
            in_fence = !in_fence;
            continue;
        }
        if !in_fence
            && (trimmed.starts_with(INLINE_VISUALIZATION_DIRECTIVE)
                || trimmed.starts_with(VISUALIZE_DIRECTIVE))
        {
            self.add_promotion(
                items,
                turn_id,
                item_id,
                BemItemPresentation::InlineVisualization {
                    index: visualization_index,
                },
            );
            visualization_index += 1;
        }
    }
}

内联 Markdown 要求去掉前导空白及可选的首行方括号头后,第一条有效行恰好等于 ::codex-realtime-inline{},并且原文本已经包含换行。首个 chunk 只有指令时不提升,换行和后续内容到达后才提升。

这里会跳过任意首行方括号前缀,不检查它是否等于 [ANALYSIS],也不读取自定义 BEM 前缀 map。于是上游测试中的小写 [analysis] 可以触发展示,却不能被第 4 节的默认 BEM 路由识别。这正是两个消费者必须分开阅读的原因。

可视化按行扫描,只识别围栏外以 ::codex-inline-vis{ 或 visualize{ 开始的行,依次编号 0、1、2。in_fence 遇到以三个反引号或波浪号开头的行就翻转;它不是完整 CommonMark 解析器,不检查围栏种类或长度配对,也不在这里验证参数 JSON。

每次普通消息 delta 都追加到缓存,再从头扫描累计文本。在 N 个字符逐步到达时,多次全量扫描的总工作量可以达到平方量级;如果分析性能,需要把扫描次数与累计长度同时考虑,不能只看到单次循环就称为整条流的 O(N)。

9.3 去重与续段 ​

最终添加展示项的入口先查 Turn 归属,再登记去重键:

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

相关函数/类型:RealtimeHistoryState::add_promotion(L391–L423,摘录)

rust
// 先登记展示键,再封口语音并生成引用项;真正保存由 effects 消费者执行。
fn add_promotion(
    &mut self,
    items: &mut Vec<RealtimeItem>,
    turn_id: &str,
    item_id: &str,
    presentation: BemItemPresentation,
) {
    let Some(realtime_session_id) = self.realtime_session_by_bem_turn.get(turn_id).cloned()
    else {
        return;
    };
    let presentation_key = match &presentation {
        BemItemPresentation::WholeItem => format!("{item_id}:whole-item"),
        BemItemPresentation::InlineMarkdown => format!("{item_id}:inline-markdown"),
        BemItemPresentation::InlineVisualization { index } => {
            format!("{item_id}:inline-visualization:{index}")
        }
    };
    // 同一 item 与展示形态只追加一次;这里还没有发生存储 I/O。
    if !self.promoted_bem_presentation_keys.insert(presentation_key) {
        return;
    }
    self.seal_segments(items, Continuation::Continue);
    items.push(RealtimeItem {
        id: Uuid::now_v7().to_string(),
        realtime_session_id,
        content: RealtimeItemContent::BemItemPromoted {
            turn_id: turn_id.to_string(),
            item_id: item_id.to_string(),
            presentation,
        },
    });
}

键由普通 item ID 与 presentation 组成;不同可视化索引是不同键。同一指令在多次 delta 扫描以及 completed 中反复出现,也只生成一次展示引用。键本身不包含 session ID 或 turn ID,因此它依赖普通 item ID 的区分性,不能当作按会话隔离的内容哈希。

seal_segments(Continue) 让正在累积的语音先封口,随后再追加产物;封口后还会创建一个空续段:

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

相关函数/类型:RealtimeHistoryState::seal_segments、RealtimeHistoryState::seal_active_segment(L483–L529,摘录)

rust
// Continue 保留空续段,记录本轮语音已被分段提交,避免 done 重放已提交文本。
fn seal_segments(&mut self, items: &mut Vec<RealtimeItem>, continuation: Continuation) {
    let mut segments = std::mem::take(&mut self.active_segments);
    let roles = match segments.first_active_role {
        Some(RealtimeTranscriptRole::Assistant) => [
            RealtimeTranscriptRole::Assistant,
            RealtimeTranscriptRole::User,
        ],
        Some(RealtimeTranscriptRole::User) | None => [
            RealtimeTranscriptRole::User,
            RealtimeTranscriptRole::Assistant,
        ],
    };
    for role in roles {
        if let Some(segment) = segments.take(role) {
            self.seal_active_segment(items, segment, continuation);
        }
    }
}

fn seal_active_segment(
    &mut self,
    items: &mut Vec<RealtimeItem>,
    segment: ActiveSegment,
    continuation: Continuation,
) {
    if !segment.text.is_empty() {
        items.push(RealtimeItem {
            id: segment.id,
            realtime_session_id: segment.session_id.clone(),
            content: RealtimeItemContent::TranscriptSegment {
                role: segment.role,
                text: segment.text,
            },
        });
    }
    // 续段使用新 UUID;空文本本身不会写为持久化转录项。
    if matches!(continuation, Continuation::Continue) {
        self.active_segments
            .first_active_role
            .get_or_insert(segment.role);
        *self.active_segments.slot_mut(segment.role) = Some(ActiveSegment {
            session_id: segment.session_id,
            id: Uuid::now_v7().to_string(),
            role: segment.role,
            text: String::new(),
        });
    }
}

保留空续段不是多余分配。假设已经播放 Already spoken,随后显示一个工具结果,上游最后发回 Already spoken, with a revision。若提升时直接删除语音槽位,finish_segment 会误认为“只收到最终转录”,重新创建整段,导致已提交语音再次出现。空续段保留了“这轮转录前半段已处理”的事实,done 只结束这个槽位。

下图展示一次“语音 → 产物 → done”的实际顺序。客户端 notifications 与 rollout 项分别标注,避免把它们混成同一个列表。

  • 语音在产物之前封口,所以显示顺序不依赖未来 done 的到达时间。
  • done 仍可触发传统转录通知,图中“无新 item”只指历史观察器的结果。
  • 若续段又收到新 delta,新语音会使用新的 UUID,与前半段分开保存。

does_not_replay_a_final_transcript_after_a_promotion_split 正是这个反例测试;它在提升后送入带修订的完整终值,断言没有新 items,也没有补发 transcript stream。promotes_backing_agent_artifacts_once_without_a_client_request 则分两次送入展示指令,在出现换行时提升,随后用完整 item 重放同一内容,断言不再新增。

去重键在保存操作之前已经进入 set。保存失败没有本层自动回滚;继续重放同一普通 item 也可能被视为已提升。处理“转录已显示但产物的 completed 没来”时,要检查线程存储 append 的错误,而不能仅凭第二次扫描无结果就认定模型没有生成产物。

前缀判定、UTF-8 分片与历史投影的实现没有按操作系统分支,因此这一算法层的平台差异不适用。设备录放音与 WebRTC 媒体支持属于外部客户端及连接层,不能由这些字符串与事件测试推出。

10. 差异定位 ​

把这篇的状态放回一个具体场景:V3、bemTags、自定义 commentary 为 [PROGRESS],模型输出 [PRO 后继续输出 GRESS]检查中,之后生成一张图,最后给出 [DONE]完成。排查时应分别回答三组问题,而不是只检查 WebSocket 是否连通。

现象先检查的条件或字段可得出的判断
普通回答在增长,sideband 无追加mode、前缀精确值、ChannelParser.phase未匹配文本在 item 完成前留在 parser
final 被送进 commentarymode 是否固定 commentary、自定义 analysis 是否为宽前缀模式优先;bemTags 下再按固定顺序匹配
BEM mode 正确却未流式回传streams_handoff_append、活动 handoff、item 是否注册client-managed 或 as-items 改变交付路径
超过约半个预算后暂停追加streamable_text_bytes、truncated、剩余缓存头部额度耗尽,最终块负责剩余文本或尾部
item 完成后再次发送全文finish 的返回值、是否曾注册及排出文本检查是否意外落入完整输出后备路径
实时转录有通知,时间线无记录Thread 的 historyModeLegacy 不运行历史观察器
done.text 与 completed.text 不同API 转录缓存与历史片段各自的终值规则一个可能修订缓存,一个保存已累积的 delta
展示指令只有字面文本首行、换行、围栏、普通 Turn 的会话映射该指令可能尚未满足 promotion 条件
产物保存失败后重试仍不显示promoted_bem_presentation_keys 与 append 错误内存去重先于保存,本层没有自动回滚

可以在源码仓库根目录用下面的搜索复述两条消费者路径:

bash
rg -n 'message_phase|struct ChannelParser|fn push|fn finish' \
  codex-rs/core/src/realtime_conversation/bem.rs
rg -n 'register_handoff_stream_item|finish_handoff_stream_item|v3_output_writer' \
  codex-rs/core/src/realtime_conversation.rs
rg -n 'finish_segment|observe_assistant_message|add_promotion|seal_active_segment' \
  codex-rs/app-server/src/realtime_history.rs

第一组定位分类与半标签;第二组定位普通 item 的归属和出站 channel;第三组定位同一文本的显示结果。注意搜索结果不会都落在名为 BEM 的文件里,状态跨越 API、Core 和 App Server。

源码 checkout 已按仓库要求准备 Rust 工具链后,可从 codex-rs 运行以下定向测试:

bash
just test --locked -p codex-core --lib \
  -E 'test(realtime_conversation::bem::tests::) | test(realtime_conversation::tests::streamed_)'
just test --locked -p codex-core --test all \
  -E 'test(conversation_flushes_assistant_deltas_every_200ms_for_v3_handoff) | test(conversation_handoff_persists_across_item_done_until_turn_complete)'
just test --locked -p codex-api --lib \
  -E 'test(protocol_frameless_bidi::tests::) | test(methods_frameless_bidi::tests::) | test(context_append)'
just test --locked -p codex-app-server --lib \
  -E 'test(realtime_history::tests::)'
just test --locked -p codex-app-server --test all \
  -E 'test(websocket_v3_routes_handoffs_by_session_mode) | test(realtime_timeline_splits_accepted_steering_and_persists_promoted_artifacts)'

BEM 单元测试验证匹配、覆盖与缓冲;Core 集成测试把 SSE 和 WebSocket 接起来;API 测试验证两种委派归一化、分片与字段编码;App Server 单元测试验证稳定 ID、终值修订和展示去重。socket 集成测试含 skip_if_no_network!,环境禁网时会提前返回,不能仅凭绿色结果推断本地 socket 场景实际执行过。

websocket_v3_routes_handoffs_by_session_mode 在四组设置下依次生成 analysis、commentary、final 和未知文本,对每次实际 sideband JSON 比较 channel 与未剥离的原文,再验证显式 speech 固定为 speakable。时间线集成测试则创建 Paginated Thread,用 gate 暂停普通输出,插入语音与被接受的 steering,再读取 thread/timeline/list,断言语音项在 steering 前,且存在关联 promoted-message 的展示项。这分别验证模型输入与客户端历史两个出口,不能互相替代。

最后尝试解释两个改动的后果:把 ChannelParser.finish 的缓冲直接丢弃,会让哪类回答丢失;把 seal_segments(Continue) 改成删除活动槽位,又会让哪类 done 重复出现。答案应沿着本文的真实调用者走到出站消息或时间线项,而不是只看修改位置附近的几行代码。需要继续追查普通文本如何变成 AgentMessageContentDelta 时,回到 流式推理与消息Delta;涉及连接关闭和发送失败后的重放,则接着阅读 实时会话架构 的断线恢复部分。