//! Per-peer jitter buffer with Opus packet-loss concealment. //! //! Incoming audio arrives as unreliable QUIC datagrams that can be reordered, //! duplicated, or dropped on real networks. Each packet carries a monotonic //! sequence number (assigned by the sender). This buffer reorders packets by //! sequence, holds a small fixed playout delay to absorb jitter, and — when a //! sequence is missing but later packets have already arrived — synthesizes a //! concealment frame via Opus PLC instead of emitting a click of silence. use crate::codec::{AudioDecoder, CodecError, opus_impl::OpusDecoder}; use opus::Channels; use std::collections::BTreeMap; /// Samples per channel in one transmitted frame (20ms @ 48kHz mono). pub const FRAME_SAMPLES: usize = 960; /// How many frames to buffer before playout begins (~60ms). This is the /// tolerance window for reordering and jitter; larger = more resilient but /// more latency. const TARGET_DELAY_FRAMES: usize = 3; /// Hard cap on buffered frames (~640ms). If we ever exceed this we've fallen /// badly behind, so we drop the oldest and resync rather than grow unbounded. const MAX_BUFFERED_FRAMES: usize = 32; pub struct JitterBuffer { decoder: OpusDecoder, /// Reorder window: sequence number -> encoded Opus payload. packets: BTreeMap>, /// Next sequence we expect to play. `None` means idle/buffering: we are /// waiting to accumulate `TARGET_DELAY_FRAMES` before (re)starting playout. next_seq: Option, } /// Wrapping-aware "is `a` strictly before `b`" for sequence numbers. fn seq_before(a: u32, b: u32) -> bool { a != b && b.wrapping_sub(a) < (1 << 31) } impl JitterBuffer { pub fn new() -> Result { Ok(Self { decoder: OpusDecoder::new(48000, Channels::Mono, FRAME_SAMPLES)?, packets: BTreeMap::new(), next_seq: None, }) } /// Store a received packet, dropping ones we've already played past and /// bounding total depth. pub fn insert(&mut self, seq: u32, payload: Vec) { // Too late: this sequence has already been played (or concealed). if let Some(next) = self.next_seq && seq_before(seq, next) { return; } self.packets.insert(seq, payload); while self.packets.len() > MAX_BUFFERED_FRAMES { let oldest = *self.packets.keys().next().expect("non-empty"); self.packets.remove(&oldest); // We've discarded backlog; resync the playout head to the new front. self.next_seq = self.packets.keys().next().copied(); } } /// Produce the next 20ms PCM frame for playout, or `None` when idle or /// still buffering (the caller should treat `None` as silence). pub fn pop_frame(&mut self) -> Option> { match self.next_seq { None => { // Buffering: start playout once we have enough to absorb jitter. if self.packets.len() >= TARGET_DELAY_FRAMES { self.next_seq = self.packets.keys().next().copied(); self.pop_frame() } else { None } } Some(next) => { if let Some(payload) = self.packets.remove(&next) { self.next_seq = Some(next.wrapping_add(1)); self.decoder.decode(Some(&payload)).ok() } else if self.packets.is_empty() { // Underrun: the talker has gone quiet (or stopped). Go idle // and re-buffer before resuming, rather than concealing forever. self.next_seq = None; None } else { // Gap with later packets already buffered: a packet was lost // or reordered out of window. Conceal this frame via Opus PLC. self.next_seq = Some(next.wrapping_add(1)); self.decoder.decode(None).ok() } } } } /// True when nothing is buffered and playout is idle (talker silent). pub fn is_idle(&self) -> bool { self.next_seq.is_none() && self.packets.is_empty() } } #[cfg(test)] mod tests { use super::*; use crate::codec::AudioEncoder; use crate::codec::opus_impl::OpusEncoder; use opus::{Application, Channels}; /// A real, decodable Opus packet for one 20ms mono frame at amplitude `amp`. fn frame(enc: &mut OpusEncoder, amp: i16) -> Vec { let pcm: Vec = (0..FRAME_SAMPLES) .map(|i| if i % 2 == 0 { amp } else { -amp }) .collect(); enc.encode(&pcm).unwrap() } #[test] fn buffers_then_plays_in_order() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Below the target delay, playout hasn't primed yet. jb.insert(0, frame(&mut enc, 1000)); assert!(jb.pop_frame().is_none()); // Reaching the target delay primes playout and yields the first frame. jb.insert(1, frame(&mut enc, 1000)); jb.insert(2, frame(&mut enc, 1000)); assert_eq!(jb.pop_frame().map(|f| f.len()), Some(FRAME_SAMPLES)); assert_eq!(jb.pop_frame().map(|f| f.len()), Some(FRAME_SAMPLES)); assert_eq!(jb.pop_frame().map(|f| f.len()), Some(FRAME_SAMPLES)); // Drained: idle again. assert!(jb.pop_frame().is_none()); assert!(jb.is_idle()); } #[test] fn reorders_out_of_order_arrivals() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Arrive scrambled but within the buffering window. jb.insert(2, frame(&mut enc, 800)); jb.insert(0, frame(&mut enc, 800)); jb.insert(1, frame(&mut enc, 800)); // Three real frames come out (in sequence order), then idle. assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_none()); } #[test] fn conceals_gap_when_later_packets_present() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Seq 2 is missing, but 0,1,3 arrive — enough to prime. jb.insert(0, frame(&mut enc, 1200)); jb.insert(1, frame(&mut enc, 1200)); jb.insert(3, frame(&mut enc, 1200)); assert!(jb.pop_frame().is_some()); // seq 0 assert!(jb.pop_frame().is_some()); // seq 1 // seq 2 missing but seq 3 buffered -> Opus PLC produces a concealment frame. let concealed = jb.pop_frame(); assert_eq!(concealed.map(|f| f.len()), Some(FRAME_SAMPLES)); assert!(jb.pop_frame().is_some()); // seq 3 assert!(jb.pop_frame().is_none()); } #[test] fn drops_packets_already_played() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); jb.insert(5, frame(&mut enc, 600)); jb.insert(6, frame(&mut enc, 600)); jb.insert(7, frame(&mut enc, 600)); assert!(jb.pop_frame().is_some()); // primes at seq 5, plays 5 assert!(jb.pop_frame().is_some()); // 6 // A straggler for an already-played sequence must be discarded. jb.insert(5, frame(&mut enc, 600)); assert_eq!(jb.packets.len(), 1); // only seq 7 remains buffered } #[test] fn test_seq_before_ordering() { // Basic ordering assert!(seq_before(0, 1)); assert!(!seq_before(1, 0)); assert!(!seq_before(5, 5)); assert!(seq_before(100, 101)); assert!(!seq_before(101, 100)); // Wraparound assert!(seq_before(u32::MAX, 0)); assert!(!seq_before(0, u32::MAX)); // Half-range boundary assert!(seq_before(0, 0x7FFF_FFFF)); assert!(!seq_before(0, 0x8000_0000)); } #[test] fn test_jitter_buffer_overflow_resync() { let mut jb = JitterBuffer::new().unwrap(); assert!(jb.next_seq.is_none()); let count = MAX_BUFFERED_FRAMES + 1; for seq in 0..count { jb.insert(seq as u32, vec![0u8]); } assert_eq!(jb.packets.len(), MAX_BUFFERED_FRAMES); assert!(!jb.packets.contains_key(&0)); assert_eq!(jb.next_seq, Some(1)); } #[test] fn reprimes_after_underrun_idle() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Prime with seq 0, 1, 2 jb.insert(0, frame(&mut enc, 1000)); jb.insert(1, frame(&mut enc, 1000)); jb.insert(2, frame(&mut enc, 1000)); // pop_frame() 3x -> 3 Some assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); // A 4th pop_frame() -> None, and jb.is_idle() is true assert!(jb.pop_frame().is_none()); assert!(jb.is_idle()); // Insert ONE new frame (seq 3): pop_frame() must still be None (must re-accumulate TARGET_DELAY_FRAMES) jb.insert(3, frame(&mut enc, 1000)); assert!(jb.pop_frame().is_none()); assert!(!jb.is_idle()); // Insert seq 4 and 5 (now 3 buffered) -> pop_frame() yields Some (re-primed) jb.insert(4, frame(&mut enc, 1000)); jb.insert(5, frame(&mut enc, 1000)); assert!(jb.pop_frame().is_some()); } #[test] fn duplicate_insert_does_not_grow_buffer() { let mut jb = JitterBuffer::new().unwrap(); jb.insert(0, vec![0u8]); jb.insert(0, vec![1u8]); assert_eq!(jb.packets.len(), 1); } #[test] fn is_idle_reflects_buffer_state() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Fresh buffer assert!(jb.is_idle()); // After a single insert jb.insert(0, frame(&mut enc, 1000)); assert!(!jb.is_idle()); // Prime (3 frames) jb.insert(1, frame(&mut enc, 1000)); jb.insert(2, frame(&mut enc, 1000)); // Drain past the end so it underruns assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_none()); assert!(jb.is_idle()); } #[test] fn overflow_resync_while_playing() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Prime with seq 0, 1, 2 jb.insert(0, frame(&mut enc, 1000)); jb.insert(1, frame(&mut enc, 1000)); jb.insert(2, frame(&mut enc, 1000)); // pop_frame() twice (now next_seq == Some(2)) assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert_eq!(jb.next_seq, Some(2)); // Insert a contiguous run of higher sequences to exceed MAX_BUFFERED_FRAMES let start = 3; let end = 3 + MAX_BUFFERED_FRAMES + 2; for seq in start..end { jb.insert(seq as u32, frame(&mut enc, 1000)); } assert_eq!(jb.packets.len(), MAX_BUFFERED_FRAMES); assert_eq!(jb.next_seq, jb.packets.keys().next().copied()); } }