From 2d66c1a97488764fa9b0daae789a8cd229ce7927 Mon Sep 17 00:00:00 2001 From: Alex Gamble Date: Tue, 16 Jun 2026 12:15:03 -0700 Subject: [PATCH] fix(core): keep realtime startup nonblocking --- codex-rs/core/src/realtime_conversation.rs | 87 ++++++++++--------- .../core/tests/suite/realtime_conversation.rs | 10 ++- 2 files changed, 54 insertions(+), 43 deletions(-) diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index 03d0172f43..ec3ef9eedc 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -948,48 +948,6 @@ async fn handle_start_inner( .await; } - let mut pending_events = Vec::new(); - let realtime_session_id = loop { - match events_rx.recv().await { - Ok(RealtimeEvent::SessionUpdated { - realtime_session_id, - instructions, - }) => { - let started_realtime_session_id = realtime_session_id.clone(); - pending_events.push(RealtimeEvent::SessionUpdated { - realtime_session_id, - instructions, - }); - break started_realtime_session_id; - } - Ok(RealtimeEvent::Error(message)) => { - sess.conversation.finish_if_active(&realtime_active).await; - return Err(CodexErr::Stream(message, None)); - } - Ok(event) => pending_events.push(event), - Err(_) if !realtime_active.load(Ordering::Relaxed) => return Ok(()), - Err(_) => { - sess.conversation.finish_if_active(&realtime_active).await; - return Err(CodexErr::Stream( - "realtime conversation ended before the upstream session was created" - .to_string(), - None, - )); - } - } - }; - - info!(%realtime_session_id, "realtime conversation started"); - - sess.send_event_raw(Event { - id: sub_id.to_string(), - msg: EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { - realtime_session_id: Some(realtime_session_id), - version, - }), - }) - .await; - let sess_clone = Arc::clone(sess); let sub_id = sub_id.to_string(); let fanout_realtime_active = Arc::clone(&realtime_active); @@ -999,10 +957,55 @@ async fn handle_start_inner( msg, }; let mut end = RealtimeConversationEnd::TransportClosed; + let mut pending_events = Vec::new(); + let realtime_session_id = loop { + match events_rx.recv().await { + Ok(RealtimeEvent::SessionUpdated { + realtime_session_id, + instructions, + }) => { + let started_realtime_session_id = realtime_session_id.clone(); + pending_events.push(RealtimeEvent::SessionUpdated { + realtime_session_id, + instructions, + }); + break Some(started_realtime_session_id); + } + Ok(RealtimeEvent::Error(message)) => { + pending_events.push(RealtimeEvent::Error(message)); + end = RealtimeConversationEnd::Error; + break None; + } + Ok(event) => pending_events.push(event), + Err(_) if !fanout_realtime_active.load(Ordering::Relaxed) => return, + Err(_) => { + pending_events.push(RealtimeEvent::Error( + "realtime conversation ended before the upstream session was created" + .to_string(), + )); + end = RealtimeConversationEnd::Error; + break None; + } + } + }; + let startup_succeeded = realtime_session_id.is_some(); + if let Some(realtime_session_id) = realtime_session_id { + info!(%realtime_session_id, "realtime conversation started"); + sess_clone + .send_event_raw(ev(EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { + realtime_session_id: Some(realtime_session_id), + version, + }, + ))) + .await; + } + let mut pending_events = pending_events.into_iter(); loop { let event = match pending_events.next() { Some(event) => event, + None if !startup_succeeded => break, None => match events_rx.recv().await { Ok(event) => event, Err(_) => break, diff --git a/codex-rs/core/tests/suite/realtime_conversation.rs b/codex-rs/core/tests/suite/realtime_conversation.rs index 13870fbe71..ee6b7e3763 100644 --- a/codex-rs/core/tests/suite/realtime_conversation.rs +++ b/codex-rs/core/tests/suite/realtime_conversation.rs @@ -418,7 +418,11 @@ async fn conversation_start_defaults_to_v2_and_gpt_realtime_1_5() -> Result<()> skip_if_no_network!(Ok(())); let api_server = start_mock_server().await; - let realtime_server = start_websocket_server(vec![vec![vec![]]]).await; + let realtime_server = start_websocket_server(vec![vec![vec![json!({ + "type": "session.updated", + "session": { "id": "sess_defaults", "instructions": "backend prompt" } + })]]]) + .await; let realtime_base_url = realtime_server.uri().to_string(); let mut builder = test_codex().with_config(move |config| { config.experimental_realtime_ws_base_url = Some(realtime_base_url); @@ -449,6 +453,10 @@ async fn conversation_start_defaults_to_v2_and_gpt_realtime_1_5() -> Result<()> }) .await .expect("conversation start failed"); + assert_eq!( + started.realtime_session_id.as_deref(), + Some("sess_defaults") + ); assert!( realtime_server