实时会话架构
用户一边说话,一边让 Codex 修改文件:语音仍在往返,后台已经开始模型推理和工具执行。这里至少有两种不同的“对话”:实时服务处理音频、转录和语音响应,Codex 的普通 Turn 负责工作区里的任务。本文要追踪的就是两者如何在同一个 Thread 中交接,以及用户停止说话、重新开始会话、网络断开时,究竟结束了哪一组资源。
读者需要了解 Rust 的 async/await、Arc 与 channel。如果还不熟悉 Thread、Session、Turn 的层级,可以先读 Thread与Turn概念模型;ModelClient结构 解释普通模型请求的客户端生命周期;Session输入队列 解释输入怎样进入 Core。
本文沿 App Server 的真实入口走到 Core 管理器、WebSocket 收发任务和通知消费者,重点分析并发、状态与清理。音频设备采集、外部客户端的 WebRTC 媒体实现、BEM 标签语法和提示词裁剪算法不在本篇范围内。读完后,应能解释为什么“启动响应成功”还可能没有音频、为什么一条转录不会自动启动 Turn,以及为什么有些断线会重连、有些直接关闭。
1. 两条交互链
Thread 是用户看到的工作对话身份;Session 是 Core 中持有该 Thread 运行状态和服务的对象;Turn 是一次普通任务执行。Realtime 会话则有自己的连接、音频队列、文本队列和后台任务。它可以跨越多个普通 Turn,也可以在没有 Turn 运行时持续接收音频。
实时服务出现 HandoffRequested 时,含义是“把这段工作交给 Codex”。这里的 handoff 是任务交接,不是直接在音频处理函数里执行 shell。Core 将其包装成文本输入,提交给普通 Turn;普通 Turn 产生的回答再送回实时服务。底层 WebSocket frame、业务音频帧、实时协议里的 response 和 Codex Turn 分属不同层级,不能以一个“流”字概括。
下图区分数据的两个去向。经过 codex-api 的链路维持实时交互,经过 StartOrSteer 的链路执行普通任务;它们通过委派文本和回答事件衔接。
图中的实时路径不经过普通 ModelClientSession 的 Responses 采样循环。ModelClient 仍参与 WebRTC call 创建和认证,但它的存在不意味着两条连接共用普通模型请求的重试、完成条件或历史。
还有一个很实际的产品边界:当前 codex-rs/tui/src/chatwidget/protocol.rs 将 ThreadRealtimeStarted、ThreadRealtimeOutputAudioDelta、ThreadRealtimeClosed 等通知归入空处理分支。源码中的实时 API 支持,不能直接推导成当前 TUI 已提供麦克风采集和扬声器播放。沿本篇操作的客户端需要自己消费通知;外部客户端怎样呈现,不由 Core 管理器决定。
2. 入口门控
2.1 两层开关
App Server 的 thread/realtime/start、appendAudio、appendText、appendSpeech、stop 和 listVoices 在 codex-rs/app-server-protocol/src/protocol/common.rs 中带有实验 API 标记。连接必须先通过初始化,并在初始化能力中启用 experimentalApi,请求才会进入对应处理器。
源码文件:codex-rs/app-server/src/message_processor.rs
相关函数/类型:dispatch_initialized_client_request(源码定位,节选)
// 先完成连接初始化,再检查客户端是否显式接受实验 API。
if !session.initialized() {
return Err(invalid_request("Not initialized"));
}
if let Some(reason) = codex_request.experimental_reason()
&& !session.experimental_api_enabled()
{
return Err(invalid_request(experimental_required_message(reason)));
}这段判断保护的是 App Server 连接层。启用实验 API 只表示客户端接受实验接口;它不替代 Thread 自身的功能开关。TurnProcessor::prepare_realtime_conversation_thread 还会加载 Thread、确保事件监听器已连接,并检查 thread.enabled(Feature::RealtimeConversation),不满足时返回请求错误。
源码文件:codex-rs/features/src/lib.rs
相关函数/类型:FEATURES 中的 Feature::RealtimeConversation(源码定位,节选)
FeatureSpec {
id: Feature::RealtimeConversation,
key: "realtime_conversation",
stage: Stage::UnderDevelopment,
default_enabled: false,
},realtime_conversation 默认关闭,阶段为 UnderDevelopment。因此定位入口失败时,要分别检查连接初始化和线程配置,不能只看到一个开关已开便断定具备能力。Core 内部的直接 Op 调用与 App Server 入口不是同一层;下面的提交循环没有重复这两层 App Server 检查。
2.2 请求与事件
thread_realtime_start_inner 将公开参数转换成 ConversationStartParams,通过 submit_core_op 调用 CodexThread::submit_with_trace。成功返回空的 ThreadRealtimeStartResponse 表示操作已提交。Core 后续能否建立连接,由实时事件通知报告;这个返回值不是媒体就绪信号。
例如公开参数里的 includeStartupContext 未填写时,普通新会话默认开启,existingCall 默认关闭;clientManagedHandoffs、flushTranscriptTailOnSessionEnd 和 codexResponsesAsItems 默认都是 false。这些值在入口转换时就已确定,后续代码处理的是明确的布尔值。
Core 的提交循环依次消费操作,实时操作有独立分支:
源码文件:codex-rs/core/src/session/handlers.rs
相关函数/类型:submission_loop 的实时操作分支(源码定位,节选)
Op::RealtimeConversationStart(params) => {
if let Err(err) =
handle_realtime_conversation_start(&sess, sub.id.clone(), params).await
{
sess.send_event_raw(Event {
id: sub.id.clone(),
msg: EventMsg::Error(ErrorEvent {
message: err.to_string(),
codex_error_info: Some(CodexErrorInfo::Other),
}),
})
.await;
}
false
}
Op::RealtimeConversationAudio(params) => {
handle_realtime_conversation_audio(&sess, sub.id.clone(), params).await;
false
}
Op::RealtimeConversationText(params) => {
handle_realtime_conversation_text(&sess, sub.id.clone(), params).await;
false
}
Op::RealtimeConversationSpeech(params) => {
handle_realtime_conversation_speech(&sess, sub.id.clone(), params).await;
false
}
Op::RealtimeConversationClose => {
handle_realtime_conversation_close(&sess, sub.id.clone()).await;
false
}Start、Audio、Text、Close 的处理函数不会在这里启动一个普通模型 Turn。它们建立或访问 sess.conversation,然后继续消费下一条操作。普通 Turn 的创建发生在后文的委派路径。
顺序也带来边界:这里会 await Start 处理。若普通 WebSocket 的初次连接,或 existingCall 的首次 sideband 加入还没有返回,同一提交循环中的 Close 尚不能执行。新建 WebRTC call 的 sideband 之所以能在握手未完成时响应 Close,是因为那段连接过程被放进了后台任务;不能把它的取消测试外推到所有启动阶段。
3. 连接契约
3.1 传输与版本
先将两个独立维度分开:传输说明“建立哪种连接”,版本说明“使用哪套实时协议”。公开版本有三种,传输也有三种:
源码文件:codex-rs/protocol/src/protocol.rs
相关函数/类型:ConversationStartTransport、RealtimeConversationVersion(源码定位,节选)
#[derive(Debug, Clone, PartialEq)]
pub enum ConversationStartTransport {
Websocket,
Webrtc { sdp: String },
ExistingCall { call_id: String },
}
// ...
#[derive(Debug, Clone, Copy, Default, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)]
#[serde(rename_all = "snake_case")]
pub enum RealtimeConversationVersion {
V1,
#[default]
V2,
V3,
}Webrtc { sdp } 携带客户端的 SDP offer,Core 请求服务创建 call 并返回 SDP answer;ExistingCall { call_id } 表示客户端已经创建 call,Core 只加入其控制连接。sideband 在这里指附着到同一 WebRTC call 的 WebSocket 控制通道;本地 Core 并不因此拥有客户端的麦克风、扬声器或媒体轨道。
默认值需要沿转换代码读,不能只看 RealtimeConversationVersion 的 #[default]:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:prepare_realtime_start(源码定位,节选)
let transport = params
.transport
.clone()
.unwrap_or(ConversationStartTransport::Websocket);
// ...
let version = params.version.unwrap_or(match &transport {
ConversationStartTransport::Websocket => config.realtime.version,
ConversationStartTransport::Webrtc { .. }
| ConversationStartTransport::ExistingCall { .. } => RealtimeWsVersion::V1,
});普通 WebSocket 沿用 config.realtime.version,而 WebRTC 和 existingCall 在没有显式版本时使用 V1。validate_avas_webrtc_start 拒绝 V2,并要求 WebRTC 创建路径使用 conversational 模式。ExistingCall 同样拒绝 V2,还拒绝 prompt、model、voice、initial items 等会重配客户端会话的选项。
Core 内部另有一个只有 V1、V2 的 RealtimeSessionKind。它描述管理逻辑的分支,不能当成公开协议版本:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:start_inner(源码定位,节选)
// 内部 session_kind 只有两种;V3 仍需结合 event_parser 判断。
let event_parser = session_config.event_parser;
let session_kind = match event_parser {
RealtimeEventParser::V1 | RealtimeEventParser::FramelessBidi => RealtimeSessionKind::V1,
RealtimeEventParser::RealtimeV2 => RealtimeSessionKind::V2,
};所以 V3 的连接会进入内部 V1 分支,同时保留 event_parser = FramelessBidi。音频打断、handoff 完成、输出格式和重连会查看不同字段。修改代码时,遇到 session_kind == V1 必须继续确认它是否同时覆盖 V3,不能把它替换成“旧版协议”的口头解释。
| 公开版本 | 事件解析器 | 内部 kind | 新普通 WebSocket 的默认模型 |
|---|---|---|---|
| V1 | V1 | V1 | gpt-realtime-1.5 |
| V2 | RealtimeV2 | V2 | gpt-realtime-1.5 |
| V3 | FramelessBidi | V1 | gpt-live-1-codex |
表中的模型由 build_realtime_session_config 选择;请求级 model 优先于配置中的实验模型覆盖值,再落到默认值。它与普通工作 Turn 使用的模型是两个选择。输出 modality 也有条件:文本输出只允许 V2;initial items 只允许 V3。这些校验在创建运行时之前执行。
3.2 初始化边界
新 WebSocket 使用实时 API key 路径:realtime_api_key 依次查询 Provider 的 key、实验 bearer token、当前认证中的 API key,以及 OpenAI Provider 的环境变量回退。只有 ChatGPT 登录并不保证此路径可用。WebRTC call 创建及已有 call 附着则经过 ModelClient 的认证处理,不能将普通 WebSocket 的 key 要求套给所有传输。
start_inner 的三条分支与通知之间有如下关系:
| 传输 | Start 处理等待什么 | Core Started 时已完成什么 | 后续仍可能失败什么 |
|---|---|---|---|
| 新 WebSocket | RealtimeWebsocketClient::connect | socket 建立和初始化发送;V3 还等待 session.started | 持续收发、服务端业务事件 |
| 新 WebRTC call | ModelClient::create_realtime_call_with_headers | call 创建、获得 answer、启动 sideband task | 后台 sideband 握手与客户端媒体连接 |
| existingCall | existing_call::attach 中的首次 sideband 连接 | 已附着控制 socket,没有重建 call | 后续控制流与外部媒体生命周期 |
API 层实际决定是否发送 session 初始化消息。注意 send_session_update 是统一方法名,具体 wire 消息由 parser 决定:
源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs
相关函数/类型:connect_realtime_websocket_url(源码定位,节选)
// 初始化策略由连接用途决定,不能只看协议版本。
let initialize_session = match session_initialization {
RealtimeSessionInitialization::NewSession => true,
RealtimeSessionInitialization::LegacyWebrtcSideband => {
config.event_parser != RealtimeEventParser::FramelessBidi
}
RealtimeSessionInitialization::ExistingCall => false,
};
if initialize_session {
debug!(
session_id = config.session_id.as_deref().unwrap_or("<none>"),
"realtime websocket sending session.update"
);
connection
.writer
.send_session_update(
config.instructions,
config.initial_items,
config.session_mode,
config.output_modality,
config.voice,
config.delegation_ack_filler,
)
.await?;
}
if matches!(
session_initialization,
RealtimeSessionInitialization::NewSession
) && config.event_parser == RealtimeEventParser::FramelessBidi
{
connection.events.wait_for_session_started().await?;
}新 WebSocket 总会初始化;传统 WebRTC sideband 对 V3 跳过初始化,因为 call 已经在创建时配置;existingCall 永远不覆盖已有 session 配置。只有 新 V3 WebSocket 额外等待 wait_for_session_started。该方法要求首个可识别业务事件能映射为 SessionUpdated,再把它放回 pending events 供上层消费;如果先收到其他业务事件则报错,不会无限等待正确事件。因此“连接函数返回”在不同路径上的保证并不完全相同。
Core 发通知的顺序也直接写在调用方:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:handle_start_inner(源码定位,节选)
let start_output = sess.conversation.start(start, mode_instructions).await?;
info!("realtime conversation started");
sess.send_event_raw(Event {
id: sub_id.to_string(),
msg: EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent {
realtime_session_id: requested_realtime_session_id,
version,
}),
})
.await;
let RealtimeStartOutput {
realtime_active,
events_rx,
transcript_tail_rx,
sdp,
} = start_output;
if let Some(sdp) = sdp {
sess.send_event_raw(Event {
id: sub_id.to_string(),
msg: EventMsg::RealtimeConversationSdp(RealtimeConversationSdpEvent { sdp }),
})
.await;
}Started 先于 SDP 通知。对于新 WebRTC call,sideband task 在这之前已创建,却可能还在握手。下图将后台控制连接与通知交付分开;两者并发推进,不能画成收到 Started 后才开始连接。
这张图给出了一个诊断顺序:已收到 SDP,但没有后续实时事件,应检查 sideband;已收到音频通知但没有声音,应检查外部客户端的消费;API 请求已返回而收到 realtime/error,应检查准备或建连错误,而不是把空响应解释成整个启动过程成功。
4. 运行时所有权
4.1 状态与通道
Session 的 conversation 字段是 Arc<RealtimeConversationManager>。管理器持有当前运行时的槽位,以及在进入、退出实时模式时使用的指令。实际状态如下:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:RealtimeConversationManager、ConversationState(源码定位,节选)
// 状态锁保护的是所有权槽位;socket 位于输入任务下层。
pub(crate) struct RealtimeConversationManager {
state: Mutex<Option<ConversationState>>,
mode_instructions: Mutex<Option<RealtimeModeInstructions>>,
}
#[derive(Clone, Debug)]
pub(crate) struct RealtimeModeInstructions {
pub(crate) start: Option<String>,
pub(crate) end: Option<String>,
}
// ...
struct ConversationState {
audio_tx: Sender<RealtimeAudioFrame>,
text_tx: Sender<ConversationTextParams>,
session_kind: RealtimeSessionKind,
handoff: RealtimeHandoffState,
input_task: JoinHandle<()>,
fanout_task: Option<JoinHandle<()>>,
realtime_active: Arc<AtomicBool>,
stop_token: CancellationToken,
}Mutex<Option<ConversationState>> 让启动和停止能够取走整组资源。input_task 是输入仲裁与实时收发的业务任务;fanout_task 消费业务事件、路由委派并向外发送事件。这里没有一个直接保存 socket 的 writer 字段,socket 的所有权在 input task 持有的 API 连接下层。
mode_instructions 与活跃连接分开保存,停止时不会随 ConversationState 一起取走。这样退出实时模式时仍可读取相关指令。它不是判断会话仍在运行的依据;成功开始新会话时才覆盖这份指令。
创建通道时,容量以消息数为单位:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:start_inner(源码定位,节选)
let (audio_tx, audio_rx) =
async_channel::bounded::<RealtimeAudioFrame>(AUDIO_IN_QUEUE_CAPACITY);
let (text_tx, text_rx) =
async_channel::bounded::<ConversationTextParams>(TEXT_IN_QUEUE_CAPACITY);
let (handoff_output_tx, handoff_output_rx) =
async_channel::bounded::<RealtimeOutbound>(HANDOFF_OUT_QUEUE_CAPACITY);
let (events_tx, events_rx) =
async_channel::bounded::<RealtimeEvent>(OUTPUT_EVENTS_QUEUE_CAPACITY);
let (transcript_tail_tx, transcript_tail_rx) = async_channel::bounded::<String>(1);四个常量分别为音频 256、文本 64、handoff 输出 64、输出事件 256,另有容量为 1 的 transcript tail 通道。tail 保存结束时尚未交接的转录,由 fanout 在普通事件排空后处理。容量 256 本身不能换算成几秒音频;只有另外证明每帧时长固定,才有这种换算。
所有权图把 Core 任务与 API 收发实现连接起来。实线组合关系表示结构中持有,虚线表示任务执行时使用或消息传递。
图中的 Fanout 是闭包任务的角色名称,并不是源码里另有同名 struct。Manager、Writer、Events 分别对应前文管理器和下文两个长名称的类型。保留这些缩写是为了看清所有权,而不能据此在仓库里寻找不存在的实现文件。
槽位中有对象与对象仍活跃也要分开:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:running_state、is_running_v2(源码定位,节选)
pub(crate) async fn running_state(&self) -> Option<()> {
let state = self.state.lock().await;
state
.as_ref()
.and_then(|state| state.realtime_active.load(Ordering::Relaxed).then_some(()))
}
pub(crate) async fn is_running_v2(&self) -> bool {
let state = self.state.lock().await;
matches!(
state.as_ref(),
Some(state)
if state.realtime_active.load(Ordering::Relaxed)
&& state.session_kind == RealtimeSessionKind::V2
)
}running_state 必须同时看到 Some(state) 和原子 active 标志为真。停止和异步退出之间可能存在“对象仍被任务持有,但已不再对外运行”的阶段。锁保护状态槽位,原子标志负责活跃性判定;Relaxed 标志本身不承担发布所有字段的同步屏障。
普通 TurnContext.realtime_active 则在 codex-rs/core/src/session/turn_context.rs 构造上下文时读取一次 running_state。它是布尔快照,不是该原子对象的共享引用;不能用某个已创建 Turn 中的值替代管理器的当前状态查询。codex-rs/core/src/session/world_state.rs 结合这个快照和 mode instructions 生成普通任务需要的实时模式上下文。
4.2 收发泵
API 连接把发送端和业务事件接收端拆开:
源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs
相关函数/类型:RealtimeWebsocketConnection、RealtimeWebsocketWriter、RealtimeWebsocketEvents(源码定位,节选)
pub struct RealtimeWebsocketConnection {
writer: RealtimeWebsocketWriter,
events: RealtimeWebsocketEvents,
}
#[derive(Clone)]
pub struct RealtimeWebsocketWriter {
stream: Arc<WsStream>,
is_closed: Arc<AtomicBool>,
event_parser: RealtimeEventParser,
context_append_channel: Option<RealtimeContextAppendChannel>,
}
#[derive(Clone)]
pub struct RealtimeWebsocketEvents {
rx_message: async_channel::Receiver<Result<Message, WsError>>,
pending_events: Arc<Mutex<VecDeque<RealtimeEvent>>>,
transcript_state: RealtimeTranscriptState,
event_parser: RealtimeEventParser,
is_closed: Arc<AtomicBool>,
}克隆 writer 会共享 Arc<WsStream>;克隆 events 会共享消息接收通道、pending events、转录状态和关闭标志。async_channel::Receiver 的克隆不是广播订阅,多个消费者会竞争消息。正常业务输入任务使用一组收发端,不能为了增加监听者就随意克隆 receiver 并期待每个观察者都收到全部事件。
下层 WsStream::new 启动 pump task,由这个任务独占 WebSocketStream。发送者把命令放入容量 32 的命令通道,并等一个 oneshot 返回本次发送结果;接收的 socket 消息则写进另一条通道:
源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs
相关函数/类型:WsStream::new、WsStream::request(源码定位,节选)
// pump 独占 socket;oneshot 返回本次写入结果,接收端不持有 socket 锁。
let (tx_command, mut rx_command) = mpsc::channel::<WsCommand>(32);
let (tx_message, rx_message) = async_channel::unbounded::<Result<Message, WsError>>();
// ...
WsCommand::Send { message, tx_result } => {
debug!("realtime websocket sending message");
let result = inner.send(message).await;
let should_break = result.is_err();
if let Err(err) = &result {
error!("realtime websocket send failed: {err}");
}
let _ = tx_result.send(result);
if should_break {
break;
}
}
// ...
async fn request(
&self,
make_command: impl FnOnce(oneshot::Sender<Result<(), WsError>>) -> WsCommand,
) -> Result<(), WsError> {
let (tx_result, rx_result) = oneshot::channel();
if self.tx_command.send(make_command(tx_result)).await.is_err() {
return Err(WsError::ConnectionClosed);
}
rx_result.await.unwrap_or(Err(WsError::ConnectionClosed))
}未列出的 inner.next() 分支与发送命令一起位于同一个 tokio::select! 中;它处理 Ping/Pong,并把文本、关闭和错误交给 tx_message。这样等待一个业务事件,不会持有一把包住整个 socket 的锁而阻止另一端发送。命令通道关闭或 oneshot 发送者消失都会被映射为连接关闭。
这里还存在一个容易遗漏的容量边界:API 的原始入站消息通道是 unbounded,Core 的 256 个输出事件槽位位于解析之后。因此后端事件生产远快于消费时,不能仅凭 Core 的有界 channel 宣称整个入站路径内存有界。
反向证据可以读 realtime_ws_e2e_send_while_next_event_waits:mock 服务先等音频请求,再返回 session.updated;客户端同时执行发送和接收,并限制发送必须在 200 ms 内返回。
源码文件:codex-rs/codex-api/tests/realtime_websocket_e2e.rs
相关函数/类型:realtime_ws_e2e_send_while_next_event_waits(源码定位,节选)
let (send_result, next_result) = tokio::join!(
async {
tokio::time::timeout(
Duration::from_millis(200),
connection.send_audio_frame(RealtimeAudioFrame {
data: "AQID".to_string(),
sample_rate: 48000,
num_channels: 1,
samples_per_channel: Some(960),
item_id: None,
}),
)
.await
},
connection.next_event()
);
send_result
.expect("send should not block on next_event")
.expect("send audio");
let next_event = next_result.expect("next event").expect("event");
assert_eq!(
next_event,
RealtimeEvent::SessionUpdated {
realtime_session_id: "sess_after_send".to_string(),
instructions: Some("backend prompt".to_string()),
}
);如果 next_event() 在等待服务端回复时锁住 socket,服务端等不到音频,发送就会超时。这个测试因此能检出读等待读写互锁;它不证明真实网络上的音频延迟一定低于 200 ms,也不证明任意消费者速度下都没有背压。
5. 输入仲裁
5.1 丢帧与背压
音频和文本虽然都进入 input task,生产端策略却不同。音频使用非等待的 try_send:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:RealtimeConversationManager::audio_in(源码定位,节选)
pub(crate) async fn audio_in(&self, frame: RealtimeAudioFrame) -> CodexResult<()> {
let sender = {
let guard = self.state.lock().await;
guard.as_ref().map(|state| state.audio_tx.clone())
};
let Some(sender) = sender else {
return Err(CodexErr::InvalidRequest(
"conversation is not running".to_string(),
));
};
match sender.try_send(frame) {
Ok(()) => Ok(()),
Err(TrySendError::Full(_)) => {
warn!("dropping input audio frame due to full queue");
Ok(())
}
Err(TrySendError::Closed(_)) => Err(CodexErr::InvalidRequest(
"conversation is not running".to_string(),
)),
}
}取 sender 时只短暂持有 state 锁,然后在锁外发送。满队列时丢弃的是这次新来的帧,已有队列内容没有被替换;函数记录 warning 并返回 Ok(())。这使音频生产者不必为队列容量停住,却也意味着调用成功不能证明每一帧都到达服务端。关闭 channel 或没有状态时才报告 conversation is not running。
文本的策略是等待容量,同时在内部 V2 的用户文本前加 [USER] 前缀:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:RealtimeConversationManager::text_in(源码定位,节选)
pub(crate) async fn text_in(&self, mut params: ConversationTextParams) -> CodexResult<()> {
let sender = {
let guard = self.state.lock().await;
guard
.as_ref()
.map(|state| (state.text_tx.clone(), state.session_kind))
};
let Some((sender, session_kind)) = sender else {
return Err(CodexErr::InvalidRequest(
"conversation is not running".to_string(),
));
};
if params.role == ConversationTextRole::User {
params.text =
prefix_realtime_text(params.text, REALTIME_USER_TEXT_PREFIX, session_kind);
}
sender
.send(params)
.await
.map_err(|_| CodexErr::InvalidRequest("conversation is not running".to_string()))?;
Ok(())
}prefix_realtime_text 只对内部 V2 的非空文本添加该前缀,已经带相同前缀时也直接返回;V1 和 V3 保留原文本。Developer 角色不会被当成用户文本加前缀。这里必须在锁外 send(...).await:若持有 state 锁等待队列空位,其他任务的停止逻辑需要拿同一把锁才能取走状态,就可能互相阻塞。
handoff 输出同样通过异步 send 等待容量。音频的“丢帧但成功”、文本和回答的“等待空位”是源码明确选择的不同契约,不能统一写成“队列满后重试”。
5.2 五路输入
run_realtime_input_task 同时监听取消、远端事件、用户文本、后台回答和用户音频。下列摘录保留分支和失败出口,连接重建时的待发重放放到后文说明。
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:run_realtime_input_task(源码定位,节选)
// 五个分支共享一个输入循环;每个分支被选中后仍会 await 自己的处理函数。
loop {
let result = tokio::select! {
_ = stop_token.cancelled() => Err(RealtimeInputTaskExit::Terminal),
realtime_event = events.next_event() => {
match realtime_event {
Ok(Some(event)) => {
handle_realtime_server_event(
event,
&writer,
&events_tx,
&handoff_state,
session_kind,
&mut output_audio_state,
&mut response_create_queue,
)
.await
.map_err(classify_realtime_input_error)
}
Ok(None) => Err(RealtimeInputTaskExit::Terminal),
Err(err) => Err(RealtimeInputTaskExit::TransportLost {
err,
pending_outbound: None,
}),
}
}
// Text input that should be sent into realtime.
text = text_rx.recv() => {
let pending_outbound = text
.as_ref()
.ok()
.cloned()
.map(RealtimePendingOutbound::Text)
.map(Box::new);
handle_text_input(
text,
&writer,
)
.await
.map_err(|err| {
classify_realtime_input_error_with_pending(err, pending_outbound)
})
}
// Background agent progress or final output that should be sent back to realtime.
background_agent_output = handoff_output_rx.recv() => {
let pending_outbound = background_agent_output
.as_ref()
.ok()
.cloned()
.map(RealtimePendingOutbound::Handoff)
.map(Box::new);
handle_handoff_output(
background_agent_output,
&writer,
&handoff_state,
event_parser,
&mut response_create_queue,
)
.await
.map_err(|err| {
classify_realtime_input_error_with_pending(err, pending_outbound)
})
}
// Audio frames captured from the user microphone.
user_audio_frame = audio_rx.recv() => {
handle_user_audio_input(user_audio_frame, &writer)
.await
.map_err(classify_realtime_input_error)
}
};
if let Err(exit) = result {
break exit;
}
}这个循环没有 biased;,不能假定永远先处理服务端事件或永远优先音频。选中一个分支后,该分支仍要等待发送、事件处理或输出通道写入完成;select! 不等于每路业务都在独立任务中无限并行。
文本和 handoff 在从通道取出后先保存一份待发副本。若 socket 写入报 ApiError,这份副本能随 TransportLost 交给重连逻辑;音频没有对应的待发枚举项。这样的区别与“实时音频可以丢弃、离散文本需要尽量补送”的行为一致,但不是端到端交付保证。
远端错误也有层次:next_event() 返回 Err(ApiError) 是传输丢失;返回 Some(RealtimeEvent::Error) 是业务错误事件,先转发给消费者,再结束输入循环。后者不会因为字符串里出现网络相关词语就自动进入 sideband 重连。
6. 音频消费
6.1 字节与时长
业务音频帧是协议对象,除了 Base64 数据还带采样信息及可选 item 身份:
源码文件:codex-rs/protocol/src/protocol.rs
相关函数/类型:RealtimeAudioFrame(源码定位,节选)
#[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>,
}samples_per_channel 表示每声道样本数量,不能当成所有声道样本数。item_id 可把输出音频与特定实时 item 关联。管理器本身不在 audio_in 里校验编码、重采样或枚举设备。
向实时服务发送时,writer 根据 parser 选择消息类型,实际取用的是 frame.data:
源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs
相关函数/类型:RealtimeWebsocketWriter::send_audio_frame(源码定位,节选)
pub async fn send_audio_frame(&self, frame: RealtimeAudioFrame) -> Result<(), ApiError> {
let message = match self.event_parser {
RealtimeEventParser::V1 | RealtimeEventParser::RealtimeV2 => {
RealtimeOutboundMessage::InputAudioBufferAppend { audio: frame.data }
}
RealtimeEventParser::FramelessBidi => {
RealtimeOutboundMessage::InputAudioAppend { audio: frame.data }
}
};
self.send_json(&message).await
}这说明不能把修改 sample_rate 字段理解成 Core 自动完成了采样率转换。外部生产者必须提供符合所建实时会话要求的音频;结构体中的描述信息并不会全部随此消息发给服务端。输入队列容量也不能代替格式验证。
V2 的输出打断逻辑需要估算音频长度,计算发生在接收侧:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:audio_duration_ms、decoded_samples_per_channel(源码定位,节选)
fn audio_duration_ms(frame: &RealtimeAudioFrame) -> u32 {
let Some(samples_per_channel) = frame
.samples_per_channel
.or(decoded_samples_per_channel(frame))
else {
return 0;
};
let sample_rate = u64::from(frame.sample_rate.max(1));
((u64::from(samples_per_channel) * 1_000) / sample_rate) as u32
}
fn decoded_samples_per_channel(frame: &RealtimeAudioFrame) -> Option<u32> {
let bytes = BASE64_STANDARD.decode(&frame.data).ok()?;
let channels = usize::from(frame.num_channels.max(1));
let samples = bytes.len().checked_div(2)?.checked_div(channels)?;
u32::try_from(samples).ok()
}优先使用显式的每声道样本数,否则解码 Base64 后按每样本 2 字节、声道数求出样本数。以 24 kHz、每声道 480 个样本为例,得到 20 ms。这里的 max(1) 只是避免除零,不是完整的参数合法性检查;无效 Base64 且没有显式样本数时,时长为零。
update_output_audio_state 对同一 item_id 累加每帧时长,换 item 时替换状态,没有 item 身份或时长为零时不更新。由于每帧独立做整数除法,得到的是按已接收数据累计的时长估计,不是外部音频设备回报的已播放位置。
6.2 打断与通知
当用户开始说话,内部 V2 分支从当前输出状态取出 item,并在事件不指定 item 或指定的是同一 item 时发送截断请求:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:handle_realtime_server_event 的语音打断分支(源码定位,节选)
RealtimeEvent::InputAudioSpeechStarted(event) => {
match session_kind {
RealtimeSessionKind::V1 => {}
RealtimeSessionKind::V2 => {
if let Some(output_audio_state) = output_audio_state.take()
&& event
.item_id
.as_deref()
.is_none_or(|item_id| item_id == output_audio_state.item_id)
{
writer
.send_payload(
json!({
"type": "conversation.item.truncate",
"item_id": output_audio_state.item_id,
"content_index": 0,
"audio_end_ms": output_audio_state.audio_end_ms,
})
.to_string(),
)
.await
.map_err(anyhow::Error::from)?;
}
}
}
false
}take() 会先清掉本地累计状态;即使事件带了另一个 item,旧状态也已取走,只是不发送 truncate。截断的 audio_end_ms 来自接收累计值,不能宣称精确反映用户耳朵已经听到的位置。V1 与 V3 不执行这段 V2 截断逻辑;普通 Turn 的取消也不是这段语音打断代码的职责。
向客户端交付音频则走另一条路径:handle_realtime_server_event 将事件写入 events_tx,fanout 封装成 EventMsg::RealtimeConversationRealtime,App Server 再转换为公开通知。
源码文件:codex-rs/app-server/src/bespoke_event_handling.rs
相关函数/类型:apply_bespoke_event_handling(源码定位,节选)
RealtimeEvent::AudioOut(audio) => {
let notification = ThreadRealtimeOutputAudioDeltaNotification {
thread_id: conversation_id.to_string(),
audio: audio.into(),
};
outgoing
.send_server_notification(ServerNotification::ThreadRealtimeOutputAudioDelta(
notification,
))
.await;
}消费者得到 thread_id 和音频帧,负责播放或进一步转发。Started、SDP、转录、错误和关闭通知也在 apply_bespoke_event_handling 中分流;不要期待一个底层 RealtimeEvent 在客户端一定保持原枚举形态。Core 没有在这一链路确认扬声器播放完成。
7. 工作交接
7.1 委派输入
fanout 的作用不只是把事件复制到外部。它专门查看 HandoffRequested,从 handoff 的输入与转录构造委派文本:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:handle_start_inner 的 fanout task(源码定位,节选)
// 先提交委派,再向外转发该 RealtimeEvent;这里不等待整个 Turn 执行结束。
let maybe_routed_text = match &event {
RealtimeEvent::HandoffRequested(handoff) => {
realtime_delegation_from_handoff(handoff)
}
_ => None,
};
if let Some(text) = maybe_routed_text {
debug!(text = %text, "[realtime-text] realtime conversation text output");
let sess_for_routed_text = Arc::clone(&sess_clone);
sess_for_routed_text.route_realtime_text_input(text).await;
}
sess_clone
.send_event_raw(ev(EventMsg::RealtimeConversationRealtime(
RealtimeConversationRealtimeEvent {
payload: event.clone(),
},
)))
.await;只有得到委派文本才调用路由方法。普通 ConversationItemAdded、音频与转录增量不会在这段代码中自动变成一次 Turn 输入。这能避免实时服务回显用户文本后又被 Core 当成新指令,形成反复委派。
路由进入的是同一 Session 的普通输入处理:
源码文件:codex-rs/core/src/session/turn_input.rs
相关函数/类型:Session::route_realtime_text_input(源码定位,节选)
pub(crate) async fn route_realtime_text_input(self: &Arc<Self>, text: String) {
let submission_id = Uuid::now_v7().to_string();
let submission = handle(
self,
TurnInputRequest::user_input(vec![UserInput::Text {
text,
text_elements: Vec::new(),
}]),
TurnInputMode::StartOrSteer,
submission_id.clone(),
)
.await;
match submission {
Ok(TurnInputSubmission::Started { .. } | TurnInputSubmission::Steered { .. }) => {}
Ok(TurnInputSubmission::NotSubmitted { reason }) => {
self.send_event_raw(Event {
id: submission_id,
msg: EventMsg::Error(ErrorEvent {
message: format!("failed to submit turn input: {reason:?}"),
codex_error_info: Some(CodexErrorInfo::BadRequest),
}),
})
.await;
}
Err(error) => {
self.send_event_raw(Event {
id: submission_id,
msg: EventMsg::Error(error.to_error_event(/*message_prefix*/ None)),
})
.await;
}
}
}StartOrSteer 在空闲时启动 Turn,在已有工作时尝试注入当前 Turn;它并不天然创建一个子 Agent 或新的 Thread。函数等待的是输入接受结果,成功形式为 Started 或 Steered,而不是等待模型推理、工具执行和最终回答全部完成。若不能提交,错误事件使用这次路由新生成的 submission ID。
下图把输入接纳与普通任务执行画成不同阶段,能解释为什么工具执行期间仍可转发音频。
图中的并发建立在队列与独立任务之上,仍受各次 await 的消费速度约束。Session::send_event 先向外交付普通事件,再执行回答镜像;实时输出队列长期拥塞时,镜像等待仍可能反过来拖慢普通事件发送路径,不能把它理解成完全无背压的旁路。
inbound_handoff_request_starts_turn 使用两台 mock 服务:实时服务发出转录和 handoff,普通 Responses 服务返回回答。关键断言同时查看新 Turn 身份与真正发出的普通模型请求:
源码文件:codex-rs/core/tests/suite/realtime_conversation.rs
相关函数/类型:inbound_handoff_request_starts_turn(源码定位,节选)
let turn_id = loop {
let event = test.codex.next_event().await?;
if let EventMsg::TurnStarted(turn_started) = event.msg {
break turn_started.turn_id;
}
};
Uuid::parse_str(&turn_id).context("realtime-routed turn ID should be a UUID")?;
wait_for_event(&test.codex, |event| {
matches!(event, EventMsg::TurnComplete(_))
})
.await;
let request = response_mock.single_request();
let user_texts = request.message_input_texts("user");
assert!(user_texts.iter().any(|text| text
== "<realtime_delegation>\n <input>text from realtime</input>\n <transcript_delta>user: text from realtime</transcript_delta>\n</realtime_delegation>"));这不仅断言某个函数被调用,还检查委派文本实际进入模型请求。边界是 mock 的请求链路,不是实时模型理解自然语音的准确率。
另一个测试 inbound_conversation_item_does_not_start_turn_and_still_forwards_audio 发送 user role 的文本回显,紧接着发送音频:
源码文件:codex-rs/core/tests/suite/realtime_conversation.rs
相关函数/类型:inbound_conversation_item_does_not_start_turn_and_still_forwards_audio(源码定位,节选)
let audio_out = tokio::time::timeout(
Duration::from_millis(500),
wait_for_event_match(&test.codex, |msg| match msg {
EventMsg::RealtimeConversationRealtime(RealtimeConversationRealtimeEvent {
payload: RealtimeEvent::AudioOut(frame),
}) => Some(frame.clone()),
_ => None,
}),
)
.await
.expect("timed out waiting for realtime audio after conversation item");
assert_eq!(audio_out.data, "AQID");
let unexpected_turn_started = tokio::time::timeout(
Duration::from_millis(200),
wait_for_event_match(&test.codex, |msg| match msg {
EventMsg::TurnStarted(_) => Some(()),
_ => None,
}),
)
.await;
assert!(unexpected_turn_started.is_err());它先要求音频在 500 ms 内到达,再要求 200 ms 观察窗口内没有 TurnStarted。因此这是“回显不委派、音频仍可继续”的反例验证,不能只写成“实时会话测试通过”。超时窗口只用于本地测试判定,不是产品响应时限。
7.2 回答回灌
返回方向从普通 Session::send_event 进入。它保留普通事件的副本,调用 maybe_mirror_event_text_to_realtime:
源码文件:codex-rs/core/src/session/mod.rs
相关函数/类型:maybe_mirror_event_text_to_realtime(源码定位,节选)
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
&& 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}");
}
}ItemStarted 为回答 item 登记流状态;AgentMessageContentDelta 追加文本;ItemCompleted 尝试完成已登记的流,若已处理就不再发送一份完整回答。不能把每个 item 的完成与整个 Turn 的完成混为一谈,一个 Turn 可以产生多个回答 item 和工具步骤。
核心输出对象 RealtimeHandoffState 持有 output_tx、共享的 last_output 和流状态。流状态用 active_handoff 标识当前交接,用 items 按回答 item ID 保存增量。handoff_out 先看 client_managed_handoffs:为真就跳过自动镜像;codex_responses_as_items 则将自动回答包装成会话 item。这两项都不是“关闭整个实时会话”。
默认路径中,有 active handoff 时回答关联该 ID;没有 active handoff 时仍可产生 standalone 输出。API 消费者 handle_handoff_output 再按 parser 分流,V1、V2、V3 不能复用同一份 wire 消息解释。V3 的逐 item 文本还经过 200 ms 刷新间隔与大小限制;这里的时间控制用于回答回灌,不是音频队列的采样周期。
完整 Turn 结束时才做 handoff 的收尾:
源码文件:codex-rs/core/src/session/mod.rs
相关函数/类型:maybe_clear_realtime_handoff_for_event(源码定位,节选)
async fn maybe_clear_realtime_handoff_for_event(&self, msg: &EventMsg) {
if !matches!(msg, EventMsg::TurnComplete(_)) {
return;
}
if let Err(err) = self.conversation.handoff_complete().await {
debug!("failed to finalize realtime handoff output: {err}");
}
self.conversation.clear_active_handoff().await;
}handoff_complete 对内部 V1 直接返回,因此也覆盖 V3;内部 V2 在拥有 active ID 和 last_output 时才生成完成消息。普通 V2 默认路径先返回 function call output 的确认文本,再按需要请求下一次实时 response;codex_responses_as_items 路径使用独立 ack。随后 clear_active_handoff 清空 active ID、item 映射和最后输出,不能提前挪到任意 ItemCompleted 上。
7.3 响应合并
V2 的“请求实时服务生成响应”与普通 Turn 完成是两套状态。RealtimeResponseCreateQueue 用两个布尔值合并请求:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:RealtimeResponseCreateQueue(源码定位,节选)
// pending_create 是一个布尔值,多次请求被合并为一次待发 response.create。
#[derive(Default)]
struct RealtimeResponseCreateQueue {
active_default_response: bool,
pending_create: bool,
}
impl RealtimeResponseCreateQueue {
async fn request_create(
&mut self,
writer: &RealtimeWebsocketWriter,
reason: &str,
) -> anyhow::Result<()> {
if self.active_default_response {
self.pending_create = true;
return Ok(());
}
self.send_create_now(writer, reason).await
}
fn mark_started(&mut self) {
self.active_default_response = true;
}
async fn mark_finished(
&mut self,
writer: &RealtimeWebsocketWriter,
reason: &str,
) -> anyhow::Result<()> {
self.active_default_response = false;
if !self.pending_create {
return Ok(());
}
self.pending_create = false;
self.send_create_now(writer, reason).await
}
async fn send_create_now(
&mut self,
writer: &RealtimeWebsocketWriter,
reason: &str,
) -> anyhow::Result<()> {
if let Err(err) = writer.send_response_create().await {
if matches!(&err, ApiError::Stream(message) if message.starts_with(REALTIME_ACTIVE_RESPONSE_ERROR_PREFIX))
{
warn!("realtime response.create raced an active response; deferring");
self.active_default_response = true;
self.pending_create = true;
return Ok(());
}
warn!("failed to send {reason} response.create: {err}");
return Err(err.into());
}
self.active_default_response = true;
Ok(())
}
}已有默认 response 运行时,后续请求只将 pending_create 设为真;一次完成或取消后最多补发一次。它不是存储 N 个响应请求的 FIFO 队列。send_create_now 还识别“already has an active response”错误前缀,把竞争结果重新变成 active 加 pending;其他发送错误才返回失败。
对应的 ResponseCreated、ResponseDone 与 ResponseCancelled 来自实时服务,TurnComplete 来自普通 Core。把它们当成一个生命周期会造成两种常见误改:普通 item 一结束就清掉 handoff,或每个回答增量都发送 response.create。前者丢失关联,后者与正在运行的实时 response 竞争。
8. 停止与替换
8.1 取走状态
第二次 Start 不是追加一套并行实时运行时。完成参数预检后,管理器取出旧状态并等待旧任务收尾,再创建新状态:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:RealtimeConversationManager::start(源码定位,节选)
async fn start(
&self,
start: RealtimeStart,
mode_instructions: RealtimeModeInstructions,
) -> CodexResult<RealtimeStartOutput> {
let previous_state = {
let mut guard = self.state.lock().await;
guard.take()
};
if let Some(state) = previous_state {
stop_conversation_state(state, RealtimeFanoutTaskStop::Await).await;
}
let output = self.start_inner(start).await?;
*self.mode_instructions.lock().await = Some(mode_instructions);
Ok(output)
}这段顺序有两个不同的失败边界。准备参数失败发生在进入 manager 前,旧会话尚未被取走;通过预检、停止旧会话后,新连接又失败,则不会自动恢复旧会话。Mutex<Option<_>> 也不是覆盖整个异步 start 的互斥锁;单一当前运行时还依赖正常入口中顺序处理 Start 的调用方式。
主动关闭复用了相同的资源停止函数:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:shutdown、stop_conversation_state(源码定位,节选)
pub(crate) async fn shutdown(&self) -> CodexResult<()> {
let state = {
let mut guard = self.state.lock().await;
guard.take()
};
if let Some(state) = state {
stop_conversation_state(state, RealtimeFanoutTaskStop::Await).await;
}
Ok(())
}
// ...
async fn stop_conversation_state(
mut state: ConversationState,
fanout_task_stop: RealtimeFanoutTaskStop,
) {
state.realtime_active.store(false, Ordering::Relaxed);
state.stop_token.cancel();
let _ = state.input_task.await;
if let Some(fanout_task) = state.fanout_task.take() {
match fanout_task_stop {
RealtimeFanoutTaskStop::Await => {
let _ = fanout_task.await;
}
RealtimeFanoutTaskStop::Detach => {}
}
}
}先在短锁内 take(),后在锁外等待。取走状态阻止新调用继续拿到旧 sender;realtime_active = false 阻止旧运行时继续以活跃身份收尾;stop_token.cancel() 让输入循环和 sideband 的连接等待观察到取消。只有 input task 结束,相关 sender 和 receiver 才按所有权释放,fanout 才能看到通道关闭并结束排空。
这不是对所有任务直接 abort(),也没有固定的关闭超时。如果某个已经选中的处理分支正在等待外部发送或另一个消费者,取消要等它回到相应的检查点。源码能说明协作取消的顺序,不能保证任意网络与消费条件下 Close 都立即返回。
8.2 会话身份
网络任务异步结束时,不能只看“现在有一个 state”便将其删除:它可能属于新一轮 Start。代码把 Arc<AtomicBool> 的分配身份也作为运行时身份:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:register_fanout_task、finish_if_active(源码定位,节选)
pub(crate) async fn register_fanout_task(
&self,
realtime_active: &Arc<AtomicBool>,
fanout_task: JoinHandle<()>,
) {
let mut fanout_task = Some(fanout_task);
{
let mut guard = self.state.lock().await;
if let Some(state) = guard.as_mut()
&& Arc::ptr_eq(&state.realtime_active, realtime_active)
{
state.fanout_task = fanout_task.take();
}
}
if let Some(fanout_task) = fanout_task {
fanout_task.abort();
let _ = fanout_task.await;
}
}
pub(crate) async fn finish_if_active(&self, realtime_active: &Arc<AtomicBool>) {
let state = {
let mut guard = self.state.lock().await;
match guard.as_ref() {
Some(state) if Arc::ptr_eq(&state.realtime_active, realtime_active) => guard.take(),
_ => None,
}
};
if let Some(state) = state {
stop_conversation_state(state, RealtimeFanoutTaskStop::Detach).await;
}
}两个原子布尔值即使都为 true,也可能属于两次不同启动;Arc::ptr_eq 比较共享对象身份,阻止旧 fanout 注册或清理到新状态上。注册失败时,那个已不属于当前运行时的 fanout 会被 abort 并等待退出。
finish_if_active 使用 Detach 还有一个并发原因:调用者就是正在结束的 fanout;若它再等待自己的 JoinHandle,便形成自等待。这里 detach 的是当前收尾任务的句柄,不是宣称所有实时资源都交给后台不管;input task 仍被等待。
fanout 本身要先排空解析好的事件,再处理结束转录:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:handle_start_inner 的 fanout 收尾(源码定位,节选)
if let Ok(text) = transcript_tail_rx.recv().await {
sess_clone.route_realtime_text_input(text).await;
}
if fanout_realtime_active.swap(false, Ordering::Relaxed) {
match end {
RealtimeConversationEnd::TransportClosed => {
info!("realtime conversation transport closed");
}
RealtimeConversationEnd::Requested | RealtimeConversationEnd::Error => {}
}
sess_clone
.conversation
.finish_if_active(&fanout_realtime_active)
.await;
send_realtime_conversation_closed(&sess_clone, sub_id, end).await;
}events_rx 的循环在此前已经结束,剩余 handoff 因此先于 transcript tail 进入普通输入路径。flush_transcript_tail_on_session_end 默认关闭;开启后,结束对话仍可能把尾部转录提交给普通 Turn。停止实时会话不等于取消普通 Turn,更不等于关闭整个 Session。
主动停止先把 active 设为 false,fanout 的 swap(false) 就不会再次发送同一轮异步关闭通知;主动 Close 的调用方发送 reason = requested。自然断开和错误则由仍拥有 active 身份的 fanout 发出 transport_closed 或 error。这防止的是相互竞争的结束路径重复通知,不意味着重复调用 Close 的请求永远只有一个响应事件。
8.3 Socket 回收
Core 的 shutdown 没有调用 RealtimeWebsocketWriter::close。普通输入任务退出并释放最后一个 writer 后,下层 WsStream 的析构负责停止 pump:
源码文件:codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs
相关函数/类型:Drop for WsStream(源码定位,节选)
impl Drop for WsStream {
fn drop(&mut self) {
self.pump_task.abort();
}
}这与直接调用 API 连接的 close() 不同:后者对 V3 会尝试发送 session.close,再发送 WebSocket Close。不能因为这个 API 方法存在,就描述成所有 Core 关闭路径都进行了协议级优雅关闭。对 WebRTC sideband,释放控制 socket 也不能证明客户端的媒体 call 已挂断。
下表描述的是由字段和任务派生出的阶段,不是另一套源码枚举:
| 阶段 | Manager 槽位 | 关键动作 | 对外结果 |
|---|---|---|---|
| 尚未启动 | None | 输入不能取得 sender | Audio/Text 报未运行 |
| 运行中 | Some 且 active | 输入任务、fanout 和 API pump 工作 | 音频、转录、handoff 通知 |
| 主动停止 | 先取为 None | inactive、cancel、等待 input 与 fanout | requested |
| 自然终止 | fanout 持有该轮身份 | 排空、身份比较、取走匹配 state | transport_closed 或 error |
| 新启动替换 | 先取走旧 state | 等待旧组结束,再构造新组 | 新一轮 Started;建连失败不回滚 |
conversation_second_start_replaces_runtime 先建立旧连接,再启动不同会话 ID 的新连接,最后发一帧音频。断言要求旧连接只收到一次初始化,新连接收到初始化和音频,且两次握手使用各自的 session ID;它验证的是路由切换,不能只数一共有几个连接。
conversation_webrtc_close_while_sideband_connecting_drops_pending_join 则让 mock 服务接收握手但不回复,然后提交 Close。它要求收到 requested,取消的握手连接最终关闭,后续 Session 退出前不再出现旧任务的错误或重复 close。这个测试正好覆盖第 3 节的后台握手边界。
9. 断线恢复
9.1 退出分类
并非所有退出都可重试。input task 的结果只有正常终止与传输丢失两个方向,后者还可带一条已经从队列取出、却未成功发送完的输出:
源码文件:codex-rs/core/src/realtime_conversation.rs
相关函数/类型:RealtimeInputTaskExit、classify_realtime_input_error_with_pending(源码定位,节选)
enum RealtimeInputTaskExit {
Terminal,
TransportLost {
err: ApiError,
pending_outbound: Option<Box<RealtimePendingOutbound>>,
},
}
#[derive(Clone, Debug, PartialEq, Eq)]
enum RealtimePendingOutbound {
Text(ConversationTextParams),
Handoff(RealtimeOutbound),
}
// ...
fn classify_realtime_input_error_with_pending(
err: anyhow::Error,
pending_outbound: Option<Box<RealtimePendingOutbound>>,
) -> RealtimeInputTaskExit {
match err.downcast::<ApiError>() {
Ok(err) => RealtimeInputTaskExit::TransportLost {
err,
pending_outbound,
},
Err(err) => {
warn!("realtime input task stopped: {err}");
RealtimeInputTaskExit::Terminal
}
}
}分类使用 anyhow::Error::downcast::<ApiError>(),依据的是错误类型。输入 channel 关闭、消费者丢失以及已转发的协议错误会终止;socket/API 错误可成为 TransportLost。普通新 WebSocket 的外层 spawn_realtime_input_task 遇到 TransportLost 只报告错误并结束,它没有包上一层自动恢复循环。
V3 的底层 next_event 还区分 WebSocket 终止方式:正常 Close code 返回 None;非正常 Close、无状态码关闭、接收流意外结束返回传输错误。V1/V2 的 Close 处理不同,不能直接沿用 V3 的异常重连结论。
9.2 两层重试
恢复循环位于 codex-rs/core/src/realtime_conversation/sideband.rs,只包裹 WebRTC sideband 路径。它拿到同一 call 的连接后运行 input task,再根据退出类型处理:
源码文件:codex-rs/core/src/realtime_conversation/sideband.rs
相关函数/类型:spawn_webrtc_sideband_input_task(源码定位,节选)
// 只在匹配 TransportLost 后重连;正常终止和协议 Error 不进入这个分支。
let RealtimeInputTaskExit::TransportLost {
err,
pending_outbound: failed_outbound,
} = exit
else {
break;
};
pending_outbound = failed_outbound;
if !realtime_active.load(Ordering::Relaxed) || stop_token.is_cancelled() {
break;
}
if event_parser != RealtimeEventParser::FramelessBidi {
report_realtime_transport_loss(&events_tx, err).await;
break;
}
if connected_at.elapsed() >= STABLE_CONNECTION_DURATION {
rapid_disconnects = 0;
}
rapid_disconnects = rapid_disconnects.saturating_add(1);
let delay = reconnect_delay(rapid_disconnects);
warn!(
call_id,
delay_ms = delay.as_millis(),
"live sideband transport lost; reconnecting: {err}"
);
reconnecting = true;
tokio::select! {
biased;
_ = stop_token.cancelled() => break,
_ = sleep(delay) => {}
}恢复必须同时满足:仍活跃、没有取消、退出为 TransportLost、parser 为 FramelessBidi。这也适用于显式 V3 的 existingCall,因为它加入同一 sideband 外层任务。成功重连沿用 call ID、输入 channels、handoff state 和 transcript state,并在读取下一条新输出前先尝试补发 pending 输出。
运行期退避与握手期重试使用不同计数:
源码文件:codex-rs/core/src/realtime_conversation/sideband.rs
相关函数/类型:reconnect_delay、webrtc_sideband_session_ended(源码定位,节选)
const RECONNECT_BASE_DELAY: Duration = Duration::from_millis(200);
const RECONNECT_MAX_DELAY: Duration = Duration::from_secs(5);
const STABLE_CONNECTION_DURATION: Duration = Duration::from_secs(30);
// ...
fn webrtc_sideband_session_ended(err: &ApiError) -> bool {
matches!(
err,
ApiError::Api { status, .. }
if matches!(*status, StatusCode::NOT_FOUND | StatusCode::GONE)
)
}
fn reconnect_delay(rapid_disconnects: u32) -> Duration {
let multiplier = 2_u32.saturating_pow(rapid_disconnects.saturating_sub(1));
RECONNECT_BASE_DELAY
.saturating_mul(multiplier)
.min(RECONNECT_MAX_DELAY)
}连续短连接的恢复等待从 200 ms、400 ms、800 ms 增长,最大 5 s;一次连接持续至少 30 s 后,短连接计数先重置。这里没有固定的“最多恢复三次”。另一方面,每次实际加入 sideband 又经过 API 的 connect_sideband,它使用 provider.retry.max_attempts 和 base_delay 做有限次握手尝试。
404/410 被视为 call 已结束,API 的握手层不重试它们。外层已经处在 reconnecting 时收到这两类状态,会正常结束而不增加错误事件;首次加入就失败时则走实时错误路径。其他错误在握手层耗尽尝试后也会报告错误并结束,并不是无限循环所有失败。
下图从退出结果出发,把普通连接、sideband、协议错误与握手状态分开。
这张图中的“报告错误或取消”还需看具体原因:主动取消不会伪造一条网络错误。sideband 等待连接和退避时都把 cancellation token 纳入 select!,并使用 biased; 优先检查取消;这是恢复期间仍能退出的关键。
9.3 补送边界
补送不是“消息恰好执行一次”。一条写入失败可能发生在服务端已经收到字节、客户端尚未获知结果的时刻;重试复用文本或 handoff 对象,并没有在此层提供服务端去重确认。待发值也只有一条,不是持久化 outbox。音频不在 pending 枚举里,进程退出后这些状态不会从本地日志恢复。
单元测试直接用两类待发输出验证类型和对象保存:
源码文件:codex-rs/core/src/realtime_conversation_tests.rs
相关函数/类型:classifies_outbound_api_failures_as_transport_loss(源码定位,节选)
fn classifies_outbound_api_failures_as_transport_loss() {
for pending_outbound in [
RealtimePendingOutbound::Text(ConversationTextParams {
text: "retry me".to_string(),
role: ConversationTextRole::User,
}),
RealtimePendingOutbound::Handoff(RealtimeOutbound::StandaloneHandoff {
text: "retry this handoff".to_string(),
phase: Some(MessagePhase::FinalAnswer),
}),
] {
let exit = classify_realtime_input_error_with_pending(
ApiError::Stream("failed to send realtime request".to_string()).into(),
Some(Box::new(pending_outbound.clone())),
);
let RealtimeInputTaskExit::TransportLost {
err: ApiError::Stream(_),
pending_outbound: Some(actual_pending_outbound),
} = exit
else {
panic!("outbound API failure should preserve pending output for reconnect");
};
assert_eq!(*actual_pending_outbound, pending_outbound);
}
assert!(matches!(
classify_realtime_input_error(anyhow::anyhow!("input channel closed")),
RealtimeInputTaskExit::Terminal
));
}断言要求 ApiError::Stream 保留原来的 text 或 handoff,而普通“input channel closed”错误归为 Terminal。这能验证错误分类和重试输入不被丢弃,不能证明服务端在失败前没有执行过请求。
更完整的测试 conversation_webrtc_live_reconnects_sideband_after_unclean_disconnect 同时参数化新 WebRTC call 与 existingCall:第一条连接发送未完成的 hello wor 转录,持续推送入站事件,要求出站文本仍能发出,然后异常断开。第二条连接补齐 hello world、发送回答转录与委派;关键断言检查恢复后的 handoff:
源码文件:codex-rs/core/tests/suite/realtime_conversation.rs
相关函数/类型:conversation_webrtc_live_reconnects_sideband_after_unclean_disconnect(源码定位,节选)
assert_eq!(
handoff,
RealtimeHandoffRequested {
handoff_id: "handoff_reconnect".to_string(),
item_id: "handoff_reconnect".to_string(),
input_transcript: "hello world".to_string(),
active_transcript: vec![
RealtimeTranscriptEntry {
role: "user".to_string(),
text: "hello world".to_string(),
},
RealtimeTranscriptEntry {
role: "assistant".to_string(),
text: "after reconnect".to_string(),
}
],
}
);同一个 transcript state 必须把完成文本与已有增量协调起来,不能生成重复的用户条目。测试还让第三次握手返回 410,要求随后没有新的连接尝试,最终关闭原因是 transport_closed。它验证的是控制通道恢复和转录连续性,不证明客户端 WebRTC 媒体在断网期间保持无缝播放。
10. 故障定位
源码阅读可以从具体症状反推到边界,避免把所有失败都归因于“实时模型没响应”。
| 现象 | 先检查哪里 | 能区分的问题 |
|---|---|---|
| RPC 立即拒绝 | dispatch_initialized_client_request、prepare_realtime_conversation_thread | 初始化、实验 API 能力、Thread feature |
| Start 返回后只收到 realtime/error | handle_start、prepare_realtime_start、连接函数 | 认证或预检失败与运行时建连失败 |
| 没有 Started,却一直等 closed | conversation_start_*_failure_emits_realtime_error_only | 前置/初次建连错误只发 realtime error,不承诺 closed |
| 已有 SDP,没有后续事件 | spawn_webrtc_sideband_input_task | call 创建成功与 sideband 加入成功不同 |
| 音频调用成功但间歇缺失 | audio_in 的 Full 分支 | 主动丢弃新帧与连接关闭不同 |
| 收到音频通知却听不到 | apply_bespoke_event_handling 后的客户端消费者 | Core 交付与设备播放不同 |
| 转录可见但没有普通 Turn | fanout 的 HandoffRequested 分支 | 展示事件与任务委派不同 |
| 一个回答 item 结束后交接错乱 | maybe_clear_realtime_handoff_for_event | ItemCompleted 与 TurnComplete 不同 |
| Close 卡住或再次 Start 很慢 | stop_conversation_state 及当前分支的 await | 协作取消、事件排空与消费者阻塞 |
| 普通 WebSocket 断线没有自动恢复 | spawn_realtime_input_task 与 sideband.rs | 是否进入 V3 sideband 恢复循环 |
两项启动失败测试提供了特别有用的反例。conversation_start_preflight_failure_emits_realtime_error_only 使用缺少 API key 的认证条件,要求实时错误里包含 key 要求;conversation_start_connect_failure_emits_realtime_error_only 让连接失败。二者都在之后的观察窗口断言没有 RealtimeConversationClosed。如果客户端总以 closed 作为启动失败后的唯一完成信号,就会在这些路径上等错事件。
在源码仓库的 codex-rs 目录,可先运行下面三组测试,再按失败点缩小到一个用例:
just test --locked -p codex-core --lib -E 'test(realtime_conversation::)'
just test --locked -p codex-core --test all -E 'test(suite::realtime_conversation::)'
just test --locked -p codex-api --test realtime_websocket_e2e第一组覆盖转录包装、输出分类、BEM 边界及退避函数,第二组从 Core 操作走到真实本地 mock socket,第三组覆盖 API 收发并行、关闭和加入重试。Nextest 中因过滤而 skipped 的其他测试,不代表被选中用例失败;同时要检查 skip_if_no_network!,网络禁用环境会让相关测试提前返回,不能把这种返回当作已验证 socket 行为。
入口和消费者还可用 App Server 用例核对:
just test --locked -p codex-app-server --test all -E \
'test(realtime_conversation_requires_feature_flag) | test(existing_call_attaches_without_reinitializing) | test(websocket_v2_forwards_audio_and_text) | test(websocket_v2_tool_call_does_not_block_sideband_audio)'这些用例分别检验 feature 拒绝、附着已有 call 不覆盖其配置、音频与文本在公开 API 两端的转换,以及普通工具调用期间仍能收到实时音频。它们都使用受控服务和输入,不能代替真实媒体设备、网络抖动与服务能力验证。本篇引用的 Core 管理器与 API 收发实现没有 macOS 专属的条件编译;设备和媒体能力仍要由实际客户端确认,不能由这些控制通道测试推导。
现在给定这样一个过程:V3 existingCall 已开始,用户发出文本,控制连接异常断开,后台 Turn 仍在执行,客户端紧接着发 Stop。沿源码应能指出:pending 文本由谁保存、普通 Turn 是否被取消、哪一个 token 中止重连等待、为什么旧 fanout 不能删除新会话,以及 transcript tail 开关可能增加哪次普通输入。每个答案都应落到一个具体字段、条件或消费者上。
若要继续追踪普通任务的停止边界,可以读 Turn中断与运行中注入;整个 Core Session 如何退出,见 Session关闭流程。两处处理的资源范围都比单独停止实时连接更大,排障时先选对生命周期,再选择取消或重启操作。
