use super::sequence::SequenceExtender;
use crate::{Packet, TaggedPacket};
use std::collections::BTreeMap;
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct JitterBufferStats {
pub out_of_order: u64,
pub duplicates: u64,
pub overflow: u64,
pub late: u64,
pub foreign_ssrc: u64,
pub underflow: u64,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub enum State {
#[default]
Buffering,
Emitting,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Rejected {
Duplicate,
Overflow,
Late,
ForeignSsrc,
}
pub struct JitterBuffer {
ssrc: Option<u32>,
packets: BTreeMap<u64, TaggedPacket>,
extender: SequenceExtender,
released_through: Option<u64>,
capacity: usize,
state: State,
stats: JitterBufferStats,
}
impl std::fmt::Debug for JitterBuffer {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("JitterBuffer")
.field("ssrc", &self.ssrc)
.field("held", &self.packets.len())
.field("capacity", &self.capacity)
.field("next_due", &self.front_sequence())
.field("released_through", &self.released_through)
.field("state", &self.state)
.field("stats", &self.stats)
.finish()
}
}
impl JitterBuffer {
pub fn new(capacity: usize) -> Self {
Self {
ssrc: None,
packets: BTreeMap::new(),
extender: SequenceExtender::new(),
released_through: None,
capacity: capacity.max(1),
state: State::Buffering,
stats: JitterBufferStats::default(),
}
}
pub fn ssrc(&self) -> Option<u32> {
self.ssrc
}
pub fn len(&self) -> usize {
self.packets.len()
}
pub fn is_empty(&self) -> bool {
self.packets.is_empty()
}
pub fn stats(&self) -> JitterBufferStats {
self.stats
}
pub fn push(&mut self, packet: TaggedPacket) -> Result<u64, Rejected> {
let Packet::Rtp(rtp) = &packet.message else {
self.stats.foreign_ssrc += 1;
return Err(Rejected::ForeignSsrc);
};
let ssrc = rtp.header.ssrc;
match self.ssrc {
None => self.ssrc = Some(ssrc),
Some(bound) if bound != ssrc => {
self.stats.foreign_ssrc += 1;
return Err(Rejected::ForeignSsrc);
}
Some(_) => {}
}
let sequence_number = rtp.header.sequence_number;
let previous_highest = self.extender.highest();
let extended = self.extender.extend(sequence_number);
if previous_highest.is_some_and(|highest| extended < highest) {
self.stats.out_of_order += 1;
}
if self
.released_through
.is_some_and(|released| extended <= released)
{
self.stats.late += 1;
return Err(Rejected::Late);
}
if self.packets.contains_key(&extended) {
self.stats.duplicates += 1;
return Err(Rejected::Duplicate);
}
if self.packets.len() >= self.capacity {
let oldest = self
.packets
.keys()
.next()
.copied()
.expect("capacity is at least 1, so a full buffer is non-empty");
if extended < oldest {
self.stats.overflow += 1;
return Err(Rejected::Overflow);
}
self.packets.remove(&oldest);
self.stats.overflow += 1;
}
self.packets.insert(extended, packet);
Ok(extended)
}
pub fn front_sequence(&self) -> Option<u64> {
self.packets.keys().next().copied()
}
pub fn peek(&self) -> Option<&TaggedPacket> {
self.packets.values().next()
}
pub fn state(&self) -> State {
self.state
}
pub fn begin_emitting(&mut self) {
self.state = State::Emitting;
}
pub fn begin_buffering(&mut self) {
self.state = State::Buffering;
}
pub fn pop(&mut self) -> Option<TaggedPacket> {
if self.state == State::Buffering {
return None;
}
let Some((&extended, _)) = self.packets.iter().next() else {
self.stats.underflow += 1;
return None;
};
let packet = self.packets.remove(&extended)?;
self.mark_released(extended);
Some(packet)
}
pub fn pop_at_sequence(&mut self, sequence_number: u16) -> Option<TaggedPacket> {
self.pop_at(self.extended_of(sequence_number)?)
}
pub fn peek_at_sequence(&self, sequence_number: u16) -> Option<&TaggedPacket> {
self.find(self.extended_of(sequence_number)?)
}
pub fn pop_at_timestamp(&mut self, timestamp: u32) -> Option<TaggedPacket> {
let extended =
self.packets
.iter()
.find_map(|(&extended, packet)| match &packet.message {
Packet::Rtp(rtp) if rtp.header.timestamp == timestamp => Some(extended),
_ => None,
})?;
self.pop_at(extended)
}
pub fn pop_at(&mut self, extended: u64) -> Option<TaggedPacket> {
let packet = self.packets.remove(&extended)?;
self.mark_released(extended);
Some(packet)
}
pub fn find(&self, extended: u64) -> Option<&TaggedPacket> {
self.packets.get(&extended)
}
fn extended_of(&self, sequence_number: u16) -> Option<u64> {
self.packets
.iter()
.find_map(|(&extended, packet)| match &packet.message {
Packet::Rtp(rtp) if rtp.header.sequence_number == sequence_number => Some(extended),
_ => None,
})
}
pub fn front_timestamp(&self) -> Option<u32> {
match &self.peek()?.message {
Packet::Rtp(rtp) => Some(rtp.header.timestamp),
_ => None,
}
}
pub fn reset(&mut self) {
self.packets.clear();
self.extender = SequenceExtender::new();
self.released_through = None;
self.state = State::Buffering;
}
fn mark_released(&mut self, extended: u64) {
self.released_through = Some(match self.released_through {
Some(previous) => previous.max(extended),
None => extended,
});
}
}
#[cfg(test)]
mod tests {
use super::*;
use shared::TransportContext;
use std::time::Instant;
fn packet(ssrc: u32, sequence_number: u16, timestamp: u32) -> TaggedPacket {
TaggedPacket {
now: Instant::now(),
transport: TransportContext::default(),
message: Packet::Rtp(rtp::Packet {
header: rtp::header::Header {
ssrc,
sequence_number,
timestamp,
..Default::default()
},
..Default::default()
}),
}
}
fn sequence_of(packet: &TaggedPacket) -> u16 {
match &packet.message {
Packet::Rtp(rtp) => rtp.header.sequence_number,
_ => panic!("not RTP"),
}
}
fn drain(buffer: &mut JitterBuffer) -> Vec<u16> {
buffer.begin_emitting();
let mut out = Vec::new();
while let Some(packet) = buffer.pop() {
out.push(sequence_of(&packet));
}
out
}
#[test]
fn packets_come_out_in_sequence_order_however_they_went_in() {
let mut buffer = JitterBuffer::new(64);
for sequence_number in [3u16, 1, 4, 2, 0] {
buffer.push(packet(1, sequence_number, 0)).expect("push");
}
assert_eq!(vec![0, 1, 2, 3, 4], drain(&mut buffer));
}
#[test]
fn reordering_is_counted() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 0, 0)).expect("push");
buffer.push(packet(1, 2, 0)).expect("push");
assert_eq!(0, buffer.stats().out_of_order);
buffer.push(packet(1, 1, 0)).expect("push");
assert_eq!(1, buffer.stats().out_of_order);
}
#[test]
fn a_duplicate_is_rejected_rather_than_stored_twice() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 7, 0)).expect("push");
assert_eq!(Err(Rejected::Duplicate), buffer.push(packet(1, 7, 0)));
assert_eq!(1, buffer.len());
assert_eq!(1, buffer.stats().duplicates);
assert_eq!(vec![7], drain(&mut buffer));
}
#[test]
fn ordering_survives_a_wrap_around() {
let mut buffer = JitterBuffer::new(64);
for sequence_number in [65534u16, 65535, 0, 1] {
buffer.push(packet(1, sequence_number, 0)).expect("push");
}
assert_eq!(vec![65534, 65535, 0, 1], drain(&mut buffer));
}
#[test]
fn ordering_survives_two_wrap_arounds_with_reordering() {
let mut buffer = JitterBuffer::new(4096);
let mut expected = Vec::new();
let mut sequence_number = 65000u16;
let mut pushed = Vec::new();
for _ in 0..(2 * 65536 / 2) {
let a = sequence_number;
let b = sequence_number.wrapping_add(1);
pushed.push((b, a)); expected.push(a);
expected.push(b);
sequence_number = sequence_number.wrapping_add(2);
}
let mut emitted = Vec::new();
buffer.begin_emitting();
for (b, a) in pushed {
let _ = buffer.push(packet(1, b, 0));
let _ = buffer.push(packet(1, a, 0));
while buffer.len() > 2 {
emitted.push(sequence_of(&buffer.pop().expect("non-empty")));
}
}
emitted.extend(drain(&mut buffer));
assert_eq!(
expected.len(),
emitted.len(),
"every packet came out exactly once across two wraps"
);
assert_eq!(expected, emitted, "and in sequence order throughout");
}
#[test]
fn a_second_ssrc_cannot_interleave_into_this_buffer() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 10, 0)).expect("push");
assert_eq!(
Err(Rejected::ForeignSsrc),
buffer.push(packet(2, 5, 0)),
"a different stream's packet must not sort against this stream's sequence numbers"
);
assert_eq!(1, buffer.stats().foreign_ssrc);
assert_eq!(Some(1), buffer.ssrc());
assert_eq!(vec![10], drain(&mut buffer));
}
#[test]
fn rtcp_is_not_stored() {
let mut buffer = JitterBuffer::new(64);
let rtcp = TaggedPacket {
now: Instant::now(),
transport: TransportContext::default(),
message: Packet::Rtcp(vec![]),
};
assert_eq!(Err(Rejected::ForeignSsrc), buffer.push(rtcp));
assert!(buffer.is_empty());
}
#[test]
fn a_packet_arriving_after_its_position_was_released_is_rejected() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 5, 0)).expect("push");
buffer.begin_emitting();
buffer.pop().expect("release 5");
assert_eq!(
Err(Rejected::Late),
buffer.push(packet(1, 5, 0)),
"the same packet again would be emitted twice"
);
assert_eq!(
Err(Rejected::Late),
buffer.push(packet(1, 4, 0)),
"an older straggler would be emitted out of order"
);
assert_eq!(2, buffer.stats().late);
assert!(buffer.is_empty());
}
#[test]
fn a_full_buffer_drops_its_oldest_packet() {
let mut buffer = JitterBuffer::new(3);
for sequence_number in 0..3u16 {
buffer.push(packet(1, sequence_number, 0)).expect("push");
}
buffer.push(packet(1, 3, 0)).expect("push displaces");
assert_eq!(3, buffer.len(), "capacity is respected");
assert_eq!(1, buffer.stats().overflow);
assert_eq!(vec![1, 2, 3], drain(&mut buffer), "the oldest gave way");
}
#[test]
fn a_full_buffer_drops_an_arrival_older_than_everything_held() {
let mut buffer = JitterBuffer::new(3);
for sequence_number in [10u16, 11, 12] {
buffer.push(packet(1, sequence_number, 0)).expect("push");
}
assert_eq!(Err(Rejected::Overflow), buffer.push(packet(1, 9, 0)));
assert_eq!(vec![10, 11, 12], drain(&mut buffer));
}
#[test]
fn pop_at_and_find_address_packets_by_extended_sequence_number() {
let mut buffer = JitterBuffer::new(64);
let first = buffer.push(packet(1, 100, 0)).expect("push");
let second = buffer.push(packet(1, 101, 0)).expect("push");
assert!(buffer.find(second).is_some());
assert!(buffer.find(second + 10).is_none());
assert_eq!(101, sequence_of(&buffer.pop_at(second).expect("pop_at")));
assert!(buffer.find(second).is_none());
assert_eq!(100, sequence_of(&buffer.pop_at(first).expect("pop_at")));
assert!(buffer.is_empty());
}
#[test]
fn the_front_reports_the_next_packet_due() {
let mut buffer = JitterBuffer::new(64);
assert_eq!(None, buffer.front_sequence());
assert_eq!(None, buffer.front_timestamp());
buffer.push(packet(1, 8, 9000)).expect("push");
buffer.push(packet(1, 7, 3000)).expect("push");
assert_eq!(Some(7), buffer.front_sequence());
assert_eq!(Some(3000), buffer.front_timestamp());
assert_eq!(7, sequence_of(buffer.peek().expect("peek")));
assert_eq!(2, buffer.len(), "peeking does not remove");
}
#[test]
fn nothing_is_handed_out_while_buffering() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 1, 0)).expect("push");
assert_eq!(
State::Buffering,
buffer.state(),
"a fresh buffer fills first"
);
assert!(buffer.pop().is_none(), "no packets while buffering");
assert_eq!(1, buffer.len(), "and the packet is still held, not dropped");
assert_eq!(
0,
buffer.stats().underflow,
"declining to emit while buffering is not an underflow"
);
buffer.begin_emitting();
assert_eq!(1, sequence_of(&buffer.pop().expect("now emitting")));
}
#[test]
fn asking_an_empty_emitting_buffer_for_a_packet_is_an_underflow() {
let mut buffer = JitterBuffer::new(64);
buffer.begin_emitting();
assert!(buffer.pop().is_none());
assert_eq!(1, buffer.stats().underflow, "playout asked and got nothing");
buffer.begin_buffering();
assert!(buffer.pop().is_none());
assert_eq!(
1,
buffer.stats().underflow,
"but a buffering stream is not underflowing — it has not started"
);
}
#[test]
fn a_gap_does_not_stall_playout() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 1, 0)).expect("push");
buffer.push(packet(1, 3, 0)).expect("push");
assert_eq!(vec![1, 3], drain(&mut buffer));
}
#[test]
fn packets_are_addressable_by_their_wire_sequence_number() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 100, 0)).expect("push");
buffer.push(packet(1, 101, 0)).expect("push");
assert_eq!(
101,
sequence_of(buffer.peek_at_sequence(101).expect("peek"))
);
assert!(buffer.peek_at_sequence(999).is_none());
assert_eq!(2, buffer.len(), "peeking does not remove");
assert_eq!(101, sequence_of(&buffer.pop_at_sequence(101).expect("pop")));
assert!(buffer.peek_at_sequence(101).is_none());
assert!(buffer.pop_at_sequence(101).is_none(), "already taken");
assert_eq!(1, buffer.len());
}
#[test]
fn addressing_by_wire_sequence_number_works_across_a_wrap() {
let mut buffer = JitterBuffer::new(64);
for sequence_number in [65534u16, 65535, 0, 1] {
buffer.push(packet(1, sequence_number, 0)).expect("push");
}
assert_eq!(
0,
sequence_of(buffer.peek_at_sequence(0).expect("post-wrap"))
);
assert_eq!(
65535,
sequence_of(buffer.peek_at_sequence(65535).expect("pre-wrap"))
);
assert_eq!(0, sequence_of(&buffer.pop_at_sequence(0).expect("pop 0")));
assert_eq!(vec![65534, 65535, 1], drain(&mut buffer));
}
#[test]
fn packets_are_addressable_by_rtp_timestamp() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 10, 9000)).expect("push");
buffer.push(packet(1, 11, 9000)).expect("push"); buffer.push(packet(1, 12, 12000)).expect("push");
assert_eq!(
10,
sequence_of(&buffer.pop_at_timestamp(9000).expect("first of the frame")),
"the lowest sequence number of that timestamp comes out first"
);
assert_eq!(
11,
sequence_of(&buffer.pop_at_timestamp(9000).expect("second of the frame"))
);
assert!(
buffer.pop_at_timestamp(9000).is_none(),
"the frame is fully released"
);
assert!(buffer.pop_at_timestamp(4242).is_none(), "no such timestamp");
assert_eq!(1, buffer.len(), "the next frame is untouched");
assert_eq!(vec![12], drain(&mut buffer));
}
#[test]
fn addressed_removal_also_marks_the_position_released() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 5, 0)).expect("push");
buffer.push(packet(1, 6, 0)).expect("push");
buffer.pop_at_sequence(6).expect("pop 6");
assert_eq!(
Err(Rejected::Late),
buffer.push(packet(1, 6, 0)),
"6 has been played out"
);
assert_eq!(
Err(Rejected::Late),
buffer.push(packet(1, 5, 0)),
"and so has everything before it"
);
}
#[test]
fn reset_clears_the_contents_but_keeps_the_stream_and_counters() {
let mut buffer = JitterBuffer::new(64);
buffer.push(packet(1, 5, 0)).expect("push");
buffer.begin_emitting();
buffer.pop().expect("release");
buffer.push(packet(1, 5, 0)).ok();
buffer.reset();
assert!(buffer.is_empty());
assert_eq!(Some(1), buffer.ssrc(), "still this stream's buffer");
assert_eq!(1, buffer.stats().late, "history is not erased");
assert_eq!(
State::Buffering,
buffer.state(),
"a restarted stream re-accumulates before playout resumes"
);
assert!(buffer.push(packet(1, 5, 0)).is_ok());
}
}