Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions codex-rs/app-server/tests/suite/v2/realtime_conversation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,9 @@ const V2_HANDOFF_COMPLETE_ACKNOWLEDGEMENT: &str =
const RESPONSE_ITEM_PREFIX: &str =
"Use the following context to inform future responses, but do not speak it to the user.";

#[path = "realtime_transcript_tests.rs"]
mod transcript_tests;

#[derive(Debug, Clone, Copy)]
enum StartupContextConfig<'a> {
Generated,
Expand Down
115 changes: 115 additions & 0 deletions codex-rs/app-server/tests/suite/v2/realtime_transcript_tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
//! Desktop-facing live notifications and reconciled voice handoff context.

use super::*;
use pretty_assertions::assert_eq;
use test_case::test_case;

#[test_case(
"A bunch ", "of stuff", "A bunch of stuff, including coding.",
"A bunch of stuff, including coding."; "final expands partial"
)]
#[test_case(
"Okay.", "I will check", "Okay.", "Okay.I will check";
"delayed final preserves newer speech"
)]
#[tokio::test]
async fn websocket_v3_reconciles_handoff_transcripts_without_changing_live_notifications(
first_delta: &str,
second_delta: &str,
assistant_final: &str,
assistant_transcript: &str,
) -> Result<()> {
skip_if_no_network!(Ok(()));

let mut harness = RealtimeE2eHarness::new(
RealtimeTestVersion::V1,
main_loop_responses(vec![create_final_assistant_message_sse_response("Done.")?]),
realtime_sideband(vec![open_realtime_sideband_connection(vec![vec![
session_started("interleaved"),
json!({"type": "input_transcript.added", "item": {"text": "What can you do?"}}),
json!({"type": "output_transcript.added", "item": {"text": first_delta}}),
json!({"type": "output_transcript.added", "item": {"text": second_delta}}),
json!({"type": "turn.done", "turn": {"role": "user", "transcript": "What can you do?"}}),
json!({"type": "turn.done", "turn": {"role": "assistant", "transcript": assistant_final}}),
json!({
"type": "delegation.created",
"item": {
"id": "handoff",
"type": "delegation",
"target": "client",
"content": [{"type": "input_text", "text": "What can you do?"}]
}
}),
]])]),
).await?;
harness
.start_frameless_bidi_realtime(
/*codex_response_handoff_mode*/ None,
/*codex_response_handoff_channel_prefixes*/ None, /*initial_items*/ None,
)
.await?;

for (role, delta) in [
("user", "What can you do?"),
("assistant", first_delta),
("assistant", second_delta),
] {
assert_eq!(
harness
.read_notification::<ThreadRealtimeTranscriptDeltaNotification>(
"thread/realtime/transcript/delta"
)
.await?,
ThreadRealtimeTranscriptDeltaNotification {
thread_id: harness.thread_id.clone(),
role: role.into(),
delta: delta.into()
},
);
}
for (role, text) in [("user", "What can you do?"), ("assistant", assistant_final)] {
assert_eq!(
harness
.read_notification::<ThreadRealtimeTranscriptDoneNotification>(
"thread/realtime/transcript/done"
)
.await?,
ThreadRealtimeTranscriptDoneNotification {
thread_id: harness.thread_id.clone(),
role: role.into(),
text: text.into()
},
);
}
assert_eq!(
harness
.read_notification::<ThreadRealtimeItemAddedNotification>("thread/realtime/itemAdded")
.await?,
ThreadRealtimeItemAddedNotification {
thread_id: harness.thread_id.clone(),
item: json!({
"type": "handoff_request",
"handoff_id": "handoff",
"item_id": "handoff",
"input_transcript": "What can you do?",
"active_transcript": [
{"role": "user", "text": "What can you do?"},
{"role": "assistant", "text": assistant_transcript},
],
}),
},
);
harness
.read_notification::<TurnCompletedNotification>("turn/completed")
.await?;
let requests = harness.main_loop_responses_requests().await?;
assert_eq!(requests.len(), 1);
assert!(response_request_contains_text(
&requests[0],
&format!(
"<realtime_delegation>\n <input>What can you do?</input>\n <transcript_delta>user: What can you do?\nassistant: {assistant_transcript}</transcript_delta>\n</realtime_delegation>"
),
));
harness.shutdown().await;
Ok(())
}
68 changes: 54 additions & 14 deletions codex-rs/codex-api/src/endpoint/realtime_websocket/methods.rs
Original file line number Diff line number Diff line change
Expand Up @@ -603,8 +603,16 @@ impl RealtimeWebsocketEvents {
}
RealtimeEvent::InputTranscriptDelta(RealtimeTranscriptDelta { delta, .. }) => {
let force_new = active_transcript.new_input_entry;
append_transcript_delta(&mut active_transcript.entries, "user", delta, force_new);
active_transcript.new_input_entry = false;
append_transcript_delta(
&mut active_transcript.entries,
"user",
delta,
force_new,
self.event_parser,
);
if !delta.is_empty() || self.event_parser != RealtimeEventParser::FramelessBidi {
active_transcript.new_input_entry = false;
}
}
RealtimeEvent::OutputTranscriptDelta(RealtimeTranscriptDelta { delta, .. }) => {
let force_new = active_transcript.new_output_entry;
Expand All @@ -613,8 +621,11 @@ impl RealtimeWebsocketEvents {
"assistant",
delta,
force_new,
self.event_parser,
);
active_transcript.new_output_entry = false;
if !delta.is_empty() || self.event_parser != RealtimeEventParser::FramelessBidi {
active_transcript.new_output_entry = false;
}
}
RealtimeEvent::InputTranscriptDone(done) => {
let force_new = active_transcript.new_input_entry;
Expand All @@ -623,8 +634,10 @@ impl RealtimeWebsocketEvents {
"user",
&done.text,
force_new,
self.event_parser,
);
active_transcript.new_input_entry = false;
active_transcript.new_input_entry =
self.event_parser == RealtimeEventParser::FramelessBidi;
}
RealtimeEvent::OutputTranscriptDone(done) => {
let force_new = active_transcript.new_output_entry;
Expand All @@ -633,8 +646,10 @@ impl RealtimeWebsocketEvents {
"assistant",
&done.text,
force_new,
self.event_parser,
);
active_transcript.new_output_entry = false;
active_transcript.new_output_entry =
self.event_parser == RealtimeEventParser::FramelessBidi;
}
RealtimeEvent::HandoffRequested(handoff) => {
append_handoff_input(&mut active_transcript.entries, &handoff.input_transcript);
Expand Down Expand Up @@ -704,15 +719,23 @@ fn append_transcript_delta(
role: &str,
delta: &str,
force_new: bool,
event_parser: RealtimeEventParser,
) {
if delta.is_empty() {
return;
}

if !force_new
&& let Some(last_entry) = entries.last_mut()
&& last_entry.role == role
{
// V3 interleaves the two speakers' transcripts. Legacy protocols retain
// adjacent speaker runs and their existing response lifecycle semantics.
let entry = match event_parser {
RealtimeEventParser::FramelessBidi => {
entries.iter_mut().rev().find(|entry| entry.role == role)
}
RealtimeEventParser::V1 | RealtimeEventParser::RealtimeV2 => {
entries.last_mut().filter(|entry| entry.role == role)
}
};
if !force_new && let Some(last_entry) = entry {
last_entry.text.push_str(delta);
return;
}
Expand All @@ -728,16 +751,29 @@ fn apply_transcript_done(
role: &str,
text: &str,
force_new: bool,
event_parser: RealtimeEventParser,
) {
if text.is_empty() {
return;
}

if !force_new
&& let Some(last_entry) = entries.last_mut()
&& last_entry.role == role
{
last_entry.text = text.to_string();
// V3 interleaves the two speakers' transcripts. Legacy protocols retain
// adjacent speaker runs and their existing response lifecycle semantics.
let entry = match event_parser {
RealtimeEventParser::FramelessBidi => {
entries.iter_mut().rev().find(|entry| entry.role == role)
}
RealtimeEventParser::V1 | RealtimeEventParser::RealtimeV2 => {
entries.last_mut().filter(|entry| entry.role == role)
}
};
if !force_new && let Some(last_entry) = entry {
// A delayed V3 final can arrive after deltas from the next utterance.
// Preserve accumulated speech unless the final extends it.
if event_parser != RealtimeEventParser::FramelessBidi || text.starts_with(&last_entry.text)
{
last_entry.text = text.to_string();
}
return;
}

Expand Down Expand Up @@ -1212,6 +1248,10 @@ fn normalize_realtime_path(url: &mut Url, event_parser: RealtimeEventParser) {
}
}

#[cfg(test)]
#[path = "transcript_tests.rs"]
mod transcript_tests;

#[cfg(test)]
mod tests {
use super::*;
Expand Down
131 changes: 131 additions & 0 deletions codex-rs/codex-api/src/endpoint/realtime_websocket/transcript_tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,131 @@
//! V3 transcript reconciliation across interleaved speakers, reconnects, and turn boundaries.

use super::*;
use codex_protocol::protocol::RealtimeTranscriptDone;
use pretty_assertions::assert_eq;

fn events(transcript_state: RealtimeTranscriptState) -> RealtimeWebsocketEvents {
let (_, rx_message) = async_channel::unbounded();
RealtimeWebsocketEvents {
rx_message,
pending_events: Arc::new(Mutex::new(VecDeque::new())),
transcript_state,
event_parser: RealtimeEventParser::FramelessBidi,
is_closed: Arc::new(AtomicBool::new(false)),
}
}

fn entry(role: &str, text: &str) -> RealtimeTranscriptEntry {
RealtimeTranscriptEntry {
role: role.to_string(),
text: text.to_string(),
}
}

async fn observe(events: &RealtimeWebsocketEvents, mut event: RealtimeEvent) {
let original = event.clone();
events.update_active_transcript(&mut event).await;
assert_eq!(
event, original,
"live transcript events must pass through unchanged"
);
}

#[tokio::test]
async fn interleaved_transcripts_reconcile_across_reconnect_in_either_completion_order() {
for assistant_finishes_first in [false, true] {
let state = RealtimeTranscriptState::default();
let first = events(state.clone());
for event in [
RealtimeEvent::InputTranscriptDelta(RealtimeTranscriptDelta {
delta: "What can ".into(),
}),
RealtimeEvent::OutputTranscriptDelta(RealtimeTranscriptDelta {
delta: "A bunch ".into(),
}),
RealtimeEvent::InputTranscriptDelta(RealtimeTranscriptDelta {
delta: "you do?".into(),
}),
RealtimeEvent::OutputTranscriptDelta(RealtimeTranscriptDelta {
delta: "of stuff".into(),
}),
] {
observe(&first, event).await;
}
let reconnected = events(state);
let user = RealtimeEvent::InputTranscriptDone(RealtimeTranscriptDone {
text: "What can you do?".into(),
});
let assistant = RealtimeEvent::OutputTranscriptDone(RealtimeTranscriptDone {
text: "A bunch of stuff, including coding.".into(),
});
let finals = if assistant_finishes_first {
[assistant, user]
} else {
[user, assistant]
};
for event in finals {
observe(&reconnected, event).await;
}
assert_eq!(
reconnected.take_transcript_tail().await,
vec![
entry("user", "What can you do?"),
entry("assistant", "A bunch of stuff, including coding."),
]
);
assert_eq!(reconnected.take_transcript_tail().await, vec![]);
}
}

#[tokio::test]
async fn completed_utterances_keep_repeated_speech_and_final_only_turns_distinct() {
let events = events(RealtimeTranscriptState::default());
for role in ["user", "assistant"] {
for with_delta in [true, true, false] {
if with_delta {
let delta = RealtimeTranscriptDelta {
delta: "Again".into(),
};
observe(
&events,
if role == "user" {
RealtimeEvent::InputTranscriptDelta(delta)
} else {
RealtimeEvent::OutputTranscriptDelta(delta)
},
)
.await;
}
let done = RealtimeTranscriptDone {
text: "Again.".into(),
};
observe(
&events,
if role == "user" {
RealtimeEvent::InputTranscriptDone(done)
} else {
RealtimeEvent::OutputTranscriptDone(done)
},
)
.await;
// An empty delta must not reopen the utterance we just completed.
let empty = RealtimeTranscriptDelta {
delta: String::new(),
};
observe(
&events,
if role == "user" {
RealtimeEvent::InputTranscriptDelta(empty)
} else {
RealtimeEvent::OutputTranscriptDelta(empty)
},
)
.await;
}
assert_eq!(
events.take_transcript_tail().await,
vec![entry(role, "Again."); 3]
);
}
}
Loading