Skip to content

实时会话架构

从实验 API 追踪实时连接、音频与文本队列、普通 Turn 交接和事件消费,解释会话替换、协作取消与 V3 sideband 断线恢复的所有权边界。

基于rust-v0.150.0
CodexRealtimeRust并发

实时会话架构 ​

用户一边说话,一边让 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(源码定位,节选)

rust
// 先完成连接初始化,再检查客户端是否显式接受实验 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(源码定位,节选)

rust
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 的实时操作分支(源码定位,节选)

rust
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(源码定位,节选)

rust
#[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(源码定位,节选)

rust
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(源码定位,节选)

rust
// 内部 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 的默认模型
V1V1V1gpt-realtime-1.5
V2RealtimeV2V2gpt-realtime-1.5
V3FramelessBidiV1gpt-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 时已完成什么后续仍可能失败什么
新 WebSocketRealtimeWebsocketClient::connectsocket 建立和初始化发送;V3 还等待 session.started持续收发、服务端业务事件
新 WebRTC callModelClient::create_realtime_call_with_headerscall 创建、获得 answer、启动 sideband task后台 sideband 握手与客户端媒体连接
existingCallexisting_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(源码定位,节选)

rust
// 初始化策略由连接用途决定,不能只看协议版本。
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(源码定位,节选)

rust
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(源码定位,节选)

rust
// 状态锁保护的是所有权槽位;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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
// 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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
// 五个分支共享一个输入循环;每个分支被选中后仍会 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(源码定位,节选)

rust
#[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(源码定位,节选)

rust
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(源码定位,节选)

rust
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 的语音打断分支(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
// 先提交委派,再向外转发该 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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
// 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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
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 收尾(源码定位,节选)

rust
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(源码定位,节选)

rust
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输入不能取得 senderAudio/Text 报未运行
运行中Some 且 active输入任务、fanout 和 API pump 工作音频、转录、handoff 通知
主动停止先取为 Noneinactive、cancel、等待 input 与 fanoutrequested
自然终止fanout 持有该轮身份排空、身份比较、取走匹配 statetransport_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(源码定位,节选)

rust
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(源码定位,节选)

rust
// 只在匹配 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(源码定位,节选)

rust
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(源码定位,节选)

rust
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(源码定位,节选)

rust
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/errorhandle_start、prepare_realtime_start、连接函数认证或预检失败与运行时建连失败
没有 Started,却一直等 closedconversation_start_*_failure_emits_realtime_error_only前置/初次建连错误只发 realtime error,不承诺 closed
已有 SDP,没有后续事件spawn_webrtc_sideband_input_taskcall 创建成功与 sideband 加入成功不同
音频调用成功但间歇缺失audio_in 的 Full 分支主动丢弃新帧与连接关闭不同
收到音频通知却听不到apply_bespoke_event_handling 后的客户端消费者Core 交付与设备播放不同
转录可见但没有普通 Turnfanout 的 HandoffRequested 分支展示事件与任务委派不同
一个回答 item 结束后交接错乱maybe_clear_realtime_handoff_for_eventItemCompleted 与 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 目录,可先运行下面三组测试,再按失败点缩小到一个用例:

bash
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 用例核对:

bash
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关闭流程。两处处理的资源范围都比单独停止实时连接更大,排障时先选对生命周期,再选择取消或重启操作。