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
20 changes: 15 additions & 5 deletions codex-rs/voice-host/src/audio_track.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
//! A single Opus RTP track. Muted capture sends generated silence to keep the peer alive;
//! elapsed time still advances its clock without inventing packet loss.
//! elapsed time advances in whole packets so capture jitter and mute transitions preserve
//! the receiver's 20 ms source grid without inventing packet loss.

use std::sync::Arc;
use std::time::Duration;
Expand Down Expand Up @@ -115,18 +116,28 @@ impl AudioTrack {

pub(crate) async fn send(&mut self, frame: EncodedAudio) -> Result<(), &'static str> {
tokio::time::timeout(SEND_TIMEOUT, async {
let mut gap = self
let duration = Duration::from_millis(/*millis*/ 20);
let elapsed = self
.end
.map(|end| frame.at.saturating_duration_since(end))
.unwrap_or_default();
// RTP counts audio samples, not callback wall-clock jitter. Skipping
// fractional packets changes the source grid and causes receivers
// that assemble fixed-duration frames to reject subsequent audio.
let skipped = Duration::from_millis(
(elapsed.as_millis() / duration.as_millis() * duration.as_millis()) as u64,
);
let mut gap = skipped;
while !gap.is_zero() {
let duration = gap.min(Duration::from_secs(/*secs*/ 3600));
self.track
.write_sample(
self.ssrc,
/*payload_type*/ OPUS_PAYLOAD_TYPE,
&Sample {
duration,
// write_sample truncates floating-point seconds to RTP ticks.
// Bias empty gaps by 1 ns to keep whole packets exact.
duration: duration + Duration::from_nanos(/*nanos*/ 1),
..Default::default()
},
&[],
Expand All @@ -135,7 +146,6 @@ impl AudioTrack {
.map_err(|_| "failed to advance voice clock")?;
gap -= duration;
}
let duration = Duration::from_millis(/*millis*/ 20);
self.track
.write_sample(
self.ssrc,
Expand All @@ -149,7 +159,7 @@ impl AudioTrack {
)
.await
.map_err(|_| "failed to send voice audio")?;
self.end = Some(self.end.unwrap_or(frame.at).max(frame.at) + duration);
self.end = Some(self.end.unwrap_or(frame.at) + skipped + duration);
Ok(())
})
.await
Expand Down
6 changes: 3 additions & 3 deletions codex-rs/voice-host/src/audio_track_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,8 @@ async fn timestamp_jitter_does_not_accumulate_but_mute_advances_the_clock() {
let at = start + Duration::from_millis(tick * 20 + tick % 2);
track.send(EncodedAudio { data: vec![], at }).await.unwrap();
}
assert_eq!(track.end, Some(start + Duration::from_millis(10_001)));
let at = start + Duration::from_secs(30);
assert_eq!(track.end, Some(start + Duration::from_millis(10_000)));
let at = start + Duration::from_millis(30_007);
track.send(EncodedAudio { data: vec![], at }).await.unwrap();
assert_eq!(track.end, Some(at + Duration::from_millis(20)));
assert_eq!(track.end, Some(start + Duration::from_millis(30_020)));
}
15 changes: 12 additions & 3 deletions codex-rs/voice-host/src/capture_worker_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -267,13 +267,20 @@ async fn capture_reaches_remote_rtp_and_mute_discards_queued_and_partial_audio()
worker.set_controls(muted()).unwrap();
worker.set_controls(unmuted()).unwrap();
let old_generation = worker.buffers.microphone.load(Ordering::Acquire);
now = Instant::now().max(start + Duration::from_millis(/*millis*/ 100));
now = start + Duration::from_millis(before.len() as u64 * 20 + 580);
enqueue(&worker, /*samples*/ 720, old_generation, now);
assert_eq!(worker.service(&mut local.audio, || now).await.unwrap(), 0);
enqueue(&worker, /*samples*/ 2_880, old_generation, now);
worker.set_controls(muted()).unwrap();
assert_eq!(worker.service(&mut local.audio, || now).await.unwrap(), 0);
let silence = received.recv().await.unwrap();
assert_eq!(
silence
.header
.timestamp
.wrapping_sub(before[0].header.timestamp),
before.len() as u32 * 960 + 27_840
);
let mut decoder = opus::Decoder::new(/*sample_rate*/ 48_000, opus::Channels::Mono).unwrap();
let mut decoded = [1.0; 960];
assert_eq!(
Expand Down Expand Up @@ -316,10 +323,12 @@ async fn capture_reaches_remote_rtp_and_mute_discards_queued_and_partial_audio()
.header
.timestamp
.wrapping_sub(before[0].header.timestamp);
assert_eq!(clock_gap % 960, 0);
let expected_gap = (resumed.duration_since(start).as_secs_f64() * 48_000.0) as u32;
// Allow 1 ms for sample rounding without imposing a wall-clock throughput requirement.
// Real gaps advance by whole packets, keeping wall time within one
// packet without imposing a wall-clock throughput requirement.
assert!(
clock_gap.abs_diff(expected_gap) <= 48,
clock_gap.abs_diff(expected_gap) <= 960,
"received RTP clock gap: {clock_gap}, expected {expected_gap}"
);
assert_eq!(
Expand Down
Loading