Reduce realtime self-interruptions during playback

Co-authored-by: Codex <noreply@openai.com>
This commit is contained in:
Ahmed Ibrahim
2026-03-15 15:19:39 -07:00
parent bc1a7cf2f4
commit 011795d705
3 changed files with 263 additions and 33 deletions

View File

@@ -5,6 +5,8 @@ use codex_protocol::protocol::RealtimeConversationClosedEvent;
use codex_protocol::protocol::RealtimeConversationRealtimeEvent;
use codex_protocol::protocol::RealtimeConversationStartedEvent;
use codex_protocol::protocol::RealtimeEvent;
#[cfg(not(target_os = "linux"))]
use std::sync::atomic::AtomicUsize;
const REALTIME_CONVERSATION_PROMPT: &str = "You are in a realtime voice conversation in the Codex TUI. Respond conversationally and concisely.";
@@ -30,6 +32,8 @@ pub(super) struct RealtimeConversationUiState {
capture: Option<crate::voice::VoiceCapture>,
#[cfg(not(target_os = "linux"))]
audio_player: Option<crate::voice::RealtimeAudioPlayer>,
#[cfg(not(target_os = "linux"))]
playback_queued_samples: Arc<AtomicUsize>,
}
impl RealtimeConversationUiState {
@@ -293,8 +297,11 @@ impl ChatWidget {
#[cfg(not(target_os = "linux"))]
{
if self.realtime_conversation.audio_player.is_none() {
self.realtime_conversation.audio_player =
crate::voice::RealtimeAudioPlayer::start(&self.config).ok();
self.realtime_conversation.audio_player = crate::voice::RealtimeAudioPlayer::start(
&self.config,
Arc::clone(&self.realtime_conversation.playback_queued_samples),
)
.ok();
}
if let Some(player) = &self.realtime_conversation.audio_player
&& let Err(err) = player.enqueue_frame(frame)
@@ -331,6 +338,7 @@ impl ChatWidget {
let capture = match crate::voice::VoiceCapture::start_realtime(
&self.config,
self.app_event_tx.clone(),
Arc::clone(&self.realtime_conversation.playback_queued_samples),
) {
Ok(capture) => capture,
Err(err) => {
@@ -349,8 +357,11 @@ impl ChatWidget {
self.realtime_conversation.capture_stop_flag = Some(stop_flag.clone());
self.realtime_conversation.capture = Some(capture);
if self.realtime_conversation.audio_player.is_none() {
self.realtime_conversation.audio_player =
crate::voice::RealtimeAudioPlayer::start(&self.config).ok();
self.realtime_conversation.audio_player = crate::voice::RealtimeAudioPlayer::start(
&self.config,
Arc::clone(&self.realtime_conversation.playback_queued_samples),
)
.ok();
}
std::thread::spawn(move || {
@@ -388,7 +399,10 @@ impl ChatWidget {
}
RealtimeAudioDeviceKind::Speaker => {
self.stop_realtime_speaker();
match crate::voice::RealtimeAudioPlayer::start(&self.config) {
match crate::voice::RealtimeAudioPlayer::start(
&self.config,
Arc::clone(&self.realtime_conversation.playback_queued_samples),
) {
Ok(player) => {
self.realtime_conversation.audio_player = Some(player);
}

View File

@@ -139,6 +139,7 @@ mod voice {
use std::sync::Mutex;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicU16;
use std::sync::atomic::AtomicUsize;
pub struct RecordedAudio {
pub data: Vec<i16>,
@@ -157,7 +158,11 @@ mod voice {
Err("voice input is unavailable in this build".to_string())
}
pub fn start_realtime(_config: &Config, _tx: AppEventSender) -> Result<Self, String> {
pub fn start_realtime(
_config: &Config,
_tx: AppEventSender,
_playback_queued_samples: Arc<AtomicUsize>,
) -> Result<Self, String> {
Err("voice input is unavailable in this build".to_string())
}
@@ -197,7 +202,10 @@ mod voice {
}
impl RealtimeAudioPlayer {
pub(crate) fn start(_config: &Config) -> Result<Self, String> {
pub(crate) fn start(
_config: &Config,
_queued_samples: Arc<AtomicUsize>,
) -> Result<Self, String> {
Err("voice output is unavailable in this build".to_string())
}

View File

@@ -23,7 +23,10 @@ use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::AtomicU16;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::time::Duration;
use std::time::Instant;
use tracing::error;
use tracing::info;
use tracing::trace;
@@ -31,6 +34,8 @@ use tracing::trace;
const AUDIO_MODEL: &str = "gpt-4o-mini-transcribe";
const MODEL_AUDIO_SAMPLE_RATE: u32 = 24_000;
const MODEL_AUDIO_CHANNELS: u16 = 1;
const REALTIME_INTERRUPT_INPUT_PEAK_THRESHOLD: u16 = 4_000;
const REALTIME_INTERRUPT_GRACE_PERIOD: Duration = Duration::from_millis(900);
struct TranscriptionAuthContext {
mode: AuthMode,
@@ -79,7 +84,11 @@ impl VoiceCapture {
})
}
pub fn start_realtime(config: &Config, tx: AppEventSender) -> Result<Self, String> {
pub fn start_realtime(
config: &Config,
tx: AppEventSender,
playback_queued_samples: Arc<AtomicUsize>,
) -> Result<Self, String> {
let (device, config) = select_realtime_input_device_and_config(config)?;
let sample_rate = config.sample_rate().0;
@@ -94,6 +103,7 @@ impl VoiceCapture {
sample_rate,
channels,
tx,
playback_queued_samples,
last_peak.clone(),
)?;
stream
@@ -343,17 +353,33 @@ fn build_realtime_input_stream(
sample_rate: u32,
channels: u16,
tx: AppEventSender,
playback_queued_samples: Arc<AtomicUsize>,
last_peak: Arc<AtomicU16>,
) -> Result<cpal::Stream, String> {
match config.sample_format() {
cpal::SampleFormat::F32 => device
.build_input_stream(
&config.clone().into(),
move |input: &[f32], _| {
let peak = peak_f32(input);
last_peak.store(peak, Ordering::Relaxed);
let samples = input.iter().copied().map(f32_to_i16).collect::<Vec<_>>();
send_realtime_audio_chunk(&tx, samples, sample_rate, channels);
{
let playback_queued_samples = Arc::clone(&playback_queued_samples);
let last_peak = Arc::clone(&last_peak);
let tx = tx;
let mut allow_input_until = None;
move |input: &[f32], _| {
let peak = peak_f32(input);
if !should_send_realtime_input(
peak,
&playback_queued_samples,
&mut allow_input_until,
Instant::now(),
) {
last_peak.store(0, Ordering::Relaxed);
return;
}
last_peak.store(peak, Ordering::Relaxed);
let samples = input.iter().copied().map(f32_to_i16).collect::<Vec<_>>();
send_realtime_audio_chunk(&tx, samples, sample_rate, channels);
}
},
move |err| error!("audio input error: {err}"),
None,
@@ -362,10 +388,25 @@ fn build_realtime_input_stream(
cpal::SampleFormat::I16 => device
.build_input_stream(
&config.clone().into(),
move |input: &[i16], _| {
let peak = peak_i16(input);
last_peak.store(peak, Ordering::Relaxed);
send_realtime_audio_chunk(&tx, input.to_vec(), sample_rate, channels);
{
let playback_queued_samples = Arc::clone(&playback_queued_samples);
let last_peak = Arc::clone(&last_peak);
let tx = tx;
let mut allow_input_until = None;
move |input: &[i16], _| {
let peak = peak_i16(input);
if !should_send_realtime_input(
peak,
&playback_queued_samples,
&mut allow_input_until,
Instant::now(),
) {
last_peak.store(0, Ordering::Relaxed);
return;
}
last_peak.store(peak, Ordering::Relaxed);
send_realtime_audio_chunk(&tx, input.to_vec(), sample_rate, channels);
}
},
move |err| error!("audio input error: {err}"),
None,
@@ -374,11 +415,26 @@ fn build_realtime_input_stream(
cpal::SampleFormat::U16 => device
.build_input_stream(
&config.clone().into(),
move |input: &[u16], _| {
let mut samples = Vec::with_capacity(input.len());
let peak = convert_u16_to_i16_and_peak(input, &mut samples);
last_peak.store(peak, Ordering::Relaxed);
send_realtime_audio_chunk(&tx, samples, sample_rate, channels);
{
let playback_queued_samples = Arc::clone(&playback_queued_samples);
let last_peak = Arc::clone(&last_peak);
let tx = tx;
let mut allow_input_until = None;
move |input: &[u16], _| {
let mut samples = Vec::with_capacity(input.len());
let peak = convert_u16_to_i16_and_peak(input, &mut samples);
if !should_send_realtime_input(
peak,
&playback_queued_samples,
&mut allow_input_until,
Instant::now(),
) {
last_peak.store(0, Ordering::Relaxed);
return;
}
last_peak.store(peak, Ordering::Relaxed);
send_realtime_audio_chunk(&tx, samples, sample_rate, channels);
}
},
move |err| error!("audio input error: {err}"),
None,
@@ -487,24 +543,31 @@ fn convert_u16_to_i16_and_peak(input: &[u16], out: &mut Vec<i16>) -> u16 {
pub(crate) struct RealtimeAudioPlayer {
_stream: cpal::Stream,
queue: Arc<Mutex<VecDeque<i16>>>,
queued_samples: Arc<AtomicUsize>,
output_sample_rate: u32,
output_channels: u16,
}
impl RealtimeAudioPlayer {
pub(crate) fn start(config: &Config) -> Result<Self, String> {
pub(crate) fn start(config: &Config, queued_samples: Arc<AtomicUsize>) -> Result<Self, String> {
let (device, config) =
crate::audio_device::select_configured_output_device_and_config(config)?;
let output_sample_rate = config.sample_rate().0;
let output_channels = config.channels();
let queue = Arc::new(Mutex::new(VecDeque::new()));
let stream = build_output_stream(&device, &config, Arc::clone(&queue))?;
let stream = build_output_stream(
&device,
&config,
Arc::clone(&queue),
Arc::clone(&queued_samples),
)?;
stream
.play()
.map_err(|e| format!("failed to start output stream: {e}"))?;
Ok(Self {
_stream: stream,
queue,
queued_samples,
output_sample_rate,
output_channels,
})
@@ -539,12 +602,15 @@ impl RealtimeAudioPlayer {
.lock()
.map_err(|_| "failed to lock output audio queue".to_string())?;
// TODO(aibrahim): Cap or trim this queue if we observe producer bursts outrunning playback.
self.queued_samples
.fetch_add(converted.len(), Ordering::Relaxed);
guard.extend(converted);
Ok(())
}
pub(crate) fn clear(&self) {
if let Ok(mut guard) = self.queue.lock() {
self.queued_samples.store(0, Ordering::Relaxed);
guard.clear();
}
}
@@ -554,13 +620,14 @@ fn build_output_stream(
device: &cpal::Device,
config: &cpal::SupportedStreamConfig,
queue: Arc<Mutex<VecDeque<i16>>>,
queued_samples: Arc<AtomicUsize>,
) -> Result<cpal::Stream, String> {
let config_any: cpal::StreamConfig = config.clone().into();
match config.sample_format() {
cpal::SampleFormat::F32 => device
.build_output_stream(
&config_any,
move |output: &mut [f32], _| fill_output_f32(output, &queue),
move |output: &mut [f32], _| fill_output_f32(output, &queue, &queued_samples),
move |err| error!("audio output error: {err}"),
None,
)
@@ -568,7 +635,7 @@ fn build_output_stream(
cpal::SampleFormat::I16 => device
.build_output_stream(
&config_any,
move |output: &mut [i16], _| fill_output_i16(output, &queue),
move |output: &mut [i16], _| fill_output_i16(output, &queue, &queued_samples),
move |err| error!("audio output error: {err}"),
None,
)
@@ -576,7 +643,7 @@ fn build_output_stream(
cpal::SampleFormat::U16 => device
.build_output_stream(
&config_any,
move |output: &mut [u16], _| fill_output_u16(output, &queue),
move |output: &mut [u16], _| fill_output_u16(output, &queue, &queued_samples),
move |err| error!("audio output error: {err}"),
None,
)
@@ -585,38 +652,109 @@ fn build_output_stream(
}
}
fn fill_output_i16(output: &mut [i16], queue: &Arc<Mutex<VecDeque<i16>>>) {
fn fill_output_i16(
output: &mut [i16],
queue: &Arc<Mutex<VecDeque<i16>>>,
queued_samples: &Arc<AtomicUsize>,
) {
if let Ok(mut guard) = queue.lock() {
let mut consumed = 0usize;
for sample in output {
*sample = guard.pop_front().unwrap_or(0);
*sample = if let Some(next) = guard.pop_front() {
consumed += 1;
next
} else {
0
};
}
if consumed > 0 {
consume_output_samples(queued_samples, consumed);
}
return;
}
output.fill(0);
}
fn fill_output_f32(output: &mut [f32], queue: &Arc<Mutex<VecDeque<i16>>>) {
fn fill_output_f32(
output: &mut [f32],
queue: &Arc<Mutex<VecDeque<i16>>>,
queued_samples: &Arc<AtomicUsize>,
) {
if let Ok(mut guard) = queue.lock() {
let mut consumed = 0usize;
for sample in output {
let v = guard.pop_front().unwrap_or(0);
let v = if let Some(next) = guard.pop_front() {
consumed += 1;
next
} else {
0
};
*sample = (v as f32) / (i16::MAX as f32);
}
if consumed > 0 {
consume_output_samples(queued_samples, consumed);
}
return;
}
output.fill(0.0);
}
fn fill_output_u16(output: &mut [u16], queue: &Arc<Mutex<VecDeque<i16>>>) {
fn fill_output_u16(
output: &mut [u16],
queue: &Arc<Mutex<VecDeque<i16>>>,
queued_samples: &Arc<AtomicUsize>,
) {
if let Ok(mut guard) = queue.lock() {
let mut consumed = 0usize;
for sample in output {
let v = guard.pop_front().unwrap_or(0);
let v = if let Some(next) = guard.pop_front() {
consumed += 1;
next
} else {
0
};
*sample = (v as i32 + 32768).clamp(0, u16::MAX as i32) as u16;
}
if consumed > 0 {
consume_output_samples(queued_samples, consumed);
}
return;
}
output.fill(32768);
}
fn consume_output_samples(queued_samples: &Arc<AtomicUsize>, consumed: usize) {
let _ = queued_samples.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |queued| {
Some(queued.saturating_sub(consumed))
});
}
fn should_send_realtime_input(
peak: u16,
playback_queued_samples: &Arc<AtomicUsize>,
allow_input_until: &mut Option<Instant>,
now: Instant,
) -> bool {
if playback_queued_samples.load(Ordering::Relaxed) == 0 {
*allow_input_until = None;
return true;
}
if let Some(deadline) = *allow_input_until {
if now < deadline {
return true;
}
*allow_input_until = None;
}
if peak >= REALTIME_INTERRUPT_INPUT_PEAK_THRESHOLD {
*allow_input_until = Some(now + REALTIME_INTERRUPT_GRACE_PERIOD);
return true;
}
false
}
fn convert_pcm16(
input: &[i16],
input_sample_rate: u32,
@@ -878,8 +1016,14 @@ mod tests {
use super::RecordedAudio;
use super::convert_pcm16;
use super::encode_wav_normalized;
use super::should_send_realtime_input;
use pretty_assertions::assert_eq;
use std::io::Cursor;
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering;
use std::time::Duration;
use std::time::Instant;
#[test]
fn convert_pcm16_downmixes_and_resamples_for_model_input() {
@@ -908,4 +1052,68 @@ mod tests {
assert_eq!(spec.sample_rate, 24_000);
assert_eq!(samples, vec![8_426, 29_490]);
}
#[test]
fn realtime_input_suppresses_low_peaks_while_playback_is_active() {
let playback_queued_samples = Arc::new(AtomicUsize::new(128));
let mut allow_input_until = None;
let should_send = should_send_realtime_input(
500,
&playback_queued_samples,
&mut allow_input_until,
Instant::now(),
);
assert!(!should_send);
assert_eq!(allow_input_until, None);
}
#[test]
fn realtime_input_allows_barge_in_after_loud_peak() {
let playback_queued_samples = Arc::new(AtomicUsize::new(128));
let now = Instant::now();
let mut allow_input_until = None;
let first_chunk = should_send_realtime_input(
6_000,
&playback_queued_samples,
&mut allow_input_until,
now,
);
let follow_up_chunk = should_send_realtime_input(
200,
&playback_queued_samples,
&mut allow_input_until,
now + Duration::from_millis(200),
);
assert!(first_chunk);
assert!(follow_up_chunk);
assert!(allow_input_until.is_some());
}
#[test]
fn realtime_input_clears_barge_in_window_after_playback_stops() {
let playback_queued_samples = Arc::new(AtomicUsize::new(128));
let now = Instant::now();
let mut allow_input_until = None;
let _ = should_send_realtime_input(
6_000,
&playback_queued_samples,
&mut allow_input_until,
now,
);
playback_queued_samples.store(0, Ordering::Relaxed);
let should_send = should_send_realtime_input(
100,
&playback_queued_samples,
&mut allow_input_until,
now + Duration::from_secs(1),
);
assert!(should_send);
assert_eq!(allow_input_until, None);
}
}