1use super::sequence::SequenceExtender;
7use crate::{Packet, TaggedPacket};
8use std::collections::BTreeMap;
9
10#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
16pub struct JitterBufferStats {
17 pub out_of_order: u64,
19 pub duplicates: u64,
21 pub overflow: u64,
23 pub late: u64,
26 pub foreign_ssrc: u64,
28 pub underflow: u64,
30}
31
32#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
39pub enum State {
40 #[default]
42 Buffering,
43 Emitting,
45}
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub enum Rejected {
50 Duplicate,
52 Overflow,
54 Late,
56 ForeignSsrc,
58}
59
60pub struct JitterBuffer {
73 ssrc: Option<u32>,
75 packets: BTreeMap<u64, TaggedPacket>,
77 extender: SequenceExtender,
78 released_through: Option<u64>,
80 capacity: usize,
82 state: State,
83 stats: JitterBufferStats,
84}
85
86impl std::fmt::Debug for JitterBuffer {
89 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90 f.debug_struct("JitterBuffer")
91 .field("ssrc", &self.ssrc)
92 .field("held", &self.packets.len())
93 .field("capacity", &self.capacity)
94 .field("next_due", &self.front_sequence())
95 .field("released_through", &self.released_through)
96 .field("state", &self.state)
97 .field("stats", &self.stats)
98 .finish()
99 }
100}
101
102impl JitterBuffer {
103 pub fn new(capacity: usize) -> Self {
105 Self {
106 ssrc: None,
107 packets: BTreeMap::new(),
108 extender: SequenceExtender::new(),
109 released_through: None,
110 capacity: capacity.max(1),
111 state: State::Buffering,
112 stats: JitterBufferStats::default(),
113 }
114 }
115
116 pub fn ssrc(&self) -> Option<u32> {
118 self.ssrc
119 }
120
121 pub fn len(&self) -> usize {
123 self.packets.len()
124 }
125
126 pub fn is_empty(&self) -> bool {
128 self.packets.is_empty()
129 }
130
131 pub fn stats(&self) -> JitterBufferStats {
133 self.stats
134 }
135
136 pub fn push(&mut self, packet: TaggedPacket) -> Result<u64, Rejected> {
141 let Packet::Rtp(rtp) = &packet.message.packet else {
142 self.stats.foreign_ssrc += 1;
143 return Err(Rejected::ForeignSsrc);
144 };
145
146 let ssrc = rtp.header.ssrc;
147 match self.ssrc {
148 None => self.ssrc = Some(ssrc),
149 Some(bound) if bound != ssrc => {
150 self.stats.foreign_ssrc += 1;
151 return Err(Rejected::ForeignSsrc);
152 }
153 Some(_) => {}
154 }
155
156 let sequence_number = rtp.header.sequence_number;
157 let previous_highest = self.extender.highest();
158 let extended = self.extender.extend(sequence_number);
159
160 if previous_highest.is_some_and(|highest| extended < highest) {
161 self.stats.out_of_order += 1;
162 }
163
164 if self
166 .released_through
167 .is_some_and(|released| extended <= released)
168 {
169 self.stats.late += 1;
170 return Err(Rejected::Late);
171 }
172
173 if self.packets.contains_key(&extended) {
174 self.stats.duplicates += 1;
175 return Err(Rejected::Duplicate);
176 }
177
178 if self.packets.len() >= self.capacity {
179 let oldest = self
182 .packets
183 .keys()
184 .next()
185 .copied()
186 .expect("capacity is at least 1, so a full buffer is non-empty");
187 if extended < oldest {
188 self.stats.overflow += 1;
189 return Err(Rejected::Overflow);
190 }
191 self.packets.remove(&oldest);
192 self.stats.overflow += 1;
193 }
194
195 self.packets.insert(extended, packet);
196 Ok(extended)
197 }
198
199 pub fn front_sequence(&self) -> Option<u64> {
201 self.packets.keys().next().copied()
202 }
203
204 pub fn peek(&self) -> Option<&TaggedPacket> {
206 self.packets.values().next()
207 }
208
209 pub fn state(&self) -> State {
211 self.state
212 }
213
214 pub fn begin_emitting(&mut self) {
219 self.state = State::Emitting;
220 }
221
222 pub fn begin_buffering(&mut self) {
227 self.state = State::Buffering;
228 }
229
230 pub fn pop(&mut self) -> Option<TaggedPacket> {
244 if self.state == State::Buffering {
245 return None;
246 }
247 let Some((&extended, _)) = self.packets.iter().next() else {
248 self.stats.underflow += 1;
249 return None;
250 };
251 let packet = self.packets.remove(&extended)?;
252 self.mark_released(extended);
253 Some(packet)
254 }
255
256 pub fn pop_at_sequence(&mut self, sequence_number: u16) -> Option<TaggedPacket> {
261 self.pop_at(self.extended_of(sequence_number)?)
262 }
263
264 pub fn peek_at_sequence(&self, sequence_number: u16) -> Option<&TaggedPacket> {
266 self.find(self.extended_of(sequence_number)?)
267 }
268
269 pub fn pop_at_timestamp(&mut self, timestamp: u32) -> Option<TaggedPacket> {
276 let extended =
277 self.packets
278 .iter()
279 .find_map(|(&extended, packet)| match &packet.message.packet {
280 Packet::Rtp(rtp) if rtp.header.timestamp == timestamp => Some(extended),
281 _ => None,
282 })?;
283 self.pop_at(extended)
284 }
285
286 pub fn pop_at(&mut self, extended: u64) -> Option<TaggedPacket> {
291 let packet = self.packets.remove(&extended)?;
292 self.mark_released(extended);
293 Some(packet)
294 }
295
296 pub fn find(&self, extended: u64) -> Option<&TaggedPacket> {
298 self.packets.get(&extended)
299 }
300
301 fn extended_of(&self, sequence_number: u16) -> Option<u64> {
307 self.packets
308 .iter()
309 .find_map(|(&extended, packet)| match &packet.message.packet {
310 Packet::Rtp(rtp) if rtp.header.sequence_number == sequence_number => Some(extended),
311 _ => None,
312 })
313 }
314
315 pub fn front_timestamp(&self) -> Option<u32> {
317 match &self.peek()?.message.packet {
318 Packet::Rtp(rtp) => Some(rtp.header.timestamp),
319 _ => None,
320 }
321 }
322
323 pub fn reset(&mut self) {
328 self.packets.clear();
329 self.extender = SequenceExtender::new();
330 self.released_through = None;
331 self.state = State::Buffering;
332 }
333
334 fn mark_released(&mut self, extended: u64) {
335 self.released_through = Some(match self.released_through {
336 Some(previous) => previous.max(extended),
337 None => extended,
338 });
339 }
340}
341
342#[cfg(test)]
343mod tests {
344 use super::*;
345 use crate::AttributedPacket;
346 use shared::TransportContext;
347 use std::time::Instant;
348
349 fn packet(ssrc: u32, sequence_number: u16, timestamp: u32) -> TaggedPacket {
350 TaggedPacket {
351 now: Instant::now(),
352 transport: TransportContext::default(),
353 message: AttributedPacket::new(Packet::Rtp(rtp::Packet {
354 header: rtp::header::Header {
355 ssrc,
356 sequence_number,
357 timestamp,
358 ..Default::default()
359 },
360 ..Default::default()
361 })),
362 }
363 }
364
365 fn sequence_of(packet: &TaggedPacket) -> u16 {
366 match &packet.message.packet {
367 Packet::Rtp(rtp) => rtp.header.sequence_number,
368 _ => panic!("not RTP"),
369 }
370 }
371
372 fn drain(buffer: &mut JitterBuffer) -> Vec<u16> {
375 buffer.begin_emitting();
376 let mut out = Vec::new();
377 while let Some(packet) = buffer.pop() {
378 out.push(sequence_of(&packet));
379 }
380 out
381 }
382
383 #[test]
384 fn packets_come_out_in_sequence_order_however_they_went_in() {
385 let mut buffer = JitterBuffer::new(64);
386 for sequence_number in [3u16, 1, 4, 2, 0] {
387 buffer.push(packet(1, sequence_number, 0)).expect("push");
388 }
389 assert_eq!(vec![0, 1, 2, 3, 4], drain(&mut buffer));
390 }
391
392 #[test]
393 fn reordering_is_counted() {
394 let mut buffer = JitterBuffer::new(64);
395 buffer.push(packet(1, 0, 0)).expect("push");
396 buffer.push(packet(1, 2, 0)).expect("push");
397 assert_eq!(0, buffer.stats().out_of_order);
398 buffer.push(packet(1, 1, 0)).expect("push");
399 assert_eq!(1, buffer.stats().out_of_order);
400 }
401
402 #[test]
405 fn a_duplicate_is_rejected_rather_than_stored_twice() {
406 let mut buffer = JitterBuffer::new(64);
407 buffer.push(packet(1, 7, 0)).expect("push");
408 assert_eq!(Err(Rejected::Duplicate), buffer.push(packet(1, 7, 0)));
409 assert_eq!(1, buffer.len());
410 assert_eq!(1, buffer.stats().duplicates);
411 assert_eq!(vec![7], drain(&mut buffer));
412 }
413
414 #[test]
416 fn ordering_survives_a_wrap_around() {
417 let mut buffer = JitterBuffer::new(64);
418 for sequence_number in [65534u16, 65535, 0, 1] {
419 buffer.push(packet(1, sequence_number, 0)).expect("push");
420 }
421 assert_eq!(vec![65534, 65535, 0, 1], drain(&mut buffer));
422 }
423
424 #[test]
425 fn ordering_survives_two_wrap_arounds_with_reordering() {
426 let mut buffer = JitterBuffer::new(4096);
427
428 let mut expected = Vec::new();
430 let mut sequence_number = 65000u16;
431 let mut pushed = Vec::new();
432 for _ in 0..(2 * 65536 / 2) {
433 let a = sequence_number;
434 let b = sequence_number.wrapping_add(1);
435 pushed.push((b, a)); expected.push(a);
437 expected.push(b);
438 sequence_number = sequence_number.wrapping_add(2);
439 }
440
441 let mut emitted = Vec::new();
443 buffer.begin_emitting();
444 for (b, a) in pushed {
445 let _ = buffer.push(packet(1, b, 0));
446 let _ = buffer.push(packet(1, a, 0));
447 while buffer.len() > 2 {
448 emitted.push(sequence_of(&buffer.pop().expect("non-empty")));
449 }
450 }
451 emitted.extend(drain(&mut buffer));
452
453 assert_eq!(
454 expected.len(),
455 emitted.len(),
456 "every packet came out exactly once across two wraps"
457 );
458 assert_eq!(expected, emitted, "and in sequence order throughout");
459 }
460
461 #[test]
463 fn a_second_ssrc_cannot_interleave_into_this_buffer() {
464 let mut buffer = JitterBuffer::new(64);
465 buffer.push(packet(1, 10, 0)).expect("push");
466
467 assert_eq!(
468 Err(Rejected::ForeignSsrc),
469 buffer.push(packet(2, 5, 0)),
470 "a different stream's packet must not sort against this stream's sequence numbers"
471 );
472 assert_eq!(1, buffer.stats().foreign_ssrc);
473 assert_eq!(Some(1), buffer.ssrc());
474 assert_eq!(vec![10], drain(&mut buffer));
475 }
476
477 #[test]
478 fn rtcp_is_not_stored() {
479 let mut buffer = JitterBuffer::new(64);
480 let rtcp = TaggedPacket {
481 now: Instant::now(),
482 transport: TransportContext::default(),
483 message: AttributedPacket::new(Packet::Rtcp(vec![])),
484 };
485 assert_eq!(Err(Rejected::ForeignSsrc), buffer.push(rtcp));
486 assert!(buffer.is_empty());
487 }
488
489 #[test]
490 fn a_packet_arriving_after_its_position_was_released_is_rejected() {
491 let mut buffer = JitterBuffer::new(64);
492 buffer.push(packet(1, 5, 0)).expect("push");
493 buffer.begin_emitting();
494 buffer.pop().expect("release 5");
495
496 assert_eq!(
497 Err(Rejected::Late),
498 buffer.push(packet(1, 5, 0)),
499 "the same packet again would be emitted twice"
500 );
501 assert_eq!(
502 Err(Rejected::Late),
503 buffer.push(packet(1, 4, 0)),
504 "an older straggler would be emitted out of order"
505 );
506 assert_eq!(2, buffer.stats().late);
507 assert!(buffer.is_empty());
508 }
509
510 #[test]
511 fn a_full_buffer_drops_its_oldest_packet() {
512 let mut buffer = JitterBuffer::new(3);
513 for sequence_number in 0..3u16 {
514 buffer.push(packet(1, sequence_number, 0)).expect("push");
515 }
516 buffer.push(packet(1, 3, 0)).expect("push displaces");
517
518 assert_eq!(3, buffer.len(), "capacity is respected");
519 assert_eq!(1, buffer.stats().overflow);
520 assert_eq!(vec![1, 2, 3], drain(&mut buffer), "the oldest gave way");
521 }
522
523 #[test]
527 fn a_full_buffer_drops_an_arrival_older_than_everything_held() {
528 let mut buffer = JitterBuffer::new(3);
529 for sequence_number in [10u16, 11, 12] {
530 buffer.push(packet(1, sequence_number, 0)).expect("push");
531 }
532
533 assert_eq!(Err(Rejected::Overflow), buffer.push(packet(1, 9, 0)));
534 assert_eq!(vec![10, 11, 12], drain(&mut buffer));
535 }
536
537 #[test]
538 fn pop_at_and_find_address_packets_by_extended_sequence_number() {
539 let mut buffer = JitterBuffer::new(64);
540 let first = buffer.push(packet(1, 100, 0)).expect("push");
541 let second = buffer.push(packet(1, 101, 0)).expect("push");
542
543 assert!(buffer.find(second).is_some());
544 assert!(buffer.find(second + 10).is_none());
545
546 assert_eq!(101, sequence_of(&buffer.pop_at(second).expect("pop_at")));
547 assert!(buffer.find(second).is_none());
548 assert_eq!(100, sequence_of(&buffer.pop_at(first).expect("pop_at")));
549 assert!(buffer.is_empty());
550 }
551
552 #[test]
553 fn the_front_reports_the_next_packet_due() {
554 let mut buffer = JitterBuffer::new(64);
555 assert_eq!(None, buffer.front_sequence());
556 assert_eq!(None, buffer.front_timestamp());
557
558 buffer.push(packet(1, 8, 9000)).expect("push");
559 buffer.push(packet(1, 7, 3000)).expect("push");
560
561 assert_eq!(Some(7), buffer.front_sequence());
562 assert_eq!(Some(3000), buffer.front_timestamp());
563 assert_eq!(7, sequence_of(buffer.peek().expect("peek")));
564 assert_eq!(2, buffer.len(), "peeking does not remove");
565 }
566
567 #[test]
570 fn nothing_is_handed_out_while_buffering() {
571 let mut buffer = JitterBuffer::new(64);
572 buffer.push(packet(1, 1, 0)).expect("push");
573
574 assert_eq!(
575 State::Buffering,
576 buffer.state(),
577 "a fresh buffer fills first"
578 );
579 assert!(buffer.pop().is_none(), "no packets while buffering");
580 assert_eq!(1, buffer.len(), "and the packet is still held, not dropped");
581 assert_eq!(
582 0,
583 buffer.stats().underflow,
584 "declining to emit while buffering is not an underflow"
585 );
586
587 buffer.begin_emitting();
588 assert_eq!(1, sequence_of(&buffer.pop().expect("now emitting")));
589 }
590
591 #[test]
592 fn asking_an_empty_emitting_buffer_for_a_packet_is_an_underflow() {
593 let mut buffer = JitterBuffer::new(64);
594 buffer.begin_emitting();
595
596 assert!(buffer.pop().is_none());
597 assert_eq!(1, buffer.stats().underflow, "playout asked and got nothing");
598
599 buffer.begin_buffering();
600 assert!(buffer.pop().is_none());
601 assert_eq!(
602 1,
603 buffer.stats().underflow,
604 "but a buffering stream is not underflowing — it has not started"
605 );
606 }
607
608 #[test]
612 fn a_gap_does_not_stall_playout() {
613 let mut buffer = JitterBuffer::new(64);
614 buffer.push(packet(1, 1, 0)).expect("push");
615 buffer.push(packet(1, 3, 0)).expect("push"); assert_eq!(vec![1, 3], drain(&mut buffer));
618 }
619
620 #[test]
621 fn packets_are_addressable_by_their_wire_sequence_number() {
622 let mut buffer = JitterBuffer::new(64);
623 buffer.push(packet(1, 100, 0)).expect("push");
624 buffer.push(packet(1, 101, 0)).expect("push");
625
626 assert_eq!(
627 101,
628 sequence_of(buffer.peek_at_sequence(101).expect("peek"))
629 );
630 assert!(buffer.peek_at_sequence(999).is_none());
631 assert_eq!(2, buffer.len(), "peeking does not remove");
632
633 assert_eq!(101, sequence_of(&buffer.pop_at_sequence(101).expect("pop")));
634 assert!(buffer.peek_at_sequence(101).is_none());
635 assert!(buffer.pop_at_sequence(101).is_none(), "already taken");
636 assert_eq!(1, buffer.len());
637 }
638
639 #[test]
642 fn addressing_by_wire_sequence_number_works_across_a_wrap() {
643 let mut buffer = JitterBuffer::new(64);
644 for sequence_number in [65534u16, 65535, 0, 1] {
645 buffer.push(packet(1, sequence_number, 0)).expect("push");
646 }
647
648 assert_eq!(
649 0,
650 sequence_of(buffer.peek_at_sequence(0).expect("post-wrap"))
651 );
652 assert_eq!(
653 65535,
654 sequence_of(buffer.peek_at_sequence(65535).expect("pre-wrap"))
655 );
656 assert_eq!(0, sequence_of(&buffer.pop_at_sequence(0).expect("pop 0")));
657 assert_eq!(vec![65534, 65535, 1], drain(&mut buffer));
658 }
659
660 #[test]
663 fn packets_are_addressable_by_rtp_timestamp() {
664 let mut buffer = JitterBuffer::new(64);
665 buffer.push(packet(1, 10, 9000)).expect("push");
666 buffer.push(packet(1, 11, 9000)).expect("push"); buffer.push(packet(1, 12, 12000)).expect("push"); assert_eq!(
670 10,
671 sequence_of(&buffer.pop_at_timestamp(9000).expect("first of the frame")),
672 "the lowest sequence number of that timestamp comes out first"
673 );
674 assert_eq!(
675 11,
676 sequence_of(&buffer.pop_at_timestamp(9000).expect("second of the frame"))
677 );
678 assert!(
679 buffer.pop_at_timestamp(9000).is_none(),
680 "the frame is fully released"
681 );
682 assert!(buffer.pop_at_timestamp(4242).is_none(), "no such timestamp");
683
684 assert_eq!(1, buffer.len(), "the next frame is untouched");
685 assert_eq!(vec![12], drain(&mut buffer));
686 }
687
688 #[test]
691 fn addressed_removal_also_marks_the_position_released() {
692 let mut buffer = JitterBuffer::new(64);
693 buffer.push(packet(1, 5, 0)).expect("push");
694 buffer.push(packet(1, 6, 0)).expect("push");
695
696 buffer.pop_at_sequence(6).expect("pop 6");
697 assert_eq!(
698 Err(Rejected::Late),
699 buffer.push(packet(1, 6, 0)),
700 "6 has been played out"
701 );
702 assert_eq!(
703 Err(Rejected::Late),
704 buffer.push(packet(1, 5, 0)),
705 "and so has everything before it"
706 );
707 }
708
709 #[test]
710 fn reset_clears_the_contents_but_keeps_the_stream_and_counters() {
711 let mut buffer = JitterBuffer::new(64);
712 buffer.push(packet(1, 5, 0)).expect("push");
713 buffer.begin_emitting();
714 buffer.pop().expect("release");
715 buffer.push(packet(1, 5, 0)).ok(); buffer.reset();
718
719 assert!(buffer.is_empty());
720 assert_eq!(Some(1), buffer.ssrc(), "still this stream's buffer");
721 assert_eq!(1, buffer.stats().late, "history is not erased");
722 assert_eq!(
723 State::Buffering,
724 buffer.state(),
725 "a restarted stream re-accumulates before playout resumes"
726 );
727 assert!(buffer.push(packet(1, 5, 0)).is_ok());
729 }
730}