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 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 {
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 {
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 {
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 shared::TransportContext;
346 use std::time::Instant;
347
348 fn packet(ssrc: u32, sequence_number: u16, timestamp: u32) -> TaggedPacket {
349 TaggedPacket {
350 now: Instant::now(),
351 transport: TransportContext::default(),
352 message: Packet::Rtp(rtp::Packet {
353 header: rtp::header::Header {
354 ssrc,
355 sequence_number,
356 timestamp,
357 ..Default::default()
358 },
359 ..Default::default()
360 }),
361 }
362 }
363
364 fn sequence_of(packet: &TaggedPacket) -> u16 {
365 match &packet.message {
366 Packet::Rtp(rtp) => rtp.header.sequence_number,
367 _ => panic!("not RTP"),
368 }
369 }
370
371 fn drain(buffer: &mut JitterBuffer) -> Vec<u16> {
374 buffer.begin_emitting();
375 let mut out = Vec::new();
376 while let Some(packet) = buffer.pop() {
377 out.push(sequence_of(&packet));
378 }
379 out
380 }
381
382 #[test]
383 fn packets_come_out_in_sequence_order_however_they_went_in() {
384 let mut buffer = JitterBuffer::new(64);
385 for sequence_number in [3u16, 1, 4, 2, 0] {
386 buffer.push(packet(1, sequence_number, 0)).expect("push");
387 }
388 assert_eq!(vec![0, 1, 2, 3, 4], drain(&mut buffer));
389 }
390
391 #[test]
392 fn reordering_is_counted() {
393 let mut buffer = JitterBuffer::new(64);
394 buffer.push(packet(1, 0, 0)).expect("push");
395 buffer.push(packet(1, 2, 0)).expect("push");
396 assert_eq!(0, buffer.stats().out_of_order);
397 buffer.push(packet(1, 1, 0)).expect("push");
398 assert_eq!(1, buffer.stats().out_of_order);
399 }
400
401 #[test]
404 fn a_duplicate_is_rejected_rather_than_stored_twice() {
405 let mut buffer = JitterBuffer::new(64);
406 buffer.push(packet(1, 7, 0)).expect("push");
407 assert_eq!(Err(Rejected::Duplicate), buffer.push(packet(1, 7, 0)));
408 assert_eq!(1, buffer.len());
409 assert_eq!(1, buffer.stats().duplicates);
410 assert_eq!(vec![7], drain(&mut buffer));
411 }
412
413 #[test]
415 fn ordering_survives_a_wrap_around() {
416 let mut buffer = JitterBuffer::new(64);
417 for sequence_number in [65534u16, 65535, 0, 1] {
418 buffer.push(packet(1, sequence_number, 0)).expect("push");
419 }
420 assert_eq!(vec![65534, 65535, 0, 1], drain(&mut buffer));
421 }
422
423 #[test]
424 fn ordering_survives_two_wrap_arounds_with_reordering() {
425 let mut buffer = JitterBuffer::new(4096);
426
427 let mut expected = Vec::new();
429 let mut sequence_number = 65000u16;
430 let mut pushed = Vec::new();
431 for _ in 0..(2 * 65536 / 2) {
432 let a = sequence_number;
433 let b = sequence_number.wrapping_add(1);
434 pushed.push((b, a)); expected.push(a);
436 expected.push(b);
437 sequence_number = sequence_number.wrapping_add(2);
438 }
439
440 let mut emitted = Vec::new();
442 buffer.begin_emitting();
443 for (b, a) in pushed {
444 let _ = buffer.push(packet(1, b, 0));
445 let _ = buffer.push(packet(1, a, 0));
446 while buffer.len() > 2 {
447 emitted.push(sequence_of(&buffer.pop().expect("non-empty")));
448 }
449 }
450 emitted.extend(drain(&mut buffer));
451
452 assert_eq!(
453 expected.len(),
454 emitted.len(),
455 "every packet came out exactly once across two wraps"
456 );
457 assert_eq!(expected, emitted, "and in sequence order throughout");
458 }
459
460 #[test]
462 fn a_second_ssrc_cannot_interleave_into_this_buffer() {
463 let mut buffer = JitterBuffer::new(64);
464 buffer.push(packet(1, 10, 0)).expect("push");
465
466 assert_eq!(
467 Err(Rejected::ForeignSsrc),
468 buffer.push(packet(2, 5, 0)),
469 "a different stream's packet must not sort against this stream's sequence numbers"
470 );
471 assert_eq!(1, buffer.stats().foreign_ssrc);
472 assert_eq!(Some(1), buffer.ssrc());
473 assert_eq!(vec![10], drain(&mut buffer));
474 }
475
476 #[test]
477 fn rtcp_is_not_stored() {
478 let mut buffer = JitterBuffer::new(64);
479 let rtcp = TaggedPacket {
480 now: Instant::now(),
481 transport: TransportContext::default(),
482 message: Packet::Rtcp(vec![]),
483 };
484 assert_eq!(Err(Rejected::ForeignSsrc), buffer.push(rtcp));
485 assert!(buffer.is_empty());
486 }
487
488 #[test]
489 fn a_packet_arriving_after_its_position_was_released_is_rejected() {
490 let mut buffer = JitterBuffer::new(64);
491 buffer.push(packet(1, 5, 0)).expect("push");
492 buffer.begin_emitting();
493 buffer.pop().expect("release 5");
494
495 assert_eq!(
496 Err(Rejected::Late),
497 buffer.push(packet(1, 5, 0)),
498 "the same packet again would be emitted twice"
499 );
500 assert_eq!(
501 Err(Rejected::Late),
502 buffer.push(packet(1, 4, 0)),
503 "an older straggler would be emitted out of order"
504 );
505 assert_eq!(2, buffer.stats().late);
506 assert!(buffer.is_empty());
507 }
508
509 #[test]
510 fn a_full_buffer_drops_its_oldest_packet() {
511 let mut buffer = JitterBuffer::new(3);
512 for sequence_number in 0..3u16 {
513 buffer.push(packet(1, sequence_number, 0)).expect("push");
514 }
515 buffer.push(packet(1, 3, 0)).expect("push displaces");
516
517 assert_eq!(3, buffer.len(), "capacity is respected");
518 assert_eq!(1, buffer.stats().overflow);
519 assert_eq!(vec![1, 2, 3], drain(&mut buffer), "the oldest gave way");
520 }
521
522 #[test]
526 fn a_full_buffer_drops_an_arrival_older_than_everything_held() {
527 let mut buffer = JitterBuffer::new(3);
528 for sequence_number in [10u16, 11, 12] {
529 buffer.push(packet(1, sequence_number, 0)).expect("push");
530 }
531
532 assert_eq!(Err(Rejected::Overflow), buffer.push(packet(1, 9, 0)));
533 assert_eq!(vec![10, 11, 12], drain(&mut buffer));
534 }
535
536 #[test]
537 fn pop_at_and_find_address_packets_by_extended_sequence_number() {
538 let mut buffer = JitterBuffer::new(64);
539 let first = buffer.push(packet(1, 100, 0)).expect("push");
540 let second = buffer.push(packet(1, 101, 0)).expect("push");
541
542 assert!(buffer.find(second).is_some());
543 assert!(buffer.find(second + 10).is_none());
544
545 assert_eq!(101, sequence_of(&buffer.pop_at(second).expect("pop_at")));
546 assert!(buffer.find(second).is_none());
547 assert_eq!(100, sequence_of(&buffer.pop_at(first).expect("pop_at")));
548 assert!(buffer.is_empty());
549 }
550
551 #[test]
552 fn the_front_reports_the_next_packet_due() {
553 let mut buffer = JitterBuffer::new(64);
554 assert_eq!(None, buffer.front_sequence());
555 assert_eq!(None, buffer.front_timestamp());
556
557 buffer.push(packet(1, 8, 9000)).expect("push");
558 buffer.push(packet(1, 7, 3000)).expect("push");
559
560 assert_eq!(Some(7), buffer.front_sequence());
561 assert_eq!(Some(3000), buffer.front_timestamp());
562 assert_eq!(7, sequence_of(buffer.peek().expect("peek")));
563 assert_eq!(2, buffer.len(), "peeking does not remove");
564 }
565
566 #[test]
569 fn nothing_is_handed_out_while_buffering() {
570 let mut buffer = JitterBuffer::new(64);
571 buffer.push(packet(1, 1, 0)).expect("push");
572
573 assert_eq!(
574 State::Buffering,
575 buffer.state(),
576 "a fresh buffer fills first"
577 );
578 assert!(buffer.pop().is_none(), "no packets while buffering");
579 assert_eq!(1, buffer.len(), "and the packet is still held, not dropped");
580 assert_eq!(
581 0,
582 buffer.stats().underflow,
583 "declining to emit while buffering is not an underflow"
584 );
585
586 buffer.begin_emitting();
587 assert_eq!(1, sequence_of(&buffer.pop().expect("now emitting")));
588 }
589
590 #[test]
591 fn asking_an_empty_emitting_buffer_for_a_packet_is_an_underflow() {
592 let mut buffer = JitterBuffer::new(64);
593 buffer.begin_emitting();
594
595 assert!(buffer.pop().is_none());
596 assert_eq!(1, buffer.stats().underflow, "playout asked and got nothing");
597
598 buffer.begin_buffering();
599 assert!(buffer.pop().is_none());
600 assert_eq!(
601 1,
602 buffer.stats().underflow,
603 "but a buffering stream is not underflowing — it has not started"
604 );
605 }
606
607 #[test]
611 fn a_gap_does_not_stall_playout() {
612 let mut buffer = JitterBuffer::new(64);
613 buffer.push(packet(1, 1, 0)).expect("push");
614 buffer.push(packet(1, 3, 0)).expect("push"); assert_eq!(vec![1, 3], drain(&mut buffer));
617 }
618
619 #[test]
620 fn packets_are_addressable_by_their_wire_sequence_number() {
621 let mut buffer = JitterBuffer::new(64);
622 buffer.push(packet(1, 100, 0)).expect("push");
623 buffer.push(packet(1, 101, 0)).expect("push");
624
625 assert_eq!(
626 101,
627 sequence_of(buffer.peek_at_sequence(101).expect("peek"))
628 );
629 assert!(buffer.peek_at_sequence(999).is_none());
630 assert_eq!(2, buffer.len(), "peeking does not remove");
631
632 assert_eq!(101, sequence_of(&buffer.pop_at_sequence(101).expect("pop")));
633 assert!(buffer.peek_at_sequence(101).is_none());
634 assert!(buffer.pop_at_sequence(101).is_none(), "already taken");
635 assert_eq!(1, buffer.len());
636 }
637
638 #[test]
641 fn addressing_by_wire_sequence_number_works_across_a_wrap() {
642 let mut buffer = JitterBuffer::new(64);
643 for sequence_number in [65534u16, 65535, 0, 1] {
644 buffer.push(packet(1, sequence_number, 0)).expect("push");
645 }
646
647 assert_eq!(
648 0,
649 sequence_of(buffer.peek_at_sequence(0).expect("post-wrap"))
650 );
651 assert_eq!(
652 65535,
653 sequence_of(buffer.peek_at_sequence(65535).expect("pre-wrap"))
654 );
655 assert_eq!(0, sequence_of(&buffer.pop_at_sequence(0).expect("pop 0")));
656 assert_eq!(vec![65534, 65535, 1], drain(&mut buffer));
657 }
658
659 #[test]
662 fn packets_are_addressable_by_rtp_timestamp() {
663 let mut buffer = JitterBuffer::new(64);
664 buffer.push(packet(1, 10, 9000)).expect("push");
665 buffer.push(packet(1, 11, 9000)).expect("push"); buffer.push(packet(1, 12, 12000)).expect("push"); assert_eq!(
669 10,
670 sequence_of(&buffer.pop_at_timestamp(9000).expect("first of the frame")),
671 "the lowest sequence number of that timestamp comes out first"
672 );
673 assert_eq!(
674 11,
675 sequence_of(&buffer.pop_at_timestamp(9000).expect("second of the frame"))
676 );
677 assert!(
678 buffer.pop_at_timestamp(9000).is_none(),
679 "the frame is fully released"
680 );
681 assert!(buffer.pop_at_timestamp(4242).is_none(), "no such timestamp");
682
683 assert_eq!(1, buffer.len(), "the next frame is untouched");
684 assert_eq!(vec![12], drain(&mut buffer));
685 }
686
687 #[test]
690 fn addressed_removal_also_marks_the_position_released() {
691 let mut buffer = JitterBuffer::new(64);
692 buffer.push(packet(1, 5, 0)).expect("push");
693 buffer.push(packet(1, 6, 0)).expect("push");
694
695 buffer.pop_at_sequence(6).expect("pop 6");
696 assert_eq!(
697 Err(Rejected::Late),
698 buffer.push(packet(1, 6, 0)),
699 "6 has been played out"
700 );
701 assert_eq!(
702 Err(Rejected::Late),
703 buffer.push(packet(1, 5, 0)),
704 "and so has everything before it"
705 );
706 }
707
708 #[test]
709 fn reset_clears_the_contents_but_keeps_the_stream_and_counters() {
710 let mut buffer = JitterBuffer::new(64);
711 buffer.push(packet(1, 5, 0)).expect("push");
712 buffer.begin_emitting();
713 buffer.pop().expect("release");
714 buffer.push(packet(1, 5, 0)).ok(); buffer.reset();
717
718 assert!(buffer.is_empty());
719 assert_eq!(Some(1), buffer.ssrc(), "still this stream's buffer");
720 assert_eq!(1, buffer.stats().late, "history is not erased");
721 assert_eq!(
722 State::Buffering,
723 buffer.state(),
724 "a restarted stream re-accumulates before playout resumes"
725 );
726 assert!(buffer.push(packet(1, 5, 0)).is_ok());
728 }
729}