[stack 2/4] Align main realtime v2 wire and runtime flow (#14830)

## Stack Position
2/4. Built on top of #14828.

## Base
- #14828

## Unblocks
- #14829
- #14827

## Scope
- Port the realtime v2 wire parsing, session, app-server, and
conversation runtime behavior onto the split websocket-method base.
- Branch runtime behavior directly on the current realtime session kind
instead of parser-derived flow flags.
- Keep regression coverage in the existing e2e suites.

---------

Co-authored-by: Codex <noreply@openai.com>
This commit is contained in:
Ahmed Ibrahim
2026-03-16 21:38:07 -07:00
committed by GitHub
Unverified
parent 1d85fe79ed
commit fbd7f9b986
28 changed files with 807 additions and 140 deletions
+3
View File
@@ -2614,6 +2614,9 @@ impl Session {
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;
}
+1
View File
@@ -2735,6 +2735,7 @@ fn submission_dispatch_span_uses_debug_for_realtime_audio() {
sample_rate: 16_000,
num_channels: 1,
samples_per_channel: Some(160),
item_id: None,
},
}),
trace: None,
+315 -51
View File
@@ -11,6 +11,8 @@ use crate::realtime_context::build_realtime_startup_context;
use async_channel::Receiver;
use async_channel::Sender;
use async_channel::TrySendError;
use base64::Engine;
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use codex_api::Provider as ApiProvider;
use codex_api::RealtimeAudioFrame;
use codex_api::RealtimeEvent;
@@ -34,6 +36,8 @@ use codex_protocol::protocol::RealtimeHandoffRequested;
use http::HeaderMap;
use http::HeaderValue;
use http::header::AUTHORIZATION;
use serde_json::Value;
use serde_json::json;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::Ordering;
@@ -49,51 +53,72 @@ const USER_TEXT_IN_QUEUE_CAPACITY: usize = 64;
const HANDOFF_OUT_QUEUE_CAPACITY: usize = 64;
const OUTPUT_EVENTS_QUEUE_CAPACITY: usize = 256;
const REALTIME_STARTUP_CONTEXT_TOKEN_BUDGET: usize = 5_000;
const ACTIVE_RESPONSE_CONFLICT_ERROR_PREFIX: &str =
"Conversation already has an active response in progress:";
pub(crate) struct RealtimeConversationManager {
state: Mutex<Option<ConversationState>>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum RealtimeSessionKind {
V1,
V2,
}
#[derive(Clone, Debug)]
struct RealtimeHandoffState {
output_tx: Sender<HandoffOutput>,
active_handoff: Arc<Mutex<Option<String>>>,
last_output_text: Arc<Mutex<Option<String>>>,
session_kind: RealtimeSessionKind,
}
#[derive(Debug, PartialEq, Eq)]
struct HandoffOutput {
handoff_id: String,
output_text: String,
enum HandoffOutput {
ImmediateAppend {
handoff_id: String,
output_text: String,
},
FinalToolCall {
handoff_id: String,
output_text: String,
},
}
#[derive(Debug, PartialEq, Eq)]
struct OutputAudioState {
item_id: String,
audio_end_ms: u32,
}
struct RealtimeInputTask {
writer: RealtimeWebsocketWriter,
events: RealtimeWebsocketEvents,
user_text_rx: Receiver<String>,
handoff_output_rx: Receiver<HandoffOutput>,
audio_rx: Receiver<RealtimeAudioFrame>,
events_tx: Sender<RealtimeEvent>,
handoff_state: RealtimeHandoffState,
session_kind: RealtimeSessionKind,
}
impl RealtimeHandoffState {
fn new(output_tx: Sender<HandoffOutput>) -> Self {
fn new(output_tx: Sender<HandoffOutput>, session_kind: RealtimeSessionKind) -> Self {
Self {
output_tx,
active_handoff: Arc::new(Mutex::new(None)),
last_output_text: Arc::new(Mutex::new(None)),
session_kind,
}
}
async fn send_output(&self, output_text: String) -> CodexResult<()> {
let Some(handoff_id) = self.active_handoff.lock().await.clone() else {
return Ok(());
};
self.output_tx
.send(HandoffOutput {
handoff_id,
output_text,
})
.await
.map_err(|_| CodexErr::InvalidRequest("conversation is not running".to_string()))?;
Ok(())
}
}
#[allow(dead_code)]
struct ConversationState {
audio_tx: Sender<RealtimeAudioFrame>,
user_text_tx: Sender<String>,
writer: RealtimeWebsocketWriter,
handoff: RealtimeHandoffState,
task: JoinHandle<()>,
realtime_active: Arc<AtomicBool>,
@@ -129,6 +154,10 @@ impl RealtimeConversationManager {
state.task.abort();
let _ = state.task.await;
}
let session_kind = match session_config.event_parser {
RealtimeEventParser::V1 => RealtimeSessionKind::V1,
RealtimeEventParser::RealtimeV2 => RealtimeSessionKind::V2,
};
let client = RealtimeWebsocketClient::new(api_provider);
let connection = client
@@ -152,21 +181,23 @@ impl RealtimeConversationManager {
async_channel::bounded::<RealtimeEvent>(OUTPUT_EVENTS_QUEUE_CAPACITY);
let realtime_active = Arc::new(AtomicBool::new(true));
let handoff = RealtimeHandoffState::new(handoff_output_tx);
let task = spawn_realtime_input_task(
writer,
let handoff = RealtimeHandoffState::new(handoff_output_tx, session_kind);
let task = spawn_realtime_input_task(RealtimeInputTask {
writer: writer.clone(),
events,
user_text_rx,
handoff_output_rx,
audio_rx,
events_tx,
handoff.clone(),
);
handoff_state: handoff.clone(),
session_kind,
});
let mut guard = self.state.lock().await;
*guard = Some(ConversationState {
audio_tx,
user_text_tx,
writer,
handoff,
task,
realtime_active: Arc::clone(&realtime_active),
@@ -228,7 +259,51 @@ impl RealtimeConversationManager {
state.handoff.clone()
};
handoff.send_output(output_text).await
let Some(handoff_id) = handoff.active_handoff.lock().await.clone() else {
return Ok(());
};
*handoff.last_output_text.lock().await = Some(output_text.clone());
if matches!(handoff.session_kind, RealtimeSessionKind::V1) {
handoff
.output_tx
.send(HandoffOutput::ImmediateAppend {
handoff_id,
output_text,
})
.await
.map_err(|_| CodexErr::InvalidRequest("conversation is not running".to_string()))?;
}
Ok(())
}
pub(crate) async fn handoff_complete(&self) -> CodexResult<()> {
let handoff = {
let guard = self.state.lock().await;
guard.as_ref().map(|state| state.handoff.clone())
};
let Some(handoff) = handoff else {
return Ok(());
};
if matches!(handoff.session_kind, RealtimeSessionKind::V1) {
return Ok(());
}
let Some(handoff_id) = handoff.active_handoff.lock().await.clone() else {
return Ok(());
};
let Some(output_text) = handoff.last_output_text.lock().await.clone() else {
return Ok(());
};
handoff
.output_tx
.send(HandoffOutput::FinalToolCall {
handoff_id,
output_text,
})
.await
.map_err(|_| CodexErr::InvalidRequest("conversation is not running".to_string()))
}
pub(crate) async fn active_handoff_id(&self) -> Option<String> {
@@ -246,6 +321,7 @@ impl RealtimeConversationManager {
};
if let Some(handoff) = handoff {
*handoff.active_handoff.lock().await = None;
*handoff.last_output_text.lock().await = None;
}
}
@@ -467,7 +543,6 @@ pub(crate) async fn handle_text(
params: ConversationTextParams,
) {
debug!(text = %params.text, "[realtime-text] appending realtime conversation text input");
if let Err(err) = sess.conversation.text_in(params.text).await {
error!("failed to append realtime text: {err}");
send_conversation_error(sess, sub_id, err.to_string(), CodexErrorInfo::BadRequest).await;
@@ -491,16 +566,23 @@ pub(crate) async fn handle_close(sess: &Arc<Session>, sub_id: String) {
}
}
fn spawn_realtime_input_task(
writer: RealtimeWebsocketWriter,
events: RealtimeWebsocketEvents,
user_text_rx: Receiver<String>,
handoff_output_rx: Receiver<HandoffOutput>,
audio_rx: Receiver<RealtimeAudioFrame>,
events_tx: Sender<RealtimeEvent>,
handoff_state: RealtimeHandoffState,
) -> JoinHandle<()> {
fn spawn_realtime_input_task(input: RealtimeInputTask) -> JoinHandle<()> {
let RealtimeInputTask {
writer,
events,
user_text_rx,
handoff_output_rx,
audio_rx,
events_tx,
handoff_state,
session_kind,
} = input;
tokio::spawn(async move {
let mut pending_response_create = false;
let mut response_in_progress = false;
let mut output_audio_state: Option<OutputAudioState> = None;
loop {
tokio::select! {
text = user_text_rx.recv() => {
@@ -511,23 +593,66 @@ fn spawn_realtime_input_task(
warn!("failed to send input text: {mapped_error}");
break;
}
if matches!(session_kind, RealtimeSessionKind::V2) {
if response_in_progress {
pending_response_create = true;
} else if let Err(err) = writer.send_response_create().await {
let mapped_error = map_api_error(err);
warn!("failed to send text response.create: {mapped_error}");
break;
} else {
pending_response_create = false;
response_in_progress = true;
}
}
}
Err(_) => break,
}
}
handoff_output = handoff_output_rx.recv() => {
match handoff_output {
Ok(HandoffOutput {
handoff_id,
output_text,
}) => {
if let Err(err) = writer
.send_conversation_handoff_append(handoff_id, output_text)
.await
{
let mapped_error = map_api_error(err);
warn!("failed to send handoff output: {mapped_error}");
break;
Ok(handoff_output) => {
match handoff_output {
HandoffOutput::ImmediateAppend {
handoff_id,
output_text,
} => {
if let Err(err) = writer
.send_conversation_handoff_append(handoff_id, output_text)
.await
{
let mapped_error = map_api_error(err);
warn!("failed to send handoff output: {mapped_error}");
break;
}
}
HandoffOutput::FinalToolCall {
handoff_id,
output_text,
} => {
if let Err(err) = writer
.send_conversation_handoff_append(handoff_id, output_text)
.await
{
let mapped_error = map_api_error(err);
warn!("failed to send handoff output: {mapped_error}");
break;
}
if matches!(session_kind, RealtimeSessionKind::V2) {
if response_in_progress {
pending_response_create = true;
} else if let Err(err) = writer.send_response_create().await {
let mapped_error = map_api_error(err);
warn!(
"failed to send handoff response.create: {mapped_error}"
);
break;
} else {
pending_response_create = false;
response_in_progress = true;
}
}
}
}
}
Err(_) => break,
@@ -536,12 +661,108 @@ fn spawn_realtime_input_task(
event = events.next_event() => {
match event {
Ok(Some(event)) => {
if let RealtimeEvent::HandoffRequested(handoff) = &event {
*handoff_state.active_handoff.lock().await =
Some(handoff.handoff_id.clone());
let mut should_stop = false;
let mut forward_event = true;
match &event {
RealtimeEvent::ConversationItemAdded(item) => {
match item.get("type").and_then(Value::as_str) {
Some("response.created")
if matches!(session_kind, RealtimeSessionKind::V2) =>
{
response_in_progress = true;
}
Some("response.done")
if matches!(session_kind, RealtimeSessionKind::V2) =>
{
response_in_progress = false;
output_audio_state = None;
if pending_response_create {
if let Err(err) = writer.send_response_create().await {
let mapped_error = map_api_error(err);
warn!(
"failed to send deferred response.create: {mapped_error}"
);
break;
}
pending_response_create = false;
response_in_progress = true;
}
}
_ => {}
}
}
RealtimeEvent::AudioOut(frame) => {
if matches!(session_kind, RealtimeSessionKind::V2) {
update_output_audio_state(&mut output_audio_state, frame);
}
}
RealtimeEvent::InputAudioSpeechStarted(event) => {
if matches!(session_kind, RealtimeSessionKind::V2)
&& 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)
&& let Err(err) = 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
{
let mapped_error = map_api_error(err);
warn!("failed to truncate realtime audio: {mapped_error}");
}
}
RealtimeEvent::ResponseCancelled(_) => {
response_in_progress = false;
output_audio_state = None;
if matches!(session_kind, RealtimeSessionKind::V2)
&& pending_response_create
{
if let Err(err) = writer.send_response_create().await {
let mapped_error = map_api_error(err);
warn!(
"failed to send deferred response.create after cancellation: {mapped_error}"
);
break;
}
pending_response_create = false;
response_in_progress = true;
}
}
RealtimeEvent::HandoffRequested(handoff) => {
*handoff_state.active_handoff.lock().await =
Some(handoff.handoff_id.clone());
*handoff_state.last_output_text.lock().await = None;
response_in_progress = false;
output_audio_state = None;
}
RealtimeEvent::Error(message)
if matches!(session_kind, RealtimeSessionKind::V2)
&& message.starts_with(ACTIVE_RESPONSE_CONFLICT_ERROR_PREFIX) =>
{
warn!(
"realtime rejected response.create because a response is already in progress; deferring follow-up response.create"
);
pending_response_create = true;
response_in_progress = true;
forward_event = false;
}
RealtimeEvent::Error(_) => {
should_stop = true;
}
RealtimeEvent::SessionUpdated { .. }
| RealtimeEvent::InputTranscriptDelta(_)
| RealtimeEvent::OutputTranscriptDelta(_)
| RealtimeEvent::ConversationItemDone { .. } => {}
}
let should_stop = matches!(&event, RealtimeEvent::Error(_));
if events_tx.send(event).await.is_err() {
if forward_event && events_tx.send(event).await.is_err() {
break;
}
if should_stop {
@@ -588,6 +809,49 @@ fn spawn_realtime_input_task(
})
}
fn update_output_audio_state(
output_audio_state: &mut Option<OutputAudioState>,
frame: &RealtimeAudioFrame,
) {
let Some(item_id) = frame.item_id.clone() else {
return;
};
let audio_end_ms = audio_duration_ms(frame);
if audio_end_ms == 0 {
return;
}
if let Some(current) = output_audio_state.as_mut()
&& current.item_id == item_id
{
current.audio_end_ms = current.audio_end_ms.saturating_add(audio_end_ms);
return;
}
*output_audio_state = Some(OutputAudioState {
item_id,
audio_end_ms,
});
}
fn audio_duration_ms(frame: &RealtimeAudioFrame) -> u32 {
let Some(samples_per_channel) = frame
.samples_per_channel
.or_else(|| 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()
}
async fn send_conversation_error(
sess: &Arc<Session>,
sub_id: String,
@@ -1,5 +1,5 @@
use super::HandoffOutput;
use super::RealtimeHandoffState;
use super::RealtimeSessionKind;
use super::realtime_text_from_handoff_request;
use async_channel::bounded;
use codex_protocol::protocol::RealtimeHandoffRequested;
@@ -57,7 +57,7 @@ fn ignores_empty_handoff_request_input_transcript() {
#[tokio::test]
async fn clears_active_handoff_explicitly() {
let (tx, _rx) = bounded(1);
let state = RealtimeHandoffState::new(tx);
let state = RealtimeHandoffState::new(tx, RealtimeSessionKind::V1);
*state.active_handoff.lock().await = Some("handoff_1".to_string());
assert_eq!(
@@ -68,47 +68,3 @@ async fn clears_active_handoff_explicitly() {
*state.active_handoff.lock().await = None;
assert_eq!(state.active_handoff.lock().await.clone(), None);
}
#[tokio::test]
async fn sends_multiple_handoff_outputs_until_cleared() {
let (tx, rx) = bounded(4);
let state = RealtimeHandoffState::new(tx);
state
.send_output("ignored".to_string())
.await
.expect("send");
assert!(rx.is_empty());
*state.active_handoff.lock().await = Some("handoff_1".to_string());
state.send_output("result".to_string()).await.expect("send");
state
.send_output("result 2".to_string())
.await
.expect("send");
let output_1 = rx.recv().await.expect("recv");
assert_eq!(
output_1,
HandoffOutput {
handoff_id: "handoff_1".to_string(),
output_text: "result".to_string(),
}
);
let output_2 = rx.recv().await.expect("recv");
assert_eq!(
output_2,
HandoffOutput {
handoff_id: "handoff_1".to_string(),
output_text: "result 2".to_string(),
}
);
*state.active_handoff.lock().await = None;
state
.send_output("ignored after clear".to_string())
.await
.expect("send");
assert!(rx.is_empty());
}