diff --git a/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst b/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst index c43dd8756d20..3ea880b4b691 100644 Binary files a/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst and b/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst differ diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 5b99a2b2d0ab..6c930f9a1f64 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -1709,6 +1709,8 @@ Event notifications are the server-initiated event stream for thread lifecycles, Thread realtime publishes thread-scoped timeline item lifecycle notifications for paginated threads alongside its existing realtime notifications. Completed timeline items are durably interleaved with ordinary turn items by `thread/timeline/list`. Neither surface changes `ThreadItem`, `thread/read`, `thread/resume`, or `thread/fork`; clients ignore notification methods they do not recognize. +Core records transcript segments, session boundaries, and backing-agent artifact promotions through its injected thread store, even without an app-server event listener. Presentation selection uses the same rules for every Core host. App-server translates Core's history events into the notifications below; it does not append those items again. Recording remains limited to paginated threads. A completed notification follows acceptance by the thread store, not an additional flush or power-loss durability barrier. + Each realtime item has an `id`, a `realtimeSessionId`, and one of four types: `realtimeSessionStarted`, `transcriptSegment`, `bemItemPromoted`, or `realtimeSessionClosed`. A `bemItemPromoted` item references an existing backing-agent item by `turnId` and `itemId`; its `presentation` is `wholeItem`, `inlineMarkdown`, or `inlineVisualization` with an `index`. Recoverable configuration and initialization warnings use the existing `configWarning` notification: `{ summary, details?, path?, range? }`. App-server may emit it during initialization for config parsing and related setup diagnostics, or to the requesting connection during `thread/start` when that thread's exec-policy rules fail to parse. diff --git a/codex-rs/app-server/src/bespoke_event_handling.rs b/codex-rs/app-server/src/bespoke_event_handling.rs index 7fcfeb8d3540..b16d931e5707 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -61,6 +61,9 @@ use codex_app_server_protocol::ThreadItem; use codex_app_server_protocol::ThreadRealtimeClosedNotification; use codex_app_server_protocol::ThreadRealtimeErrorNotification; use codex_app_server_protocol::ThreadRealtimeItemAddedNotification; +use codex_app_server_protocol::ThreadRealtimeItemCompletedNotification; +use codex_app_server_protocol::ThreadRealtimeItemStartedNotification; +use codex_app_server_protocol::ThreadRealtimeItemTranscriptDeltaNotification; use codex_app_server_protocol::ThreadRealtimeOutputAudioDeltaNotification; use codex_app_server_protocol::ThreadRealtimeSdpNotification; use codex_app_server_protocol::ThreadRealtimeStartedNotification; @@ -463,6 +466,39 @@ pub(crate) async fn apply_bespoke_event_handling( .await; } EventMsg::RealtimeConversationRealtime(event) => match event.payload { + RealtimeEvent::HistoryItemStarted(item) => { + outgoing + .send_server_notification(ServerNotification::ThreadRealtimeItemStarted( + ThreadRealtimeItemStartedNotification { + thread_id: conversation_id.to_string(), + item: item.into(), + }, + )) + .await; + } + RealtimeEvent::HistoryTranscriptDelta { item_id, delta } => { + outgoing + .send_server_notification( + ServerNotification::ThreadRealtimeItemTranscriptDelta( + ThreadRealtimeItemTranscriptDeltaNotification { + thread_id: conversation_id.to_string(), + item_id, + delta, + }, + ), + ) + .await; + } + RealtimeEvent::HistoryItemCompleted(item) => { + outgoing + .send_server_notification(ServerNotification::ThreadRealtimeItemCompleted( + ThreadRealtimeItemCompletedNotification { + thread_id: conversation_id.to_string(), + item: item.into(), + }, + )) + .await; + } RealtimeEvent::SessionUpdated { .. } => {} RealtimeEvent::InputAudioSpeechStarted(event) => { let notification = ThreadRealtimeItemAddedNotification { diff --git a/codex-rs/app-server/src/lib.rs b/codex-rs/app-server/src/lib.rs index 04991451227a..76709f869c26 100644 --- a/codex-rs/app-server/src/lib.rs +++ b/codex-rs/app-server/src/lib.rs @@ -122,8 +122,6 @@ mod models_refresh_worker; mod notification_media; mod otel_reloader; mod outgoing_message; -mod realtime_event_handling; -mod realtime_history; mod request_processors; mod request_serialization; mod server_request_error; diff --git a/codex-rs/app-server/src/realtime_event_handling.rs b/codex-rs/app-server/src/realtime_event_handling.rs deleted file mode 100644 index cc4a29e4ec4f..000000000000 --- a/codex-rs/app-server/src/realtime_event_handling.rs +++ /dev/null @@ -1,91 +0,0 @@ -use crate::outgoing_message::ThreadScopedOutgoingMessageSender; -use crate::realtime_history::RealtimeEventEffects; -use codex_app_server_protocol::ServerNotification; -use codex_app_server_protocol::ThreadRealtimeItemCompletedNotification; -use codex_app_server_protocol::ThreadRealtimeItemStartedNotification; -use codex_app_server_protocol::ThreadRealtimeItemTranscriptDeltaNotification; -use codex_core::CodexThread; -use codex_protocol::ThreadId; -use codex_protocol::realtime::RealtimeItem; -use codex_protocol::realtime::RealtimeItemContent; -use codex_rollout::RolloutItem; -use tracing::warn; - -pub(crate) async fn apply_realtime_event_effects( - conversation: &CodexThread, - outgoing: &ThreadScopedOutgoingMessageSender, - thread_id: ThreadId, - effects: RealtimeEventEffects, -) { - let thread_id = thread_id.to_string(); - - if let Some(stream) = effects.transcript_stream { - if let Some(item) = stream.started_item { - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemStarted( - ThreadRealtimeItemStartedNotification { - thread_id: thread_id.clone(), - item: item.into(), - }, - )) - .await; - } - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemTranscriptDelta( - ThreadRealtimeItemTranscriptDeltaNotification { - thread_id: thread_id.clone(), - item_id: stream.item_id, - delta: stream.delta, - }, - )) - .await; - } - - if let Err(error) = - persist_realtime_items(conversation, outgoing, &thread_id, effects.items).await - { - warn!(thread_id, "failed to persist realtime history: {error}"); - } -} - -pub(crate) async fn persist_realtime_items( - conversation: &CodexThread, - outgoing: &ThreadScopedOutgoingMessageSender, - thread_id: &str, - items: Vec, -) -> Result<(), String> { - if items.is_empty() { - return Ok(()); - } - conversation - .append_rollout_items( - &items - .iter() - .cloned() - .map(RolloutItem::RealtimeItem) - .collect::>(), - ) - .await - .map_err(|error| format!("failed to persist realtime history: {error}"))?; - for item in items { - if !matches!(&item.content, RealtimeItemContent::TranscriptSegment { .. }) { - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemStarted( - ThreadRealtimeItemStartedNotification { - thread_id: thread_id.to_string(), - item: item.clone().into(), - }, - )) - .await; - } - outgoing - .send_server_notification(ServerNotification::ThreadRealtimeItemCompleted( - ThreadRealtimeItemCompletedNotification { - thread_id: thread_id.to_string(), - item: item.into(), - }, - )) - .await; - } - Ok(()) -} diff --git a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs index 7060ec4f3254..8354a4e0d37d 100644 --- a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs +++ b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs @@ -1,12 +1,8 @@ use super::*; use crate::extensions::send_thread_warning; -use crate::realtime_event_handling::apply_realtime_event_effects; -use crate::realtime_event_handling::persist_realtime_items; -use crate::realtime_history::RealtimeEventEffects; use codex_app_server_protocol::ThreadQueueChangedNotification; use codex_extension_api::ThreadIdleCause; use codex_protocol::config_types::MultiAgentMode; -use codex_protocol::protocol::ThreadHistoryMode; pub(super) const THREAD_UNLOADING_DELAY: Duration = Duration::from_secs(30 * 60); @@ -247,8 +243,6 @@ pub(super) async fn ensure_listener_task_running( ) .await; let config_snapshot = conversation.config_snapshot().await; - let realtime_history_enabled = - matches!(config_snapshot.history_mode, ThreadHistoryMode::Paginated); let thread_settings_baseline = thread_settings_from_config_snapshot(&config_snapshot); let (mut listener_command_rx, listener_generation) = { let mut thread_state = thread_state.lock().await; @@ -331,20 +325,10 @@ pub(super) async fn ensure_listener_task_running( // Track the event before emitting any typed translations // so thread-local state such as raw event opt-in stays // synchronized with the conversation. - let (raw_events_enabled, realtime_effects) = { + let raw_events_enabled = { let mut thread_state = thread_state.lock().await; thread_state.track_current_turn_event(&event.id, &event.msg); - let realtime_effects = if realtime_history_enabled - && thread_state.realtime_history.should_observe(&event.msg) - { - let active_turn_id = thread_state.active_turn_snapshot().map(|turn| turn.id); - thread_state - .realtime_history - .observe(&event.msg, active_turn_id.as_deref()) - } else { - RealtimeEventEffects::default() - }; - (thread_state.experimental_raw_events, realtime_effects) + thread_state.experimental_raw_events }; if matches!( &event.msg, @@ -362,14 +346,6 @@ pub(super) async fn ensure_listener_task_running( conversation_id, ); - apply_realtime_event_effects( - conversation.as_ref(), - &thread_outgoing, - conversation_id, - realtime_effects, - ) - .await; - apply_bespoke_event_handling( event.clone(), conversation_id, @@ -582,32 +558,6 @@ pub(super) async fn handle_thread_listener_command( .await; let _ = completion_tx.send(()); } - ThreadListenerCommand::SealRealtimeUserInput { - input, - completion_tx, - } => { - let items = thread_state - .lock() - .await - .realtime_history - .seal_user_input(&input); - let subscribed_connection_ids = thread_state_manager - .subscribed_connection_ids(conversation_id) - .await; - let thread_outgoing = ThreadScopedOutgoingMessageSender::new( - outgoing.clone(), - subscribed_connection_ids, - conversation_id, - ); - let result = persist_realtime_items( - conversation.as_ref(), - &thread_outgoing, - &conversation_id.to_string(), - items, - ) - .await; - let _ = completion_tx.send(result); - } } } diff --git a/codex-rs/app-server/src/request_processors/turn_processor.rs b/codex-rs/app-server/src/request_processors/turn_processor.rs index 129ad3b5a829..2cd4c436f8b4 100644 --- a/codex-rs/app-server/src/request_processors/turn_processor.rs +++ b/codex-rs/app-server/src/request_processors/turn_processor.rs @@ -619,10 +619,6 @@ impl TurnRequestProcessor { }, ) .await?; - if let TurnInput::UserInput { content, .. } = &input { - self.seal_realtime_transcript_before_user_input(thread_id, content) - .await?; - } let submission = thread .start_or_steer_turn( @@ -997,12 +993,12 @@ impl TurnRequestProcessor { request_id: &ConnectionRequestId, params: TurnSteerParams, ) -> Result { - let (thread_id, thread) = - self.load_thread(¶ms.thread_id) - .await - .inspect_err(|error| { - self.track_error_response(request_id, error, /*error_type*/ None); - })?; + let (_, thread) = self + .load_thread(¶ms.thread_id) + .await + .inspect_err(|error| { + self.track_error_response(request_id, error, /*error_type*/ None); + })?; self.ensure_direct_input_allowed(request_id, thread.as_ref()) .await?; @@ -1028,9 +1024,6 @@ impl TurnRequestProcessor { .collect(); let additional_context = map_additional_context(params.additional_context); - self.seal_realtime_transcript_before_user_input(thread_id, &mapped_items) - .await?; - let submission = thread .steer_turn( TurnInputRequest::new(TurnInput::UserInput { @@ -1127,37 +1120,6 @@ impl TurnRequestProcessor { Ok(TurnSteerResponse { turn_id }) } - async fn seal_realtime_transcript_before_user_input( - &self, - thread_id: ThreadId, - input: &[CoreInputItem], - ) -> Result<(), JSONRPCErrorError> { - let thread_state = self.thread_state_manager.thread_state(thread_id).await; - if !thread_state - .lock() - .await - .realtime_history - .should_seal_user_input(input) - { - return Ok(()); - } - let listener = self - .thread_state_manager - .current_listener_command_tx(thread_id) - .ok_or_else(|| internal_error("thread listener is not running"))?; - let (completion_tx, completion_rx) = tokio::sync::oneshot::channel(); - listener - .send(ThreadListenerCommand::SealRealtimeUserInput { - input: input.to_vec(), - completion_tx, - }) - .map_err(|_| internal_error("thread listener is not running"))?; - completion_rx - .await - .map_err(|_| internal_error("thread listener stopped before sealing realtime input"))? - .map_err(internal_error) - } - async fn prepare_realtime_conversation_thread( &self, request_id: &ConnectionRequestId, diff --git a/codex-rs/app-server/src/thread_state.rs b/codex-rs/app-server/src/thread_state.rs index 402783151ed2..3795638cdd36 100644 --- a/codex-rs/app-server/src/thread_state.rs +++ b/codex-rs/app-server/src/thread_state.rs @@ -1,6 +1,5 @@ use crate::outgoing_message::ConnectionId; use crate::outgoing_message::ConnectionRequestId; -use crate::realtime_history::RealtimeHistoryState; use codex_app_server_protocol::RequestId; use codex_app_server_protocol::ThreadGoal; use codex_app_server_protocol::ThreadHistoryBuilder; @@ -18,7 +17,6 @@ use codex_protocol::items::AgentMessageContent as CoreAgentMessageContent; use codex_protocol::items::TurnItem as CoreTurnItem; use codex_protocol::models::MessagePhase; use codex_protocol::protocol::EventMsg; -use codex_protocol::user_input::UserInput; use codex_rollout::RolloutItem; use codex_rollout::state_db::StateDbHandle; use codex_utils_path_uri::LegacyAppPathString; @@ -83,10 +81,6 @@ pub(crate) enum ThreadListenerCommand { request_id: RequestId, completion_tx: oneshot::Sender<()>, }, - SealRealtimeUserInput { - input: Vec, - completion_tx: oneshot::Sender>, - }, } /// Per-conversation accumulation of the latest states e.g. error message while a turn runs. @@ -110,7 +104,6 @@ pub(crate) struct ThreadState { pub(crate) cancel_tx: Option>, pub(crate) experimental_raw_events: bool, pub(crate) listener_generation: u64, - pub(crate) realtime_history: RealtimeHistoryState, last_thread_settings: Option, listener_command_tx: Option>, current_turn_history: ThreadHistoryBuilder, diff --git a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs index 9cf86981bf08..5cc7bb3a25ea 100644 --- a/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs +++ b/codex-rs/app-server/tests/suite/v2/realtime_conversation.rs @@ -760,10 +760,39 @@ async fn realtime_conversation_streams_timeline_items() -> Result<()> { assert_eq!(completed.item.id, started.item.id); assert_eq!(delta.delta, "hello"); assert!(matches!( - completed.item.content, + &completed.item.content, ThreadRealtimeItemContent::TranscriptSegment { text, .. } if text == "hello" )); + let closed = read_notification::( + &mut mcp, + "thread/realtime/item/completed", + ) + .await?; + let _: ThreadRealtimeClosedNotification = + read_notification(&mut mcp, "thread/realtime/closed").await?; + let request = mcp + .send_thread_timeline_list_request(ThreadTimelineListParams { + thread_id: thread.thread.id, + cursor: None, + limit: Some(100), + }) + .await?; + let page: ThreadTimelineListResponse = + timeout(DEFAULT_TIMEOUT, mcp.read_response(request)).await??; + let persisted = page + .data + .into_iter() + .filter_map(|entry| match entry { + ThreadTimelineEntry::Realtime { item, .. } => Some(item), + _ => None, + }) + .collect::>(); + assert_eq!( + persisted, + vec![session_completed.item, completed.item, closed.item] + ); + realtime_server.shutdown().await; Ok(()) } @@ -1185,6 +1214,10 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa "type": "response.output_text.delta", "delta": "Spoken before steering" })], + vec![json!({ + "type": "response.output_text.delta", + "delta": " and after rejected steering" + })], ])]), ) .await?; @@ -1230,6 +1263,45 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa ) .await?; + for (input, expected_turn_id) in [ + (Vec::new(), turn.turn.id.clone()), + ( + vec![V2UserInput::Text { + text: "Rejected steering".to_string(), + text_elements: Vec::new(), + }], + "stale-turn".to_string(), + ), + ] { + let request = harness + .mcp + .send_turn_steer_request(TurnSteerParams { + thread_id: harness.thread_id.clone(), + input, + expected_turn_id, + additional_context: None, + client_user_message_id: None, + responsesapi_client_metadata: None, + }) + .await?; + let rejected = timeout( + DEFAULT_TIMEOUT, + harness + .mcp + .read_stream_until_error_message(RequestId::Integer(request)), + ) + .await??; + assert_eq!(rejected.error.code, -32600); + } + harness + .append_text(harness.thread_id.clone(), "Continue speech after rejection") + .await?; + harness + .read_notification::( + "thread/realtime/transcript/delta", + ) + .await?; + let steering_request = harness .mcp .send_turn_steer_request(TurnSteerParams { @@ -1274,7 +1346,7 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa .. }, .. - } if text == "Spoken before steering" + } if text == "Spoken before steering and after rejected steering" ) }) .context("accepted steering should seal the active transcript")?; @@ -1296,16 +1368,25 @@ async fn realtime_timeline_splits_accepted_steering_and_persists_promoted_artifa }) .context("accepted steering should be included in the timeline")?; assert!(transcript_index < steering_index); - assert!(page.data.iter().any(|entry| matches!( - entry, - ThreadTimelineEntry::Realtime { - item: ThreadRealtimeItem { - content: ThreadRealtimeItemContent::BemItemPromoted { item_id, .. }, - .. - }, - .. - } if item_id == "promoted-message" - ))); + // Inline artifacts are promoted while streaming, before their final item. + assert_eq!( + page.data + .iter() + .filter_map(|entry| match entry { + ThreadTimelineEntry::Item { item, .. } + if matches!(item.as_ref(), ThreadItem::AgentMessage { id, .. } if id == "promoted-message") => Some("artifact"), + ThreadTimelineEntry::Realtime { + item: ThreadRealtimeItem { + content: ThreadRealtimeItemContent::BemItemPromoted { item_id, .. }, + .. + }, + .. + } if item_id == "promoted-message" => Some("promotion"), + _ => None, + }) + .collect::>(), + vec!["promotion", "artifact"] + ); for entry in &page.data { if let ThreadTimelineEntry::Realtime { item, .. } = entry { assert_eq!(Uuid::parse_str(&item.id)?.get_version_num(), 7); diff --git a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs index 7dd84e536613..293305ea538a 100644 --- a/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs +++ b/codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs @@ -652,7 +652,10 @@ impl RealtimeWebsocketEvents { | RealtimeEvent::ConversationItemDone { .. } | RealtimeEvent::NoopRequested(_) | RealtimeEvent::ConversationItemAdded(_) - | RealtimeEvent::Error(_) => {} + | RealtimeEvent::Error(_) + | RealtimeEvent::HistoryItemStarted(_) + | RealtimeEvent::HistoryTranscriptDelta { .. } + | RealtimeEvent::HistoryItemCompleted(_) => {} } truncate_active_transcript(&mut active_transcript.entries); } diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 8463ae9e16e9..b4e979289267 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -11,6 +11,7 @@ mod client; mod client_common; mod realtime_context; mod realtime_conversation; +mod realtime_history; mod realtime_prompt; mod responses_metadata; mod responses_retry; diff --git a/codex-rs/core/src/realtime_conversation.rs b/codex-rs/core/src/realtime_conversation.rs index bec3a059853f..a810a48c1768 100644 --- a/codex-rs/core/src/realtime_conversation.rs +++ b/codex-rs/core/src/realtime_conversation.rs @@ -2449,7 +2449,10 @@ async fn handle_realtime_server_event( | RealtimeEvent::OutputTranscriptDelta(_) | RealtimeEvent::OutputTranscriptDone(_) | RealtimeEvent::ConversationItemAdded(_) - | RealtimeEvent::ConversationItemDone { .. } => false, + | RealtimeEvent::ConversationItemDone { .. } + | RealtimeEvent::HistoryItemStarted(_) + | RealtimeEvent::HistoryTranscriptDelta { .. } + | RealtimeEvent::HistoryItemCompleted(_) => false, }; if events_tx.send(event).await.is_err() { diff --git a/codex-rs/app-server/src/realtime_history.rs b/codex-rs/core/src/realtime_history.rs similarity index 73% rename from codex-rs/app-server/src/realtime_history.rs rename to codex-rs/core/src/realtime_history.rs index 507632f8e06d..8ba0a71f5557 100644 --- a/codex-rs/app-server/src/realtime_history.rs +++ b/codex-rs/core/src/realtime_history.rs @@ -1,12 +1,14 @@ -use codex_protocol::items::AgentMessageContent; -use codex_protocol::items::DynamicToolCallStatus; -use codex_protocol::items::McpToolCallStatus; +//! Records canonical Voice history for every Core host. Presentation rules share +//! transcript boundaries and deduplication state with the history reducer. +//! The reducer selects when its effects are persisted relative to their source event. + +mod presentation; + use codex_protocol::items::TurnItem; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::RealtimeEvent; use codex_protocol::protocol::RealtimeTranscriptDelta; use codex_protocol::protocol::RealtimeTranscriptDone; -use codex_protocol::protocol::SubAgentActivityKind; use codex_protocol::realtime::BemItemPresentation; use codex_protocol::realtime::RealtimeItem; use codex_protocol::realtime::RealtimeItemContent; @@ -18,12 +20,6 @@ use std::collections::HashSet; use std::collections::VecDeque; use uuid::Uuid; -const INLINE_MARKDOWN_DIRECTIVE: &str = "::codex-realtime-inline{}"; -const INLINE_VISUALIZATION_DIRECTIVE: &str = "::codex-inline-vis{"; -const VISUALIZE_DIRECTIVE: &str = "visualize{"; -const BACKTICK_FENCE: &str = "```"; -const TILDE_FENCE: &str = "~~~"; - #[derive(Debug, Clone, PartialEq, Eq)] struct ActiveSegment { session_id: String, @@ -75,10 +71,18 @@ pub(crate) struct RealtimeTranscriptStream { #[derive(Debug, Default)] pub(crate) struct RealtimeEventEffects { + pub(crate) order: RealtimeEventOrder, pub(crate) items: Vec, pub(crate) transcript_stream: Option, } +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub(crate) enum RealtimeEventOrder { + BeforeEvent, + #[default] + AfterEvent, +} + #[derive(Clone, Copy)] enum Continuation { Continue, @@ -86,11 +90,12 @@ enum Continuation { } /// Retains only live session state; durable history is served by the rollout index. -#[derive(Debug, Default)] +#[derive(Default)] pub(crate) struct RealtimeHistoryState { active_session_id: Option, active_segments: ActiveTranscriptSegments, streaming_agent_message: Option, + active_turn_id: Option, realtime_session_by_bem_turn: HashMap, promoted_bem_presentation_keys: HashSet, pending_handoffs: VecDeque, @@ -98,7 +103,7 @@ pub(crate) struct RealtimeHistoryState { } impl RealtimeHistoryState { - pub(crate) fn should_seal_user_input(&self, input: &[UserInput]) -> bool { + fn should_seal_user_input(&self, input: &[UserInput]) -> bool { self.active_session_id.is_some() && [&self.active_segments.user, &self.active_segments.assistant] .into_iter() @@ -111,7 +116,7 @@ impl RealtimeHistoryState { }) } - pub(crate) fn seal_user_input(&mut self, input: &[UserInput]) -> Vec { + fn seal_user_input(&mut self, input: &[UserInput]) -> Vec { if !self.should_seal_user_input(input) { return Vec::new(); } @@ -121,17 +126,22 @@ impl RealtimeHistoryState { } pub(crate) fn should_observe(&self, event: &EventMsg) -> bool { - matches!(event, EventMsg::RealtimeConversationStarted(_)) - || (self.active_session_id.is_some() - && matches!( - event, - EventMsg::RealtimeConversationRealtime(_) - | EventMsg::RealtimeConversationClosed(_) - | EventMsg::TurnStarted(_) - | EventMsg::ItemStarted(_) - | EventMsg::ItemCompleted(_) - | EventMsg::AgentMessageContentDelta(_) - )) + matches!( + event, + EventMsg::RealtimeConversationStarted(_) + | EventMsg::TurnStarted(_) + | EventMsg::TurnComplete(_) + | EventMsg::TurnAborted(_) + ) || (self.active_session_id.is_some() + && matches!( + event, + EventMsg::RealtimeConversationRealtime(_) + | EventMsg::RealtimeConversationClosed(_) + | EventMsg::TurnStarted(_) + | EventMsg::ItemStarted(_) + | EventMsg::ItemCompleted(_) + | EventMsg::AgentMessageContentDelta(_) + )) || match event { EventMsg::ItemStarted(event) => self .realtime_session_by_bem_turn @@ -142,16 +152,29 @@ impl RealtimeHistoryState { EventMsg::AgentMessageContentDelta(event) => self .realtime_session_by_bem_turn .contains_key(&event.turn_id), - EventMsg::TurnStarted(_) => !self.pending_handoffs.is_empty(), _ => false, } } - pub(crate) fn observe( - &mut self, - event: &EventMsg, - active_turn_id: Option<&str>, - ) -> RealtimeEventEffects { + pub(crate) fn observe(&mut self, event: &EventMsg) -> RealtimeEventEffects { + match event { + EventMsg::TurnStarted(event) => self.active_turn_id = Some(event.turn_id.clone()), + EventMsg::TurnComplete(event) + if self.active_turn_id.as_deref() == Some(event.turn_id.as_str()) => + { + self.active_turn_id = None; + } + EventMsg::TurnAborted(event) + if event.turn_id.is_none() || event.turn_id == self.active_turn_id => + { + self.active_turn_id = None; + } + _ => {} + } + if !self.should_observe(event) { + return RealtimeEventEffects::default(); + } + let mut order = RealtimeEventOrder::AfterEvent; let mut items = Vec::new(); let mut transcript_stream = None; match event { @@ -170,7 +193,7 @@ impl RealtimeHistoryState { content: RealtimeItemContent::RealtimeSessionStarted, }); } - if let Some(turn_id) = active_turn_id { + if let Some(turn_id) = &self.active_turn_id { self.realtime_session_by_bem_turn .insert(turn_id.to_string(), session_id); } @@ -211,6 +234,7 @@ impl RealtimeHistoryState { } EventMsg::ItemStarted(event) => { if let TurnItem::UserMessage(item) = &event.item { + order = RealtimeEventOrder::BeforeEvent; items.extend(self.seal_user_input(&item.content)); } self.observe_item( @@ -222,6 +246,7 @@ impl RealtimeHistoryState { } EventMsg::ItemCompleted(event) => { if let TurnItem::UserMessage(item) = &event.item { + order = RealtimeEventOrder::BeforeEvent; items.extend(self.seal_user_input(&item.content)); } self.observe_item( @@ -273,121 +298,12 @@ impl RealtimeHistoryState { _ => {} } RealtimeEventEffects { + order, items, transcript_stream, } } - fn observe_item( - &mut self, - items: &mut Vec, - turn_id: &str, - item: &TurnItem, - completed: bool, - ) { - match item { - TurnItem::AgentMessage(message) => { - let text = message - .content - .iter() - .map(|content| match content { - AgentMessageContent::Text { text } => text.as_str(), - }) - .collect::(); - if !completed { - self.streaming_agent_message = Some(StreamingAgentMessage { - item_id: message.id.clone(), - text: text.clone(), - }); - } - self.observe_assistant_message(items, turn_id, &message.id, &text); - } - TurnItem::ImageGeneration(image) => { - self.add_promotion(items, turn_id, &image.id, BemItemPresentation::WholeItem); - } - TurnItem::Extension(extension) - if serde_json::to_value(extension) - .is_ok_and(|item| item["kind"] == "image_gen.generation") => - { - self.add_promotion( - items, - turn_id, - extension.id(), - BemItemPresentation::WholeItem, - ); - } - TurnItem::SubAgentActivity(activity) - if completed && activity.kind == SubAgentActivityKind::Started => - { - self.add_promotion(items, turn_id, &activity.id, BemItemPresentation::WholeItem); - } - TurnItem::DynamicToolCall(call) - if self.active_session_id.is_some() - && completed - && call.status == DynamicToolCallStatus::Completed - && call.success == Some(true) => - { - self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); - } - TurnItem::McpToolCall(call) - if self.active_session_id.is_some() - && completed - && call.server == "codex_app" - && call.status == McpToolCallStatus::Completed => - { - self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); - } - _ => {} - } - } - - fn observe_assistant_message( - &mut self, - items: &mut Vec, - turn_id: &str, - item_id: &str, - text: &str, - ) { - let mut lines = text.trim_start().lines(); - let mut first = lines.next().unwrap_or_default(); - if first.starts_with('[') - && let Some((_, content)) = first.split_once(']') - { - first = content.trim_start(); - if first.is_empty() { - first = lines.next().unwrap_or_default(); - } - } - if first == INLINE_MARKDOWN_DIRECTIVE && text.contains('\n') { - self.add_promotion(items, turn_id, item_id, BemItemPresentation::InlineMarkdown); - return; - } - - let mut in_fence = false; - let mut visualization_index = 0; - for line in text.lines() { - let trimmed = line.trim_start(); - if trimmed.starts_with(BACKTICK_FENCE) || trimmed.starts_with(TILDE_FENCE) { - in_fence = !in_fence; - continue; - } - if !in_fence - && (trimmed.starts_with(INLINE_VISUALIZATION_DIRECTIVE) - || trimmed.starts_with(VISUALIZE_DIRECTIVE)) - { - self.add_promotion( - items, - turn_id, - item_id, - BemItemPresentation::InlineVisualization { - index: visualization_index, - }, - ); - visualization_index += 1; - } - } - } - fn add_promotion( &mut self, items: &mut Vec, diff --git a/codex-rs/core/src/realtime_history/presentation.rs b/codex-rs/core/src/realtime_history/presentation.rs new file mode 100644 index 000000000000..f21c0d02b523 --- /dev/null +++ b/codex-rs/core/src/realtime_history/presentation.rs @@ -0,0 +1,129 @@ +//! Selects backing-agent items for the canonical Voice timeline using shared rules. + +use super::RealtimeHistoryState; +use super::StreamingAgentMessage; +use codex_protocol::items::AgentMessageContent; +use codex_protocol::items::DynamicToolCallStatus; +use codex_protocol::items::McpToolCallStatus; +use codex_protocol::items::TurnItem; +use codex_protocol::protocol::SubAgentActivityKind; +use codex_protocol::realtime::BemItemPresentation; +use codex_protocol::realtime::RealtimeItem; + +const INLINE_MARKDOWN_DIRECTIVE: &str = "::codex-realtime-inline{}"; +const INLINE_VISUALIZATION_DIRECTIVE: &str = "::codex-inline-vis{"; +const VISUALIZE_DIRECTIVE: &str = "visualize{"; +const BACKTICK_FENCE: &str = "```"; +const TILDE_FENCE: &str = "~~~"; + +impl RealtimeHistoryState { + pub(super) fn observe_item( + &mut self, + items: &mut Vec, + turn_id: &str, + item: &TurnItem, + completed: bool, + ) { + match item { + TurnItem::AgentMessage(message) => { + let text = message + .content + .iter() + .map(|content| match content { + AgentMessageContent::Text { text } => text.as_str(), + }) + .collect::(); + if !completed { + self.streaming_agent_message = Some(StreamingAgentMessage { + item_id: message.id.clone(), + text: text.clone(), + }); + } + self.observe_assistant_message(items, turn_id, &message.id, &text); + } + TurnItem::ImageGeneration(image) => { + self.add_promotion(items, turn_id, &image.id, BemItemPresentation::WholeItem); + } + TurnItem::Extension(extension) + if serde_json::to_value(extension) + .is_ok_and(|item| item["kind"] == "image_gen.generation") => + { + self.add_promotion( + items, + turn_id, + extension.id(), + BemItemPresentation::WholeItem, + ); + } + TurnItem::SubAgentActivity(activity) + if completed && activity.kind == SubAgentActivityKind::Started => + { + self.add_promotion(items, turn_id, &activity.id, BemItemPresentation::WholeItem); + } + TurnItem::DynamicToolCall(call) + if self.active_session_id.is_some() + && completed + && call.status == DynamicToolCallStatus::Completed + && call.success == Some(true) => + { + self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); + } + TurnItem::McpToolCall(call) + if self.active_session_id.is_some() + && completed + && call.server == "codex_app" + && call.status == McpToolCallStatus::Completed => + { + self.add_promotion(items, turn_id, &call.id, BemItemPresentation::WholeItem); + } + _ => {} + } + } + + pub(super) fn observe_assistant_message( + &mut self, + items: &mut Vec, + turn_id: &str, + item_id: &str, + text: &str, + ) { + let mut lines = text.trim_start().lines(); + let mut first = lines.next().unwrap_or_default(); + if first.starts_with('[') + && let Some((_, content)) = first.split_once(']') + { + first = content.trim_start(); + if first.is_empty() { + first = lines.next().unwrap_or_default(); + } + } + if first == INLINE_MARKDOWN_DIRECTIVE && text.contains('\n') { + self.add_promotion(items, turn_id, item_id, BemItemPresentation::InlineMarkdown); + return; + } + + let mut in_fence = false; + let mut visualization_index = 0; + for line in text.lines() { + let trimmed = line.trim_start(); + if trimmed.starts_with(BACKTICK_FENCE) || trimmed.starts_with(TILDE_FENCE) { + in_fence = !in_fence; + continue; + } + if !in_fence + && (trimmed.starts_with(INLINE_VISUALIZATION_DIRECTIVE) + || trimmed.starts_with(VISUALIZE_DIRECTIVE)) + { + self.add_promotion( + items, + turn_id, + item_id, + BemItemPresentation::InlineVisualization { + index: visualization_index, + }, + ); + visualization_index += 1; + } + } + } +} diff --git a/codex-rs/app-server/src/realtime_history_tests.rs b/codex-rs/core/src/realtime_history_tests.rs similarity index 65% rename from codex-rs/app-server/src/realtime_history_tests.rs rename to codex-rs/core/src/realtime_history_tests.rs index ad16bcff4590..7e0afae7b554 100644 --- a/codex-rs/app-server/src/realtime_history_tests.rs +++ b/codex-rs/core/src/realtime_history_tests.rs @@ -1,11 +1,15 @@ use super::*; use codex_protocol::AgentPath; use codex_protocol::ThreadId; +use codex_protocol::items::AgentMessageContent; use codex_protocol::items::AgentMessageItem; use codex_protocol::items::DynamicToolCallItem; +use codex_protocol::items::DynamicToolCallStatus; use codex_protocol::items::ImageGenerationItem; use codex_protocol::items::McpToolCallItem; +use codex_protocol::items::McpToolCallStatus; use codex_protocol::items::SubAgentActivityItem; +use codex_protocol::items::UserMessageItem; use codex_protocol::protocol::AgentMessageContentDeltaEvent; use codex_protocol::protocol::ItemCompletedEvent; use codex_protocol::protocol::ItemStartedEvent; @@ -13,18 +17,30 @@ use codex_protocol::protocol::RealtimeConversationClosedEvent; use codex_protocol::protocol::RealtimeConversationRealtimeEvent; use codex_protocol::protocol::RealtimeConversationStartedEvent; use codex_protocol::protocol::RealtimeConversationVersion; +use codex_protocol::protocol::SubAgentActivityKind; +use codex_protocol::protocol::TurnAbortReason; +use codex_protocol::protocol::TurnAbortedEvent; use pretty_assertions::assert_eq; +use test_case::test_case; use uuid::Uuid; fn started_state() -> RealtimeHistoryState { let mut state = RealtimeHistoryState::default(); - state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { + state.observe(&EventMsg::TurnStarted( + codex_protocol::protocol::TurnStartedEvent { + turn_id: "turn-1".to_string(), + trace_id: None, + started_at: None, + model_context_window: None, + collaboration_mode_kind: Default::default(), + }, + )); + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { realtime_session_id: Some("voice-1".to_string()), version: RealtimeConversationVersion::V2, - }), - Some("turn-1"), - ); + }, + )); state } @@ -32,10 +48,9 @@ fn observe_realtime( state: &mut RealtimeHistoryState, payload: RealtimeEvent, ) -> RealtimeEventEffects { - state.observe( - &EventMsg::RealtimeConversationRealtime(RealtimeConversationRealtimeEvent { payload }), - /*active_turn_id*/ None, - ) + state.observe(&EventMsg::RealtimeConversationRealtime( + RealtimeConversationRealtimeEvent { payload }, + )) } fn assistant_delta(item_id: &str, delta: &str) -> EventMsg { @@ -102,20 +117,96 @@ fn contents(effects: RealtimeEventEffects) -> Vec { .collect() } +#[test_case(Some("turn-1"), None; "matching_turn")] +#[test_case(None, None; "without_turn_id")] +#[test_case(Some("another-turn"), Some("voice-2"); "unrelated_turn")] +fn interrupted_turn_is_not_associated_with_a_new_voice_session( + aborted_turn_id: Option<&str>, + expected_session_id: Option<&str>, +) { + let mut state = RealtimeHistoryState::default(); + state.observe(&EventMsg::TurnStarted( + codex_protocol::protocol::TurnStartedEvent { + turn_id: "turn-1".to_string(), + trace_id: None, + started_at: None, + model_context_window: None, + collaboration_mode_kind: Default::default(), + }, + )); + let aborted = EventMsg::TurnAborted(TurnAbortedEvent { + turn_id: aborted_turn_id.map(str::to_string), + reason: TurnAbortReason::Interrupted, + started_at: None, + completed_at: None, + duration_ms: None, + }); + assert!(state.should_observe(&aborted)); + state.observe(&aborted); + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { + realtime_session_id: Some("voice-2".to_string()), + version: RealtimeConversationVersion::V2, + }, + )); + + let promoted = state + .observe(&completed_item(dynamic_tool_call( + "late-tool", + DynamicToolCallStatus::Completed, + ))) + .items; + assert_eq!( + promoted + .iter() + .map(|item| item.realtime_session_id.as_str()) + .collect::>(), + expected_session_id.into_iter().collect::>() + ); +} + +#[test] +fn interrupted_turn_keeps_its_existing_voice_session_for_late_artifacts() { + let mut state = started_state(); + state.observe(&EventMsg::TurnAborted(TurnAbortedEvent { + turn_id: Some("turn-1".to_string()), + reason: TurnAbortReason::Interrupted, + started_at: None, + completed_at: None, + duration_ms: None, + })); + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { + realtime_session_id: Some("voice-2".to_string()), + version: RealtimeConversationVersion::V2, + }, + )); + let promoted = state + .observe(&completed_item(dynamic_tool_call( + "late-tool", + DynamicToolCallStatus::Completed, + ))) + .items; + assert_eq!( + promoted + .iter() + .map(|item| item.realtime_session_id.as_str()) + .collect::>(), + vec!["voice-1"] + ); +} + #[test] fn promotes_backing_agent_artifacts_once_without_a_client_request() { let mut state = started_state(); let first_delta = assistant_delta("message-1", "[analysis] ::codex-realtime-inline{}"); - assert!( - state - .observe(&first_delta, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&first_delta).items.is_empty()); let second_delta = assistant_delta("message-1", "\nVisible explanation"); + let effects = state.observe(&second_delta); + assert_eq!(effects.order, RealtimeEventOrder::AfterEvent); assert_eq!( - contents(state.observe(&second_delta, /*active_turn_id*/ None)), + contents(effects), vec![RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), item_id: "message-1".to_string(), @@ -132,16 +223,11 @@ fn promotes_backing_agent_artifacts_once_without_a_client_request() { memory_citation: None, delivery: None, })); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); let next_message = assistant_delta("message-2", "::codex-realtime-inline{}\nNext message"); assert_eq!( - contents(state.observe(&next_message, /*active_turn_id*/ None)), + contents(state.observe(&next_message)), vec![RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), item_id: "message-2".to_string(), @@ -169,7 +255,7 @@ fn promotes_backing_agent_artifacts_once_without_a_client_request() { })); for (event, item_id) in [(image, "image-1"), (subagent, "subagent-1")] { assert_eq!( - contents(state.observe(&event, /*active_turn_id*/ None)), + contents(state.observe(&event)), vec![RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), item_id: item_id.to_string(), @@ -177,6 +263,15 @@ fn promotes_backing_agent_artifacts_once_without_a_client_request() { }] ); } + state.observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )); + let late = state.observe(&assistant_delta( + "late-artifact", + "::codex-realtime-inline{}\npresent later", + )); + assert_eq!(late.items.len(), 1); + assert_eq!(late.items[0].realtime_session_id, "voice-1"); } #[test] @@ -185,7 +280,7 @@ fn promotes_distinct_visualizations_once_and_ignores_markdown_fences() { let text = "```\n::codex-inline-vis{file=hidden}\n```\n~~~\nvisualize{file=also-hidden}\n~~~\n::codex-inline-vis{file=first}\nvisualize{file=second}"; let message = assistant_delta("message-1", text); assert_eq!( - contents(state.observe(&message, /*active_turn_id*/ None)), + contents(state.observe(&message)), (0..2) .map(|index| RealtimeItemContent::BemItemPromoted { turn_id: "turn-1".to_string(), @@ -204,12 +299,7 @@ fn promotes_distinct_visualizations_once_and_ignores_markdown_fences() { memory_citation: None, delivery: None, })); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); } #[test] @@ -302,12 +392,9 @@ fn preserves_first_transcript_activity_order_when_both_speakers_are_active() { ); } let items = state - .observe( - &EventMsg::RealtimeConversationClosed(RealtimeConversationClosedEvent { - reason: None, - }), - /*active_turn_id*/ None, - ) + .observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )) .items; let observed_roles = items .iter() @@ -333,13 +420,7 @@ fn does_not_replay_a_final_transcript_after_a_promotion_split() { "successful-tool", DynamicToolCallStatus::Completed, )); - assert_eq!( - state - .observe(&promotion, /*active_turn_id*/ None) - .items - .len(), - 2 - ); + assert_eq!(state.observe(&promotion).items.len(), 2); let done = observe_realtime( &mut state, @@ -356,17 +437,15 @@ fn generates_distinct_boundary_ids_for_reused_realtime_sessions() { let mut state = RealtimeHistoryState::default(); let mut boundary_ids = Vec::new(); for _ in 0..2 { - let started = state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { + let started = state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { realtime_session_id: Some("reused-session".to_string()), version: RealtimeConversationVersion::V2, - }), - /*active_turn_id*/ None, - ); - let closed = state.observe( - &EventMsg::RealtimeConversationClosed(RealtimeConversationClosedEvent { reason: None }), - /*active_turn_id*/ None, - ); + }, + )); + let closed = state.observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )); for item in started.items.into_iter().chain(closed.items) { assert_eq!(item.realtime_session_id, "reused-session"); assert_eq!(Uuid::parse_str(&item.id).unwrap().get_version_num(), 7); @@ -378,13 +457,12 @@ fn generates_distinct_boundary_ids_for_reused_realtime_sessions() { .collect::>(); assert_eq!(unique_ids.len(), 4); - let fallback = state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { + let fallback = state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { realtime_session_id: None, version: RealtimeConversationVersion::V2, - }), - /*active_turn_id*/ None, - ); + }, + )); assert_eq!(fallback.items.len(), 1); assert_eq!( Uuid::parse_str(&fallback.items[0].realtime_session_id) @@ -410,10 +488,11 @@ fn promotes_successful_codex_app_mcp_calls_and_splits_active_transcripts() { ] { assert!( state - .observe( - &completed_item(mcp_tool_call("not-promoted", server, status)), - /*active_turn_id*/ None, - ) + .observe(&completed_item(mcp_tool_call( + "not-promoted", + server, + status + )),) .items .is_empty() ); @@ -424,7 +503,7 @@ fn promotes_successful_codex_app_mcp_calls_and_splits_active_transcripts() { McpToolCallStatus::Completed, )); assert_eq!( - contents(state.observe(&completed, /*active_turn_id*/ None)), + contents(state.observe(&completed)), vec![ RealtimeItemContent::TranscriptSegment { role: RealtimeTranscriptRole::Assistant, @@ -437,12 +516,7 @@ fn promotes_successful_codex_app_mcp_calls_and_splits_active_transcripts() { }, ] ); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); } #[test] @@ -461,18 +535,15 @@ fn promotes_successful_dynamic_tools_and_splits_active_transcripts() { "failed-tool", DynamicToolCallStatus::Failed, )); - assert!( - state - .observe(&failed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&failed).items.is_empty()); let completed = completed_item(dynamic_tool_call( "successful-tool", DynamicToolCallStatus::Completed, )); - let items = state.observe(&completed, /*active_turn_id*/ None).items; + let effects = state.observe(&completed); + assert_eq!(effects.order, RealtimeEventOrder::AfterEvent); + let items = effects.items; assert_eq!(items.len(), 2); assert_eq!( items[0], @@ -505,37 +576,89 @@ fn promotes_successful_dynamic_tools_and_splits_active_transcripts() { .expect("second transcript stream"); assert_ne!(first.item_id, second.item_id); assert!(second.started_item.is_some()); - assert!( - state - .observe(&completed, /*active_turn_id*/ None) - .items - .is_empty() - ); + assert!(state.observe(&completed).items.is_empty()); - let closed = state.observe( - &EventMsg::RealtimeConversationClosed(RealtimeConversationClosedEvent { reason: None }), - /*active_turn_id*/ None, - ); + let closed = state.observe(&EventMsg::RealtimeConversationClosed( + RealtimeConversationClosedEvent { reason: None }, + )); assert_eq!(closed.items[0].id, second.item_id); let late = completed_item(dynamic_tool_call( "late-tool", DynamicToolCallStatus::Completed, )); + assert!(state.observe(&late).items.is_empty()); + + state.observe(&EventMsg::RealtimeConversationStarted( + RealtimeConversationStartedEvent { + realtime_session_id: Some("voice-2".to_string()), + version: RealtimeConversationVersion::V2, + }, + )); + let resumed = state.observe(&late).items; + assert_eq!(resumed.len(), 1); + assert_eq!(resumed[0].realtime_session_id, "voice-2"); +} + +#[test_case(|item| EventMsg::ItemStarted(ItemStartedEvent { + thread_id: ThreadId::new(), + turn_id: "turn-1".to_string(), + item, + started_at_ms: 0, +}); "started")] +#[test_case(completed_item; "completed")] +fn typed_input_seals_both_roles_but_realtime_delegation_does_not( + user_event: fn(TurnItem) -> EventMsg, +) { + let mut state = started_state(); + for event in [ + RealtimeEvent::OutputTranscriptDelta(RealtimeTranscriptDelta { + delta: "assistant".to_string(), + }), + RealtimeEvent::InputTranscriptDelta(RealtimeTranscriptDelta { + delta: "user".to_string(), + }), + ] { + observe_realtime(&mut state, event); + } assert!( state - .observe(&late, /*active_turn_id*/ None) + .observe(&user_event(TurnItem::UserMessage(UserMessageItem::new(&[ + UserInput::Text { + text: "delegated".to_string(), + text_elements: Vec::new(), + } + ])))) .items .is_empty() ); - - state.observe( - &EventMsg::RealtimeConversationStarted(RealtimeConversationStartedEvent { - realtime_session_id: Some("voice-2".to_string()), - version: RealtimeConversationVersion::V2, - }), - Some("turn-1"), + let effects = state.observe(&user_event(TurnItem::UserMessage(UserMessageItem::new(&[ + UserInput::Text { + text: "typed steering".to_string(), + text_elements: Vec::new(), + }, + ])))); + assert_eq!(effects.order, RealtimeEventOrder::BeforeEvent); + assert_eq!( + contents(effects), + vec![ + RealtimeItemContent::TranscriptSegment { + role: RealtimeTranscriptRole::Assistant, + text: "assistant".to_string() + }, + RealtimeItemContent::TranscriptSegment { + role: RealtimeTranscriptRole::User, + text: "user".to_string() + }, + ] ); - let resumed = state.observe(&late, /*active_turn_id*/ None).items; - assert_eq!(resumed.len(), 1); - assert_eq!(resumed[0].realtime_session_id, "voice-2"); + for event in [ + RealtimeEvent::OutputTranscriptDone(RealtimeTranscriptDone { + text: "assistant revised".to_string(), + }), + RealtimeEvent::InputTranscriptDone(RealtimeTranscriptDone { + text: "user revised".to_string(), + }), + ] { + assert!(observe_realtime(&mut state, event).items.is_empty()); + } } diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index f7ab1271a0cb..aa0bfe12cc31 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -40,6 +40,7 @@ use crate::image_preparation::prepare_response_items as prepare_image_response_i use crate::image_preparation::unified_image_budget_enabled; use crate::parse_turn_item; use crate::realtime_conversation::RealtimeConversationManager; +use crate::realtime_history::RealtimeEventOrder; use crate::session::step_context::StepContext; use crate::session::step_settings::ResolvedStepSettings; use crate::session::step_settings::StepSettings; @@ -227,6 +228,7 @@ mod mcp_prewarm; mod mcp_refresh; mod mcp_runtime; pub(crate) mod multi_agents; +mod realtime_history; mod review; mod rollout_budget; mod rollout_reconstruction; @@ -2339,7 +2341,32 @@ impl Session { } async fn send_event_raw_with_persistence(&self, event: Event, persist: bool) { + // Keep realtime reduction, canonical append, and delivery in the same order. + // This lock must not acquire SessionState or ActiveTurn: event producers can + // already hold those locks. Host presentation policies are synchronous. + let mut realtime_history = match &self.realtime_history { + Some(history) => { + let history = history.lock().await; + history.should_observe(&event.msg).then_some(history) + } + None => None, + }; self.services.mcp_runtime.observe_event(&event.msg); + let (before_event, after_event) = match realtime_history.as_mut() { + Some(history) => { + let effects = history.observe(&event.msg); + match effects.order { + RealtimeEventOrder::BeforeEvent => (Some(effects), None), + RealtimeEventOrder::AfterEvent => (None, Some(effects)), + } + } + None => (None, None), + }; + if let Some(effects) = before_event + && let Err(error) = self.send_realtime_history_effects(&event.id, effects).await + { + warn!("failed to persist realtime history: {error}"); + } // Persist the event into rollout storage; the store applies its persistence policy. if persist { let rollout_items = vec![RolloutItem::EventMsg(event.msg.clone())]; @@ -2348,6 +2375,11 @@ impl Session { self.services .rollout_thread_trace .record_protocol_event(&event.msg); + if let Some(effects) = after_event + && let Err(error) = self.send_realtime_history_effects(&event.id, effects).await + { + warn!("failed to persist realtime history: {error}"); + } self.deliver_event_raw(event).await; } diff --git a/codex-rs/core/src/session/realtime_history.rs b/codex-rs/core/src/session/realtime_history.rs new file mode 100644 index 000000000000..5b1acf4873bb --- /dev/null +++ b/codex-rs/core/src/session/realtime_history.rs @@ -0,0 +1,56 @@ +use super::session::Session; +use crate::realtime_history::RealtimeEventEffects; +use codex_history::RolloutItem; +use codex_protocol::protocol::Event; +use codex_protocol::protocol::EventMsg; +use codex_protocol::protocol::RealtimeConversationRealtimeEvent; +use codex_protocol::protocol::RealtimeEvent; +use codex_protocol::realtime::RealtimeItemContent; + +impl Session { + pub(super) async fn send_realtime_history_effects( + &self, + submission_id: &str, + effects: RealtimeEventEffects, + ) -> anyhow::Result<()> { + let event = |payload| Event { + id: submission_id.to_string(), + msg: EventMsg::RealtimeConversationRealtime(RealtimeConversationRealtimeEvent { + payload, + }), + }; + if let Some(stream) = effects.transcript_stream { + if let Some(item) = stream.started_item { + self.deliver_event_raw(event(RealtimeEvent::HistoryItemStarted(item))) + .await; + } + self.deliver_event_raw(event(RealtimeEvent::HistoryTranscriptDelta { + item_id: stream.item_id, + delta: stream.delta, + })) + .await; + } + if effects.items.is_empty() { + return Ok(()); + } + self.live_thread_for_persistence("append realtime history")? + .append_items( + &effects + .items + .iter() + .cloned() + .map(RolloutItem::RealtimeItem) + .collect::>(), + ) + .await?; + for item in effects.items { + if !matches!(&item.content, RealtimeItemContent::TranscriptSegment { .. }) { + self.deliver_event_raw(event(RealtimeEvent::HistoryItemStarted(item.clone()))) + .await; + } + self.deliver_event_raw(event(RealtimeEvent::HistoryItemCompleted(item))) + .await; + } + Ok(()) + } +} diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index 3b3da7440ba8..63530cc31309 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -64,6 +64,7 @@ pub(crate) struct Session { pub(super) mcp_prewarm_shutdown: CancellationToken, pub(super) mcp_prewarm_task: std::sync::Mutex>>, pub(crate) conversation: Arc, + pub(crate) realtime_history: Option>, pub(crate) active_turn: Mutex>, pub(crate) async_hook_results: async_channel::Receiver, pub(crate) input_queue: InputQueue, @@ -1492,6 +1493,9 @@ impl Session { mcp_prewarm_shutdown: CancellationToken::new(), mcp_prewarm_task: std::sync::Mutex::new(None), conversation: Arc::new(RealtimeConversationManager::new()), + realtime_history: (session_configuration.history_mode == ThreadHistoryMode::Paginated + && services.live_thread.is_some()) + .then(|| Mutex::new(Default::default())), active_turn: Mutex::new(None), async_hook_results, input_queue: InputQueue::new(), diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index b8d97c87e526..268de9339131 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -6612,6 +6612,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { mcp_prewarm_shutdown: CancellationToken::new(), mcp_prewarm_task: std::sync::Mutex::new(None), conversation: Arc::new(RealtimeConversationManager::new()), + realtime_history: None, active_turn: Mutex::new(None), async_hook_results, input_queue: super::input_queue::InputQueue::new(), @@ -8892,6 +8893,7 @@ where mcp_prewarm_shutdown: CancellationToken::new(), mcp_prewarm_task: std::sync::Mutex::new(None), conversation: Arc::new(RealtimeConversationManager::new()), + realtime_history: None, active_turn: Mutex::new(None), async_hook_results, input_queue: super::input_queue::InputQueue::new(), diff --git a/codex-rs/core/src/session/turn_input_tests.rs b/codex-rs/core/src/session/turn_input_tests.rs index 3ef6797c065b..c9810df1f4cb 100644 --- a/codex-rs/core/src/session/turn_input_tests.rs +++ b/codex-rs/core/src/session/turn_input_tests.rs @@ -1,6 +1,7 @@ use super::*; use crate::config::Constrained; use crate::session::step_settings::StepSettingsUpdate; +use crate::session::tests::make_session_and_context; use crate::session::tests::make_session_and_context_with_rx; use crate::session::turn_context::TurnContext; use crate::state::TaskKind; @@ -107,6 +108,65 @@ async fn submit_steer_only( .expect("steer-only submission should be valid") } +#[tokio::test] +#[expect( + clippy::await_holding_invalid_type, + reason = "simulate an in-flight realtime append while checking input admission" +)] +async fn steering_does_not_wait_for_realtime_history() { + let (mut session, turn_context) = make_session_and_context().await; + session.realtime_history = Some(tokio::sync::Mutex::new(Default::default())); + let session = Arc::new(session); + let turn_context = Arc::new(turn_context); + session + .spawn_task( + Arc::clone(&turn_context), + Vec::new(), + NeverEndingTask { + kind: TaskKind::Regular, + listen_to_cancellation_token: true, + }, + ) + .await; + + let history = session + .realtime_history + .as_ref() + .expect("realtime history") + .lock() + .await; + for mode in [ + TurnInputMode::StartOrSteer, + TurnInputMode::Steer { + expected_turn_id: turn_context.sub_id.clone(), + }, + ] { + let submission = tokio::time::timeout( + std::time::Duration::from_secs(5), + handle( + &session, + TurnInputRequest::user_input(vec![UserInput::Text { + text: "steer without waiting for persistence".to_string(), + text_elements: Vec::new(), + }]), + mode, + "steer-submission".to_string(), + ), + ) + .await + .expect("steering must not wait for the realtime recorder") + .expect("steering should succeed"); + assert_eq!( + submission, + TurnInputSubmission::Steered { + turn_id: turn_context.sub_id.clone() + } + ); + } + drop(history); + session.abort_all_tasks(TurnAbortReason::Interrupted).await; +} + #[tokio::test] async fn accepted_input_applies_thread_settings() { let (session, turn_context, _rx) = make_session_and_context_with_rx().await; diff --git a/codex-rs/core/tests/suite/realtime_conversation.rs b/codex-rs/core/tests/suite/realtime_conversation.rs index de16d730ed57..b891fcfd6436 100644 --- a/codex-rs/core/tests/suite/realtime_conversation.rs +++ b/codex-rs/core/tests/suite/realtime_conversation.rs @@ -1,6 +1,10 @@ use anyhow::Context; use anyhow::Result; use chrono::Utc; +use codex_app_server_protocol::ThreadRealtimeItemContent; +use codex_app_server_protocol::ThreadRealtimeSessionOutcome; +use codex_app_server_protocol::ThreadRealtimeTranscriptRole; +use codex_app_server_protocol::ThreadTimelineEntry; use codex_config::config_toml::RealtimeWsMode; use codex_config::config_toml::RealtimeWsVersion; use codex_core::StartThreadOptions; @@ -33,8 +37,10 @@ use codex_protocol::protocol::RealtimeOutputModality; use codex_protocol::protocol::RealtimeTranscriptEntry; use codex_protocol::protocol::RealtimeVoice; use codex_protocol::protocol::SessionSource; +use codex_protocol::protocol::ThreadHistoryMode; use codex_protocol::protocol::ThreadSource; use codex_protocol::user_input::UserInput; +use codex_thread_store::ListTimelineParams; use core_test_support::responses; use core_test_support::responses::WebSocketConnectionConfig; use core_test_support::responses::start_mock_server; @@ -437,10 +443,148 @@ async fn conversation_start_audio_text_close_round_trip() -> Result<()> { Some("requested" | "transport_closed") )); + test.codex.ensure_rollout_materialized().await; + let history = test.codex.load_history(/*include_archived*/ false).await?; + assert!( + !history + .items + .iter() + .any(|item| matches!(item, RolloutItem::RealtimeItem(_))) + ); + server.shutdown().await; Ok(()) } +// No host consumes realtime events in this test. The injected store must receive +// canonical history even when the host has no realtime notification adapter. +#[test_case(ThreadRealtimeSessionOutcome::Ended; "transport_close")] +#[test_case(ThreadRealtimeSessionOutcome::Failed; "error_close")] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn conversation_records_history_without_an_event_observer( + outcome: ThreadRealtimeSessionOutcome, +) -> Result<()> { + skip_if_no_network!(Ok(())); + let api_server = start_mock_server().await; + let mut events = vec![ + json!({ "type": "session.updated", "session": { "id": "voice-1" } }), + json!({ "type": "response.output_text.delta", "delta": "assistant first" }), + json!({ "type": "conversation.item.input_audio_transcription.delta", "delta": "user second" }), + ]; + if outcome == ThreadRealtimeSessionOutcome::Failed { + events.push(json!({ "type": "error", "error": { "message": "fixture failure" } })); + } + let realtime_server = start_websocket_server(vec![vec![events.clone()], vec![events]]).await; + let mut builder = test_codex() + .with_history_mode(ThreadHistoryMode::Paginated) + .with_config({ + let url = realtime_server.uri().to_string(); + move |config| { + config.experimental_realtime_ws_base_url = Some(url); + config.realtime.version = RealtimeWsVersion::V2; + } + }); + let test = builder.build_with_auto_env(&api_server).await?; + test.codex.ensure_rollout_materialized().await; + let mut expected = Vec::new(); + for session_count in 1..=2 { + test.codex + .submit(Op::RealtimeConversationStart(ConversationStartParams { + client_managed_handoffs: false, + delegation_ack_filler: None, + flush_transcript_tail_on_session_end: false, + codex_responses_as_items: false, + codex_response_item_prefix: None, + codex_response_handoff_mode: + codex_protocol::protocol::CodexResponseHandoffMode::Thinking, + codex_response_handoff_channel_prefixes: None, + model: None, + output_modality: RealtimeOutputModality::Audio, + include_startup_context: false, + initial_items: Vec::new(), + realtime_start_instructions: None, + realtime_end_instructions: None, + prompt: Some(Some("fixture".to_string())), + realtime_session_id: Some("voice-1".to_string()), + transport: None, + version: None, + voice: None, + })) + .await?; + let items = timeout(Duration::from_secs(10), async { + loop { + test.thread_store + .flush_thread(test.session_configured.thread_id) + .await?; + let history = test + .thread_store + .list_timeline(ListTimelineParams { + thread_id: test.session_configured.thread_id, + cursor: None, + page_size: 100, + }) + .await?; + let items = history + .items + .into_iter() + .filter_map(|item| match item { + ThreadTimelineEntry::Realtime { item, .. } => Some(item), + _ => None, + }) + .collect::>(); + if items + .iter() + .filter(|item| { + matches!( + item.content, + ThreadRealtimeItemContent::RealtimeSessionClosed { .. } + ) + }) + .count() + == session_count + { + break Ok::<_, anyhow::Error>(items); + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .context("Core should persist history without a host observer")??; + expected.extend([ + ThreadRealtimeItemContent::RealtimeSessionStarted, + ThreadRealtimeItemContent::TranscriptSegment { + role: ThreadRealtimeTranscriptRole::Assistant, + text: "assistant first".to_string(), + }, + ThreadRealtimeItemContent::TranscriptSegment { + role: ThreadRealtimeTranscriptRole::User, + text: "user second".to_string(), + }, + ThreadRealtimeItemContent::RealtimeSessionClosed { outcome }, + ]); + assert_eq!( + items + .iter() + .map(|item| item.content.clone()) + .collect::>(), + expected + ); + let ids = items + .iter() + .map(|item| item.id.as_str()) + .collect::>(); + assert_eq!(ids.len(), items.len()); + assert!( + items + .iter() + .all(|item| item.realtime_session_id == "voice-1") + ); + } + test.codex.submit(Op::Shutdown).await?; + realtime_server.shutdown().await; + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn conversation_start_defaults_to_v2_and_gpt_realtime_1_5() -> Result<()> { skip_if_no_network!(Ok(())); @@ -3530,8 +3674,10 @@ async fn conversation_startup_context_is_truncated_and_sent_once_per_start() -> Ok(()) } +#[test_case(false; "durable")] +#[test_case(true; "ephemeral")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { +async fn conversation_user_text_turn_is_not_sent_to_realtime(ephemeral: bool) -> Result<()> { skip_if_no_network!(Ok(())); let api_server = start_mock_server().await; @@ -3545,22 +3691,32 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { .await; let realtime_server = start_websocket_server(vec![vec![ - vec![json!({ - "type": "session.updated", - "session": { "id": "sess_user_text", "instructions": "backend prompt" } - })], + vec![ + json!({ + "type": "session.updated", + "session": { "id": "sess_user_text", "instructions": "backend prompt" } + }), + json!({ + "type": "response.output_text.delta", + "delta": "spoken before typed input" + }), + ], vec![], ]]) .await; - let mut builder = test_codex().with_config({ - let realtime_base_url = realtime_server.uri().to_string(); - move |config| { - config.experimental_realtime_ws_base_url = Some(realtime_base_url); - config.experimental_realtime_ws_startup_context = Some(String::new()); - } - }); - let test = builder.build(&api_server).await?; + let mut builder = test_codex() + .with_history_mode(ThreadHistoryMode::Paginated) + .with_config({ + let realtime_base_url = realtime_server.uri().to_string(); + move |config| { + config.ephemeral = ephemeral; + config.realtime.version = RealtimeWsVersion::V2; + config.experimental_realtime_ws_base_url = Some(realtime_base_url); + config.experimental_realtime_ws_startup_context = Some(String::new()); + } + }); + let test = builder.build_with_auto_env(&api_server).await?; test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { @@ -3599,6 +3755,20 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { .await; assert_eq!(session_updated, "sess_user_text"); + wait_for_event(&test.codex, |event| { + if let EventMsg::RealtimeConversationRealtime(event) = event { + match event.payload { + RealtimeEvent::HistoryItemStarted(_) + | RealtimeEvent::HistoryTranscriptDelta { .. } + | RealtimeEvent::HistoryItemCompleted(_) => assert!(!ephemeral), + RealtimeEvent::OutputTranscriptDelta(_) => return true, + _ => {} + } + } + false + }) + .await; + let user_text = "typed follow-up for realtime"; test.codex .start_or_steer_turn(TurnInputRequest::user_input(vec![UserInput::Text { @@ -3617,6 +3787,30 @@ async fn conversation_user_text_turn_is_not_sent_to_realtime() -> Result<()> { let model_user_texts = response_mock.single_request().message_input_texts("user"); assert!(model_user_texts.iter().any(|text| text == user_text)); + if !ephemeral { + test.thread_store + .flush_thread(test.session_configured.thread_id) + .await?; + let timeline = test + .thread_store + .list_timeline(ListTimelineParams { + thread_id: test.session_configured.thread_id, + cursor: None, + page_size: 100, + }) + .await?; + assert_eq!( + timeline.items.iter().filter_map(|entry| match entry { + ThreadTimelineEntry::Realtime { item, .. } + if matches!(&item.content, ThreadRealtimeItemContent::TranscriptSegment { text, .. } if text == "spoken before typed input") => Some("speech"), + ThreadTimelineEntry::Item { item, .. } + if matches!(item.as_ref(), codex_app_server_protocol::ThreadItem::UserMessage { .. }) => Some("typed input"), + _ => None, + }).collect::>(), + vec!["speech", "typed input"] + ); + } + let realtime_connections = realtime_server.connections(); assert_eq!(realtime_connections.len(), 1); assert_eq!(realtime_connections[0].len(), 1); @@ -4246,7 +4440,7 @@ fn message_input_texts(body: &Value, role: &str) -> Vec { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn inbound_handoff_request_starts_turn() -> Result<()> { +async fn inbound_handoff_request_starts_turn_and_promotes_its_artifact() -> Result<()> { skip_if_no_network!(Ok(())); let api_server = start_mock_server().await; @@ -4254,7 +4448,7 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { &api_server, responses::sse(vec![ responses::ev_response_created("resp-1"), - responses::ev_assistant_message("msg-1", "ok"), + responses::ev_assistant_message("msg-1", "::codex-realtime-inline{}\nShared artifact"), responses::ev_completed("resp-1"), ]), ) @@ -4278,14 +4472,16 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { ]]]) .await; - let mut builder = test_codex().with_config({ - let realtime_base_url = realtime_server.uri().to_string(); - move |config| { - config.experimental_realtime_ws_base_url = Some(realtime_base_url); - config.realtime.version = RealtimeWsVersion::V1; - } - }); - let test = builder.build(&api_server).await?; + let mut builder = test_codex() + .with_history_mode(ThreadHistoryMode::Paginated) + .with_config({ + let realtime_base_url = realtime_server.uri().to_string(); + move |config| { + config.experimental_realtime_ws_base_url = Some(realtime_base_url); + config.realtime.version = RealtimeWsVersion::V1; + } + }); + let test = builder.build_with_auto_env(&api_server).await?; test.codex .submit(Op::RealtimeConversationStart(ConversationStartParams { @@ -4349,6 +4545,42 @@ async fn inbound_handoff_request_starts_turn() -> Result<()> { }) .await; + test.thread_store + .flush_thread(test.session_configured.thread_id) + .await?; + let timeline = test + .thread_store + .list_timeline(ListTimelineParams { + thread_id: test.session_configured.thread_id, + cursor: None, + page_size: 100, + }) + .await?; + let promotions = timeline + .items + .into_iter() + .filter_map(|entry| match entry { + ThreadTimelineEntry::Realtime { item, .. } + if matches!( + item.content, + ThreadRealtimeItemContent::BemItemPromoted { .. } + ) => + { + Some(item.content) + } + _ => None, + }) + .collect::>(); + assert_eq!( + promotions, + vec![ThreadRealtimeItemContent::BemItemPromoted { + turn_id, + item_id: "msg-1".to_string(), + presentation: + codex_app_server_protocol::ThreadRealtimeBemItemPresentation::InlineMarkdown, + }] + ); + let request = response_mock.single_request(); let turn_metadata: Value = serde_json::from_str( request diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index 06ba9705f7b7..d4d236d963c2 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -448,6 +448,13 @@ pub enum RealtimeEvent { ConversationItemDone { item_id: String, }, + /// Canonical display history produced by Core, separate from provider events. + HistoryItemStarted(crate::realtime::RealtimeItem), + HistoryTranscriptDelta { + item_id: String, + delta: String, + }, + HistoryItemCompleted(crate::realtime::RealtimeItem), HandoffRequested(RealtimeHandoffRequested), NoopRequested(RealtimeNoopRequested), Error(String),