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 事件到达,通道信息可以编码在回答开头;实时服务另有入站事件,时间线又有独立的数据结构。
| 对象 | 生产者与消费者 | 表达的事实 |
|---|---|---|
RealtimeEvent | codex-api 解析服务端帧,Core 消费 | 转录、音频、委派、会话或响应事件 |
EventMsg::ItemStarted / AgentMessageContentDelta / ItemCompleted | 普通 Turn 产生,Session 镜像给实时子系统 | 一个 Agent 输出项的开始、文本追加和完成 |
RealtimeOutbound | Core 镜像与刷新逻辑产生,输入任务消费 | 准备送给实时模型的追加、完整结果或独立文本 |
RealtimeItem | App 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,摘录)
// 这是适配后的内部事件;可选 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 ID | V2 响应状态协调;取消映射为客户端 itemAdded |
ConversationItemAdded / Done | 原始 item 或 item ID | Added 可通知客户端;Done 不结束普通 Turn 的交接 |
HandoffRequested | handoff ID、item ID、输入和转录 | 启动或 steer 普通 Turn |
NoopRequested | call ID、item ID | V2 为静默工具调用回填空结果 |
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,摘录)
// 适配器先按 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.updated | session.id 字符串 | SessionUpdated;instructions 可选 |
input_transcript.added | item.text 字符串 | InputTranscriptDelta;不保留 item ID |
output_transcript.added | item.text 字符串 | OutputTranscriptDelta;不保留 item ID |
turn.done | turn.role、turn.transcript | user/assistant 分别对应两类 Done;其他 role 不产生事件 |
output_audio.delta | audio 字符串 | AudioOut;适配器填 24,000 Hz、单声道,不保留 start/end 时间 |
delegation.created | client 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,摘录)
// 只接收 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,摘录)
// 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,摘录)
// 这里只摘录 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,摘录)
// 启动参数省略 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,摘录)
// 这里只摘录 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 的通道路由 |
codexResponseHandoffMode | thinking / commentary / bemTags | 省略 channel、固定 commentary、按标签分类 |
codexResponseHandoffChannelPrefixes | 可选 map | 仅覆盖配置了的 analysis/commentary/final 通道 |
clientManagedHandoffs | false | true 时跳过自动回答回传;不阻止入站委派提交普通 Turn |
codexResponsesAsItems | false | true 时不走流式 handoff append,完成文本走 conversation item 适配 |
codexResponseItemPrefix | 无 | 只用于上一行的路径;BEM 判定发生在添加这个前缀之前 |
outputModality | 请求显式提供 | 本篇 V3 使用 audio;text 模式要求 V2 |
例如,已启用实时 feature 的 Thread 可提交以下请求。threadId 应使用 thread/start 返回的 ID;若还需要后面的持久化时间线,创建 Thread 时选择 historyMode: "paginated"。
示例代码:
{
"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,摘录)
// 在 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 |
| commentary | commentary | commentary | commentary | commentary |
| bemTags | commentary | commentary | speakable | speakable |
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,摘录)
// 依次检查三个通道;配置存在时替换该通道默认前缀。
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]完成 | Commentary | analysis 先匹配;不是最长前缀优先 |
默认;[FINAL][COMMENTARY]文本 | FinalAnswer | 只认起始前缀;后面的标签是文本 |
因此自定义前缀最好互不重叠。否则同一条回答落到哪个 channel 由匹配顺序决定,增加前缀的排列长度并不能得到“最具体匹配”。
4.2 半标签缓冲
网络上的模型文本可以在任何字符边界分块。ChannelParser 为一个输出项保留尚未分类的文本,并且一旦识别,就锁定该项的 phase。
源码文件:codex-rs/core/src/realtime_conversation/bem.rs
相关函数/类型:ChannelParser(L30–L67,摘录)
// 每个输出项独占一份 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") | None | buffer 为 [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,摘录)
// 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,摘录)
// 测试分别约束半标签释放时机和未知文本的 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,摘录)
// 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_tx | start_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,摘录)
// 这是单个普通 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,摘录)
// 常规 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,摘录)
// 是否流式交付与是否按标签分类是两个独立判断。
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,摘录)
// 只有活动 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,摘录)
// 先移除 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,摘录)
// 只有普通 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,摘录)
// 追加 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,摘录)
// 锁内取走文本,锁外等待队列;延迟任务重新按 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,摘录)
// 给尾部和截断标记留空间,不能把全部预算都用于提前发送。
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,摘录)
// 第一次超过预算时决定头尾;之后只移动尾部窗口。
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,摘录)
// 常规刷新受头部额度约束;最终刷新才释放剩余全文或标记加尾部。
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,摘录)
// 省略 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,摘录)
// 这是 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,摘录)
// 节选 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,摘录)
// 第二段节选线协议 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>,
},例如某次流式结果会变成:
示例代码:
{
"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,摘录)
// 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,摘录)
// 同一次调用顺序等待每个片段;中途失败后不继续发送剩余片段。
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,摘录)
// 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,摘录)
// 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,摘录)
// 观察器保留活动槽位、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_id | Started 设置、Closed 取走;决定新转录归属 |
active_segments | user/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,摘录)
// 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,摘录)
// 按角色选槽位,同一片段的 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,摘录)
// 没有活动片段时从 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,摘录)
// 用完全不同的 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,摘录)
// 增量通知先发送,完成项随后保存;保存失败在这一层只记 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,摘录)
// 保存成功后才发送 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,摘录)
// 不同工作项有不同提升条件;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 |
| ImageGeneration | started 或 completed 均可 | WholeItem |
| Extension | 序列化后的 kind 为 image_gen.generation | WholeItem |
| SubAgentActivity | item completed 且 activity kind 是 Started | WholeItem |
| DynamicToolCall | 会话活跃、completed、状态 Completed、success 为 true | WholeItem |
| McpToolCall | 会话活跃、completed、server 是 codex_app、状态 Completed | WholeItem |
| 其他 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,摘录)
// 展示指令解析独立于 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,摘录)
// 先登记展示键,再封口语音并生成引用项;真正保存由 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,摘录)
// 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 被送进 commentary | mode 是否固定 commentary、自定义 analysis 是否为宽前缀 | 模式优先;bemTags 下再按固定顺序匹配 |
| BEM mode 正确却未流式回传 | streams_handoff_append、活动 handoff、item 是否注册 | client-managed 或 as-items 改变交付路径 |
| 超过约半个预算后暂停追加 | streamable_text_bytes、truncated、剩余缓存 | 头部额度耗尽,最终块负责剩余文本或尾部 |
| item 完成后再次发送全文 | finish 的返回值、是否曾注册及排出文本 | 检查是否意外落入完整输出后备路径 |
| 实时转录有通知,时间线无记录 | Thread 的 historyMode | Legacy 不运行历史观察器 |
| done.text 与 completed.text 不同 | API 转录缓存与历史片段各自的终值规则 | 一个可能修订缓存,一个保存已累积的 delta |
| 展示指令只有字面文本 | 首行、换行、围栏、普通 Turn 的会话映射 | 该指令可能尚未满足 promotion 条件 |
| 产物保存失败后重试仍不显示 | promoted_bem_presentation_keys 与 append 错误 | 内存去重先于保存,本层没有自动回滚 |
可以在源码仓库根目录用下面的搜索复述两条消费者路径:
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 运行以下定向测试:
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;涉及连接关闭和发送失败后的重放,则接着阅读 实时会话架构 的断线恢复部分。
