//! 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 an *adaptive* 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. //! //! ## Adaptive playout delay //! //! The playout delay (how many frames we accumulate before (re)starting //! playout) is a feedback controller, not a fixed constant. It reacts to the //! buffer's own observations, with no wall clock required: //! //! * **Grow** (jitter beat the cushion): a late-arriving packet (one for a //! sequence we already played past) or a gap that forced Opus PLC each bumps //! the target up one frame. //! * **Shrink** (comfortably ahead): a long unbroken run of real decoded frames //! shaves the target back down one frame. //! //! Growth is fast and shrink is slow (AIMD-style) so we react to badness //! immediately but reclaim latency cautiously. Benign silence — a talker //! pausing, so packets simply stop — produces none of these signals, so the //! target is left untouched across quiet stretches. 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; /// Starting (and most common) playout delay (~60ms): the number of frames to /// accumulate before playout begins. The adaptive controller moves the live /// target up and down from here within `[MIN_DELAY_FRAMES, MAX_DELAY_FRAMES]`. const DEFAULT_DELAY_FRAMES: usize = 3; /// Floor for the adaptive delay (~40ms). Below this there's no slack left to /// reorder even a single packet, so we never shrink past it. const MIN_DELAY_FRAMES: usize = 2; /// Ceiling for the adaptive delay (~240ms). Kept well under /// `MAX_BUFFERED_FRAMES` so a deep cushion still leaves reorder headroom, and /// bounded so a pathological link can't drive playout latency unboundedly. const MAX_DELAY_FRAMES: usize = 12; /// Consecutive cleanly-played real frames (~5s) required to shave one frame off /// the target. Deliberately long so we reclaim latency slowly and don't flap. const CLEAN_RUN_TO_SHRINK: usize = 250; /// While buffering, prime playout after this many `pop_frame` polls even if the /// target delay isn't met yet (~500ms; the mixer polls every 20ms). This /// rescues a short utterance that never reaches a grown target, and bounds the /// worst-case startup latency. const PRIME_TIMEOUT_TICKS: usize = 25; /// 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; /// Sequence discontinuities larger than this (~10s at 20ms/frame) are treated /// as a restarted/new stream, not ordinary packet loss or reordering. const MAX_REASONABLE_SEQ_GAP: u32 = 500; 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, /// Live adaptive playout delay, in frames. Moved by the controller within /// `[MIN_DELAY_FRAMES, MAX_DELAY_FRAMES]`. target_delay: usize, /// Consecutive cleanly-played real frames since the last disruption; drives /// the slow shrink toward `MIN_DELAY_FRAMES`. clean_run: usize, /// `pop_frame` polls spent buffering with packets present; drives the /// `PRIME_TIMEOUT_TICKS` safety prime. buffering_ticks: usize, } /// 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, target_delay: DEFAULT_DELAY_FRAMES, clean_run: 0, buffering_ticks: 0, }) } /// Current adaptive playout delay, in frames. Exposed for metrics/tests. pub fn target_delay(&self) -> usize { self.target_delay } /// Grow the playout delay one frame (bounded): jitter beat the cushion, so /// next time we (re)prime we hold a deeper buffer. Resets the clean run. fn note_disruption(&mut self) { self.target_delay = (self.target_delay + 1).min(MAX_DELAY_FRAMES); self.clean_run = 0; } /// Count one cleanly-played real frame; after a long unbroken run, shave one /// frame off the delay (bounded below) to reclaim latency on a calm link. fn note_clean(&mut self) { self.clean_run += 1; if self.clean_run >= CLEAN_RUN_TO_SHRINK { self.target_delay = self.target_delay.saturating_sub(1).max(MIN_DELAY_FRAMES); self.clean_run = 0; } } fn reset_to_stream(&mut self, seq: u32, payload: Vec) { self.packets.clear(); self.packets.insert(seq, payload); self.next_seq = None; self.clean_run = 0; self.buffering_ticks = 0; } /// 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). Its // arrival after the playout head means our cushion was too shallow. if let Some(next) = self.next_seq && seq_before(seq, next) { if next.wrapping_sub(seq) > MAX_REASONABLE_SEQ_GAP { self.reset_to_stream(seq, payload); return; } self.note_disruption(); return; } if let Some(next) = self.next_seq && seq.wrapping_sub(next) > MAX_REASONABLE_SEQ_GAP { self.reset_to_stream(seq, payload); 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(); // The resync breaks sequence continuity; restart the clean run. self.clean_run = 0; } } /// 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 => { // Idle with nothing buffered: genuinely silent, no prime pending. if self.packets.is_empty() { self.buffering_ticks = 0; return None; } // Buffering: prime once we've accumulated the adaptive target, or // after a bounded wait so a short utterance isn't held forever. self.buffering_ticks += 1; if self.packets.len() >= self.target_delay || self.buffering_ticks >= PRIME_TIMEOUT_TICKS { self.buffering_ticks = 0; self.clean_run = 0; 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)); let frame = self.decoder.decode(Some(&payload)).ok(); // A real, in-order frame played: the link is keeping up. self.note_clean(); frame } 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. // This is benign (silence), so we don't grow the delay; just // end the clean run since playout is breaking. self.next_seq = None; self.clean_run = 0; None } else { // Gap with later packets already buffered: a packet was lost // or reordered out of window. Conceal this frame via Opus PLC // and grow the cushion — the jitter beat our current delay. self.next_seq = Some(next.wrapping_add(1)); self.note_disruption(); 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 far_behind_sequence_resets_as_restarted_stream() { let mut jb = JitterBuffer::new().unwrap(); jb.next_seq = Some(5_000); jb.packets.insert(5_000, vec![9]); jb.clean_run = 12; jb.buffering_ticks = 4; jb.insert(0, vec![1]); assert_eq!(jb.next_seq, None); assert_eq!(jb.packets.len(), 1); assert_eq!(jb.packets.get(&0).map(Vec::as_slice), Some(&[1][..])); assert_eq!(jb.clean_run, 0); assert_eq!(jb.buffering_ticks, 0); } #[test] fn far_ahead_sequence_resets_to_bound_plc_run() { let mut jb = JitterBuffer::new().unwrap(); jb.next_seq = Some(10); jb.packets.insert(10, vec![9]); jb.clean_run = 12; jb.buffering_ticks = 4; let jumped_seq = 10 + MAX_REASONABLE_SEQ_GAP + 1; jb.insert(jumped_seq, vec![2]); assert_eq!(jb.next_seq, None); assert_eq!(jb.packets.len(), 1); assert_eq!( jb.packets.get(&jumped_seq).map(Vec::as_slice), Some(&[2][..]) ); assert_eq!(jb.clean_run, 0); assert_eq!(jb.buffering_ticks, 0); } #[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()); } // ---- Adaptive playout delay ------------------------------------------ #[test] fn starts_at_default_delay() { let jb = JitterBuffer::new().unwrap(); assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES); } #[test] fn grows_delay_on_late_arrival() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Prime and play two frames so the playout head sits at seq 2. jb.insert(0, frame(&mut enc, 1000)); jb.insert(1, frame(&mut enc, 1000)); jb.insert(2, frame(&mut enc, 1000)); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert_eq!(jb.next_seq, Some(2)); assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES); // A packet for an already-played sequence arrives too late: grow by one. jb.insert(0, vec![0u8]); assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES + 1); // The stale payload was dropped, not buffered. assert!(!jb.packets.contains_key(&0)); } #[test] fn grows_delay_on_gap_conceal() { 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 (clean) assert!(jb.pop_frame().is_some()); // seq 1 (clean) assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES); // Seq 2 missing with seq 3 buffered -> PLC conceal -> grow by one. assert!(jb.pop_frame().is_some()); assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES + 1); } #[test] fn silence_does_not_change_delay() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // A clean short utterance that drains to an underrun (talker stops). jb.insert(0, frame(&mut enc, 1000)); jb.insert(1, frame(&mut enc, 1000)); jb.insert(2, frame(&mut enc, 1000)); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); // Underrun + further idle polls must leave the delay untouched: a quiet // talker is not a network problem. assert!(jb.pop_frame().is_none()); assert!(jb.is_idle()); assert!(jb.pop_frame().is_none()); assert!(jb.pop_frame().is_none()); assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES); } #[test] fn shrinks_delay_after_clean_run() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Steady state: keep the cushion topped up so every pop yields a real, // in-order frame (no conceal, no underrun, no overflow). let mut next = 0u32; for _ in 0..DEFAULT_DELAY_FRAMES { jb.insert(next, frame(&mut enc, 800)); next += 1; } for _ in 0..CLEAN_RUN_TO_SHRINK { assert!(jb.pop_frame().is_some()); jb.insert(next, frame(&mut enc, 800)); next += 1; } // One clean run's worth of frames shaves exactly one off the delay. assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES - 1); } #[test] fn delay_is_bounded_above() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Prime and advance the head, then hammer late arrivals. jb.insert(0, frame(&mut enc, 500)); jb.insert(1, frame(&mut enc, 500)); jb.insert(2, frame(&mut enc, 500)); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); for _ in 0..100 { jb.insert(0, vec![0u8]); // always "too late" -> disruption } assert_eq!(jb.target_delay(), MAX_DELAY_FRAMES); } #[test] fn delay_is_bounded_below() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Many clean runs would shrink forever; it must stop at the floor. let mut next = 0u32; for _ in 0..DEFAULT_DELAY_FRAMES { jb.insert(next, frame(&mut enc, 800)); next += 1; } for _ in 0..(CLEAN_RUN_TO_SHRINK * 4) { assert!(jb.pop_frame().is_some()); jb.insert(next, frame(&mut enc, 800)); next += 1; assert!(jb.target_delay() >= MIN_DELAY_FRAMES); } assert_eq!(jb.target_delay(), MIN_DELAY_FRAMES); } #[test] fn prime_timeout_rescues_short_utterance() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); // Drive the target above what a short utterance can reach. jb.insert(0, frame(&mut enc, 500)); jb.insert(1, frame(&mut enc, 500)); jb.insert(2, frame(&mut enc, 500)); assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); while jb.target_delay() < 6 { jb.insert(0, vec![0u8]); } assert!(jb.pop_frame().is_some()); // drain seq 2 assert!(jb.pop_frame().is_none()); // underrun -> idle assert!(jb.is_idle()); // A 2-frame utterance is below the grown target of 6, so only the // timeout can start it — and it must, exactly at PRIME_TIMEOUT_TICKS. jb.insert(100, frame(&mut enc, 700)); jb.insert(101, frame(&mut enc, 700)); let mut polls = 0; loop { polls += 1; assert!(polls <= PRIME_TIMEOUT_TICKS, "must prime by the timeout"); if jb.pop_frame().is_some() { break; } } assert_eq!(polls, PRIME_TIMEOUT_TICKS); } #[test] fn grown_target_requires_deeper_reprime() { let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); let mut jb = JitterBuffer::new().unwrap(); assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES); // 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 -> head now at seq 2 assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert_eq!(jb.next_seq, Some(2)); assert_eq!(jb.target_delay(), DEFAULT_DELAY_FRAMES); // A late packet for an already-played sequence (0) arrives: grows delay to 4 jb.insert(0, vec![0u8]); assert_eq!(jb.target_delay(), 4); // Drain: plays seq 2, then underruns (goes idle) assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_none()); assert!(jb.is_idle()); // Insert three fresh contiguous frames (seq 100, 101, 102) jb.insert(100, frame(&mut enc, 1000)); jb.insert(101, frame(&mut enc, 1000)); jb.insert(102, frame(&mut enc, 1000)); // Playout must NOT prime yet (3 < grown target of 4) assert!(jb.pop_frame().is_none()); assert!(!jb.is_idle()); // Insert a fourth frame (seq 103) -> primes and plays seq 100 jb.insert(103, frame(&mut enc, 1000)); assert!(jb.pop_frame().is_some()); } #[test] fn overflow_resync_resets_clean_run() { 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)); // Play a few in-order real frames so clean_run > 0 assert!(jb.pop_frame().is_some()); assert!(jb.pop_frame().is_some()); assert!(jb.clean_run > 0); // Insert a contiguous run long enough to exceed MAX_BUFFERED_FRAMES // packets currently contains seq 2 (length 1). // Inserting seq 3..=34 (32 frames) makes total length 33, exceeding MAX_BUFFERED_FRAMES (32) for seq in 3..=34 { jb.insert(seq, frame(&mut enc, 1000)); } // Assert overflow occurred and reset clean_run assert_eq!(jb.clean_run, 0); } }