Skip to content
File

Blob: firmware/vendor/str0m/src/packet/buffer_rx.rs

rust1174 lines
1use std::collections::VecDeque;
2use std::fmt;
3use std::ops::{Range, RangeInclusive};
4use std::sync::Arc;
5use std::time::Instant;
6 
7use crate::rtp::vla::VideoLayersAllocation;
8use crate::rtp_::{ExtensionValues, MediaTime, RtpHeader, SenderInfo, SeqNo};
9 
10use super::contiguity::Contiguity;
11use super::contiguity_vp8::Vp8Contiguity;
12use super::contiguity_vp9::Vp9Contiguity;
13use super::{CodecDepacketizer, CodecExtra, Depacketizer, PacketError};
14 
15#[derive(Clone, PartialEq, Eq)]
16/// Holds metadata incoming RTP data.
17pub struct RtpMeta {
18 /// When this RTP packet was received.
19 pub received: Instant,
20 /// Media time translated from the RtpHeader time.
21 pub time: MediaTime,
22 /// Sequence number, extended from the RTPHeader.
23 pub seq_no: SeqNo,
24 /// The actual header.
25 pub header: RtpHeader,
26 /// Sender information from the most recent Sender Report(SR).
27 ///
28 /// If no Sender Report(SR) has been received this is [`None`].
29 pub last_sender_info: Option<SenderInfo>,
30}
31 
32#[derive(Clone)]
33pub struct Depacketized {
34 pub time: MediaTime,
35 pub contiguous: bool,
36 pub meta: Vec<RtpMeta>,
37 pub data: Vec<u8>,
38 pub codec_extra: CodecExtra,
39}
40 
41impl Depacketized {
42 pub fn first_network_time(&self) -> Instant {
43 self.meta
44 .iter()
45 .map(|m| m.received)
46 .min()
47 .expect("a depacketized to consist of at least one packet")
48 }
49 
50 pub fn first_sender_info(&self) -> Option<SenderInfo> {
51 self.meta
52 .iter()
53 .min_by_key(|m| m.received)
54 .map(|m| m.last_sender_info)
55 .expect("a depacketized to consist of at least one packet")
56 }
57 
58 pub fn seq_range(&self) -> RangeInclusive<SeqNo> {
59 let first = self
60 .meta
61 .first()
62 .expect("a depacketized to consist of at least one packet")
63 .seq_no;
64 let last = self
65 .meta
66 .last()
67 .expect("a depacketized to consist of at least one packet")
68 .seq_no;
69 first..=last
70 }
71 
72 pub fn start_of_talkspurt(&self) -> bool {
73 self.meta
74 .first()
75 .expect("a depacketized to consist of at least one packet")
76 .header
77 .marker
78 }
79 
80 pub fn ext_vals(&self) -> ExtensionValues {
81 let last = &self
82 .meta
83 .last()
84 .expect("depacketized video frame must contain a trailing packet")
85 .header
86 .ext_vals;
87 
88 let first = &self
89 .meta
90 .first()
91 .expect("depacketized video frame must contain a leading packet")
92 .header
93 .ext_vals;
94 
95 // We use the extensions from the last packet because certain extensions, such as video
96 // orientation, are only added on the last packet to save bytes.
97 let mut merged = last.clone();
98 
99 // str0m strictly attaches some fields to the first packet of a frame.
100 if let Some(first_val) = &first.abs_capture_time {
101 merged.abs_capture_time = Some(*first_val);
102 }
103 
104 if let Some(first_val) = first.user_values.get_arc::<VideoLayersAllocation>() {
105 merged.user_values.set_arc(first_val);
106 }
107 
108 merged
109 }
110}
111 
112#[derive(Debug)]
113struct Entry {
114 meta: RtpMeta,
115 data: Arc<[u8]>,
116 head: bool,
117 tail: bool,
118}
119 
120#[derive(Debug)]
121pub struct DepacketizingBuffer {
122 hold_back: usize,
123 depack: CodecDepacketizer,
124 queue: VecDeque<Entry>,
125 segments: Vec<(usize, usize)>,
126 last_emitted: Option<(SeqNo, CodecExtra)>,
127 max_time: Option<MediaTime>,
128 depack_cache: Option<(Range<usize>, Depacketized)>,
129 contiguity: Contiguity,
130}
131 
132impl DepacketizingBuffer {
133 pub(crate) fn new(depack: CodecDepacketizer, hold_back: usize) -> Self {
134 let contiguity = match depack {
135 CodecDepacketizer::Vp8(_) => Contiguity::Vp8(Vp8Contiguity::new()),
136 CodecDepacketizer::Vp9(_) => Contiguity::Vp9(Vp9Contiguity::new()),
137 CodecDepacketizer::H264(_)
138 | CodecDepacketizer::H265(_)
139 | CodecDepacketizer::H266(_)
140 | CodecDepacketizer::Av1(_)
141 | CodecDepacketizer::Boxed(_)
142 | CodecDepacketizer::Opus(_)
143 | CodecDepacketizer::ComfortNoise(_)
144 | CodecDepacketizer::G711(_)
145 | CodecDepacketizer::Null(_) => Contiguity::None,
146 };
147 
148 DepacketizingBuffer {
149 hold_back,
150 depack,
151 queue: VecDeque::new(),
152 segments: Vec::new(),
153 last_emitted: None,
154 max_time: None,
155 depack_cache: None,
156 contiguity,
157 }
158 }
159 
160 pub fn push(&mut self, meta: RtpMeta, data: impl Into<Arc<[u8]>>) {
161 self.push_entry(meta, data.into(), None);
162 }
163 
164 pub(crate) fn push_padding(&mut self, meta: RtpMeta) {
165 self.push_entry(meta, Arc::from([]), Some((false, false)));
166 }
167 
168 fn push_entry(&mut self, meta: RtpMeta, data: Arc<[u8]>, partition: Option<(bool, bool)>) {
169 // We're not emitting frames in the wrong order. If we receive
170 // packets that are before the last emitted, we drop.
171 //
172 // As a special case, per popular demand, if hold_back is 0, we do emit
173 // out of order packets.
174 if let Some((last, _)) = self.last_emitted {
175 if meta.seq_no <= last && self.hold_back > 0 {
176 trace!("Drop before emitted: {} <= {}", meta.seq_no, last);
177 return;
178 }
179 }
180 
181 // Record that latest seen max time (used for extending time to u64).
182 self.max_time = Some(if let Some(m) = self.max_time {
183 m.max(meta.time)
184 } else {
185 meta.time
186 });
187 
188 match self
189 .queue
190 .binary_search_by_key(&meta.seq_no, |r| r.meta.seq_no)
191 {
192 Ok(_) => {
193 // exact same seq_no found. ignore
194 trace!("Drop exactly same packet: {}", meta.seq_no);
195 }
196 Err(i) => {
197 let (head, tail) = partition.unwrap_or_else(|| {
198 (
199 self.depack.is_partition_head(data.as_ref()),
200 self.depack
201 .is_partition_tail(meta.header.marker, data.as_ref()),
202 )
203 });
204 
205 // i is insertion point to maintain order
206 let entry = Entry {
207 meta,
208 data,
209 head,
210 tail,
211 };
212 self.queue.insert(i, entry);
213 }
214 }
215 }
216 
217 pub fn pop(&mut self) -> Option<Result<Depacketized, PacketError>> {
218 self.update_segments();
219 
220 if self.segments.is_empty() {
221 self.discard_old_padding();
222 return None;
223 }
224 
225 // println!(
226 // "{:?} {:?}",
227 // self.queue.iter().map(|e| e.meta.seq_no).collect::<Vec<_>>(),
228 // self.segments
229 // );
230 
231 let (start, stop) = *self.segments.first().expect("segment exists");
232 
233 let seq = {
234 let last = self.queue.get(stop).expect("entry for stop index");
235 last.meta.seq_no
236 };
237 
238 // depack ahead, even if we may not emit right away
239 let mut dep = match self.depacketize(start, stop, seq) {
240 Ok(d) => d,
241 Err(e) => {
242 // this segment cannot be decoded correctly
243 // remove from the queue and return the error
244 self.last_emitted = Some((seq, CodecExtra::None));
245 self.queue.drain(0..=stop);
246 return Some(Err(e));
247 }
248 };
249 
250 // If we have contiguity of seq numbers we emit right away,
251 // Otherwise, we wait for retransmissions up to `hold_back` frames
252 // and re-evaluate contiguity based on codec specific information
253 
254 let more_than_hold_back = self.segments.len() >= self.hold_back;
255 let contiguous_seq = self.is_following_last(start);
256 let wait_for_contiguity = !contiguous_seq && !more_than_hold_back;
257 
258 if wait_for_contiguity {
259 // if we are not sending, cache the depacked
260 self.depack_cache = Some((start..stop, dep));
261 self.discard_old_padding();
262 return None;
263 }
264 
265 let (can_emit, contiguous_codec) = self.contiguity.check(&dep.codec_extra, contiguous_seq);
266 dep.contiguous = contiguous_codec;
267 
268 let last = self
269 .queue
270 .get(stop)
271 .expect("entry for stop index")
272 .meta
273 .seq_no;
274 
275 // We're not going to emit frames in the incorrect order, there's no point in keeping
276 // stuff before the emitted range.
277 self.queue.drain(0..=stop);
278 
279 if !can_emit {
280 return None;
281 }
282 
283 self.last_emitted = Some((last, dep.codec_extra));
284 
285 Some(Ok(dep))
286 }
287 
288 fn discard_old_padding(&mut self) {
289 let original_len = self.queue.len();
290 let is_padding = |entry: &Entry| entry.data.is_empty() && !entry.head && !entry.tail;
291 
292 if let Some((mut last, extra)) = self.last_emitted {
293 while self.queue.len() > self.hold_back {
294 let entry = self.queue.front().expect("queue exceeds hold back");
295 if !is_padding(entry) {
296 break;
297 }
298 
299 if last.is_next(entry.meta.seq_no) {
300 last = entry.meta.seq_no;
301 }
302 self.queue.pop_front();
303 }
304 
305 self.last_emitted = Some((last, extra));
306 }
307 
308 // An incomplete frame can remain at the front while another payload type
309 // contributes synthetic padding indefinitely. Keep the frame, but retain
310 // only the newest padding needed for the reordering window. Padding removed
311 // from behind media cannot advance last_emitted across that media.
312 let mut excess_padding = self
313 .queue
314 .iter()
315 .filter(|entry| is_padding(entry))
316 .count()
317 .saturating_sub(self.hold_back);
318 self.queue.retain(|entry| {
319 if excess_padding > 0 && is_padding(entry) {
320 excess_padding -= 1;
321 false
322 } else {
323 true
324 }
325 });
326 
327 if self.queue.len() != original_len {
328 self.depack_cache = None;
329 }
330 }
331 
332 fn depacketize(
333 &mut self,
334 start: usize,
335 stop: usize,
336 _seq: SeqNo,
337 ) -> Result<Depacketized, PacketError> {
338 if let Some(cached) = self.depack_cache.take() {
339 if cached.0 == (start..stop) {
340 trace!("depack cache hit for segment start {}", start);
341 return Ok(cached.1);
342 }
343 }
344 
345 let packets_size = self.queue.range(start..=stop).map(|p| p.data.len()).sum();
346 let mut data = self
347 .depack
348 .out_size_hint(packets_size)
349 .map(Vec::with_capacity)
350 .unwrap_or_else(Vec::new);
351 let mut codec_extra = CodecExtra::None;
352 
353 let time = self.queue.get(start).expect("first index exist").meta.time;
354 let mut meta = Vec::with_capacity(stop - start + 1);
355 
356 for entry in self.queue.range_mut(start..=stop) {
357 self.depack
358 .depacketize(entry.data.as_ref(), &mut data, &mut codec_extra)?;
359 meta.push(entry.meta.clone());
360 }
361 
362 Ok(Depacketized {
363 time,
364 contiguous: true, // the caller taking ownership will modify this accordingly
365 meta,
366 data,
367 codec_extra,
368 })
369 }
370 
371 fn update_segments(&mut self) -> Option<(usize, usize)> {
372 self.segments.clear();
373 
374 #[derive(Clone, Copy)]
375 struct Start {
376 index: i64,
377 time: MediaTime,
378 offset: i64,
379 }
380 
381 let mut start: Option<Start> = None;
382 
383 for (index, entry) in self.queue.iter().enumerate() {
384 let index = index as i64;
385 let iseq = *entry.meta.seq_no as i64;
386 let expected_seq = start.map(|s| s.offset.saturating_add(index));
387 
388 let is_expected_seq = expected_seq == Some(iseq);
389 let is_same_timestamp = start.map(|s| s.time) == Some(entry.meta.time);
390 let is_defacto_tail = is_expected_seq && !is_same_timestamp;
391 
392 if start.is_some() && is_defacto_tail {
393 // We found a segment that ended because the timestamp changed without
394 // a gap in the sequence number. The marker bit in the RTP packet is
395 // just indicative, this is the robust fallback.
396 let segment = (start.unwrap().index as usize, index as usize - 1);
397 self.segments.push(segment);
398 start = None;
399 }
400 
401 if start.is_some() && (!is_expected_seq || !is_same_timestamp) {
402 // Not contiguous. Start looking again.
403 start = None;
404 }
405 
406 // Each segment can have multiple is_partition_head() == true, record the first.
407 if start.is_none() && entry.head {
408 start = Some(Start {
409 index,
410 time: entry.meta.time,
411 offset: iseq.saturating_sub(index),
412 });
413 }
414 
415 if start.is_some() && entry.tail {
416 // We found a contiguous sequence of packets ending with something from
417 // the packet (like the RTP marker bit) indicating it's the tail.
418 let segment = (start.unwrap().index as usize, index as usize);
419 self.segments.push(segment);
420 start = None;
421 }
422 }
423 
424 None
425 }
426 
427 fn is_following_last(&self, start: usize) -> bool {
428 let Some((last, _)) = self.last_emitted else {
429 // First time we emit something.
430 return true;
431 };
432 
433 // track sequence numbers are sequential
434 let mut seq = last;
435 
436 // Expect all entries before start to be padding.
437 for entry in self.queue.range(0..start) {
438 if !seq.is_next(entry.meta.seq_no) {
439 // Not a sequence
440 return false;
441 }
442 // for next loop round.
443 seq = entry.meta.seq_no;
444 
445 let is_padding = entry.data.is_empty() && !entry.head && !entry.tail;
446 if !is_padding {
447 return false;
448 }
449 }
450 
451 let start_entry = self.queue.get(start).expect("entry for start index");
452 
453 seq.is_next(start_entry.meta.seq_no)
454 }
455}
456 
457impl fmt::Debug for RtpMeta {
458 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
459 f.debug_struct("RtpMeta")
460 .field("received", &self.received)
461 .field("time", &self.time)
462 .field("seq_no", &self.seq_no)
463 .field("header", &self.header)
464 .finish()
465 }
466}
467 
468impl fmt::Debug for Depacketized {
469 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
470 f.debug_struct("Depacketized")
471 .field("time", &self.time)
472 .field("meta", &self.meta)
473 .field("data", &self.data.len())
474 .finish()
475 }
476}
477 
478#[cfg(test)]
479mod test {
480 use std::time::{Duration, Instant, SystemTime};
481 
482 use super::*;
483 use crate::packet::vp9::Vp9Depacketizer;
484 use crate::rtp::UserExtensionValues;
485 use crate::rtp_::{AbsCaptureTime, Frequency, MediaTime, Pt, Ssrc, VideoOrientation};
486 
487 #[test]
488 fn end_on_marker() {
489 test(&[
490 //
491 (1, 1, &[1], &[]),
492 (2, 1, &[9], &[(1, &[1, 9])]),
493 ])
494 }
495 
496 #[test]
497 fn ext_vals_extracts_from_first_and_last_packet() {
498 let first_time = Instant::now();
499 let abs_capture_time = AbsCaptureTime {
500 capture_time: SystemTime::UNIX_EPOCH + Duration::from_secs(1),
501 clock_offset: None,
502 };
503 let mut vla = UserExtensionValues::default();
504 vla.set(VideoLayersAllocation {
505 current_simulcast_stream_index: 1,
506 simulcast_streams: vec![],
507 });
508 
509 let first_header = RtpHeader {
510 version: 2,
511 has_padding: false,
512 has_extension: true,
513 csrc_count: 0,
514 marker: false,
515 payload_type: Pt::new_with_value(98),
516 sequence_number: 1,
517 timestamp: 100,
518 ssrc: Ssrc::from(42),
519 csrc: [0; 15],
520 ext_vals: ExtensionValues {
521 abs_capture_time: Some(abs_capture_time),
522 user_values: vla.clone(),
523 ..Default::default()
524 },
525 header_len: 0,
526 };
527 
528 let last_header = RtpHeader {
529 ext_vals: ExtensionValues {
530 video_orientation: Some(VideoOrientation::Deg90),
531 ..Default::default()
532 },
533 ..first_header.clone()
534 };
535 
536 let time_value = 100_u64;
537 let dep = Depacketized {
538 time: MediaTime::new(time_value, Frequency::new(90000).unwrap()),
539 contiguous: true,
540 meta: vec![
541 RtpMeta {
542 received: first_time,
543 time: MediaTime::new(time_value, Frequency::new(90000).unwrap()),
544 seq_no: SeqNo::from(1u64),
545 header: first_header,
546 last_sender_info: None,
547 },
548 RtpMeta {
549 received: first_time + Duration::from_millis(1),
550 time: MediaTime::new(time_value, Frequency::new(90000).unwrap()),
551 seq_no: SeqNo::from(2u64),
552 header: last_header,
553 last_sender_info: None,
554 },
555 ],
556 data: Vec::new(),
557 codec_extra: CodecExtra::None,
558 };
559 
560 let merged = dep.ext_vals();
561 assert_eq!(
562 merged.abs_capture_time.unwrap().capture_time,
563 abs_capture_time.capture_time
564 );
565 assert_eq!(
566 merged
567 .user_values
568 .get::<VideoLayersAllocation>()
569 .unwrap()
570 .current_simulcast_stream_index,
571 vla.get::<VideoLayersAllocation>()
572 .unwrap()
573 .current_simulcast_stream_index,
574 );
575 assert_eq!(merged.video_orientation, Some(VideoOrientation::Deg90));
576 }
577 
578 #[test]
579 fn end_on_defacto() {
580 test(&[
581 (1, 1, &[1], &[]),
582 (2, 1, &[2], &[]),
583 (3, 2, &[3], &[(1, &[1, 2])]),
584 ])
585 }
586 
587 #[test]
588 fn skip_padding() {
589 test(&[
590 (1, 1, &[1], &[]),
591 (2, 1, &[9], &[(1, &[1, 9])]),
592 (3, 1, &[], &[]), // padding!
593 (4, 2, &[1], &[]),
594 (5, 2, &[9], &[(2, &[1, 9])]),
595 ])
596 }
597 
598 #[test]
599 fn gap_after_emit() {
600 test(&[
601 (1, 1, &[1], &[]),
602 (2, 1, &[9], &[(1, &[1, 9])]),
603 // gap
604 (4, 2, &[1], &[]),
605 (5, 2, &[9], &[]),
606 ])
607 }
608 
609 #[test]
610 fn gap_after_padding() {
611 test(&[
612 (1, 1, &[1], &[]),
613 (2, 1, &[9], &[(1, &[1, 9])]),
614 (3, 1, &[], &[]), // padding!
615 // gap
616 (5, 2, &[1], &[]),
617 (6, 2, &[9], &[]),
618 ])
619 }
620 
621 #[test]
622 fn single_packets() {
623 test(&[
624 (1, 1, &[1, 9], &[(1, &[1, 9])]),
625 (2, 2, &[1, 9], &[(2, &[1, 9])]),
626 (3, 3, &[1, 9], &[(3, &[1, 9])]),
627 (4, 4, &[1, 9], &[(4, &[1, 9])]),
628 ])
629 }
630 
631 #[test]
632 fn packets_out_of_order() {
633 test(&[
634 (1, 1, &[1], &[]),
635 (2, 1, &[9], &[(1, &[1, 9])]),
636 (4, 2, &[9], &[]),
637 (3, 2, &[1], &[(2, &[1, 9])]),
638 ])
639 }
640 
641 #[test]
642 fn packets_after_hold_out() {
643 test(&[
644 (1, 1, &[1, 9], &[(1, &[1, 9])]),
645 (3, 3, &[1, 9], &[]),
646 (4, 4, &[1, 9], &[]),
647 (5, 5, &[1, 9], &[(3, &[1, 9]), (4, &[1, 9]), (5, &[1, 9])]),
648 ])
649 }
650 
651 #[test]
652 fn packets_with_hold_0() {
653 test0(&[
654 (1, 1, &[1, 9], &[(1, &[1, 9])]),
655 (3, 3, &[1, 9], &[(3, &[1, 9])]),
656 (4, 4, &[1, 9], &[(4, &[1, 9])]),
657 (5, 5, &[1, 9], &[(5, &[1, 9])]),
658 ])
659 }
660 
661 #[test]
662 fn out_of_order_packets_with_hold_0() {
663 test0(&[
664 (3, 1, &[1, 9], &[(1, &[1, 9])]),
665 (1, 3, &[1, 9], &[(3, &[1, 9])]),
666 (5, 4, &[1, 9], &[(4, &[1, 9])]),
667 (2, 5, &[1, 9], &[(5, &[1, 9])]),
668 ])
669 }
670 
671 #[test]
672 fn padding_only_run_does_not_grow_the_queue_without_bound() {
673 let depack = CodecDepacketizer::Boxed(Box::new(TestDepack));
674 let mut buf = DepacketizingBuffer::new(depack, 3);
675 let meta = |seq: u64| RtpMeta {
676 received: Instant::now(),
677 seq_no: seq.into(),
678 time: MediaTime::from_90khz(seq),
679 last_sender_info: None,
680 header: RtpHeader {
681 sequence_number: seq as u16,
682 timestamp: seq as u32,
683 ..Default::default()
684 },
685 };
686 
687 // Emit one packet for this PT, then model a long run on another PT.
688 buf.push(meta(1), vec![1, 9]);
689 assert!(buf.pop().is_some());
690 
691 for seq in 2..=1_001 {
692 buf.push_padding(meta(seq));
693 assert!(buf.pop().is_none());
694 }
695 
696 assert!(
697 buf.queue.len() <= buf.hold_back,
698 "padding-only queue retained {} entries after a 1,000-packet run",
699 buf.queue.len()
700 );
701 
702 // Returning to this PT after compaction must still be contiguous with the last packet.
703 buf.push(meta(1_002), vec![1, 9]);
704 let dep = buf.pop().expect("packet emitted").expect("valid packet");
705 assert!(dep.contiguous);
706 }
707 
708 #[test]
709 fn padding_only_run_with_loss_stays_bounded_and_reports_discontinuity() {
710 let depack = CodecDepacketizer::Boxed(Box::new(TestDepack));
711 let mut buf = DepacketizingBuffer::new(depack, 3);
712 let meta = |seq: u64| RtpMeta {
713 received: Instant::now(),
714 seq_no: seq.into(),
715 time: MediaTime::from_90khz(seq),
716 last_sender_info: None,
717 header: RtpHeader {
718 sequence_number: seq as u16,
719 timestamp: seq as u32,
720 ..Default::default()
721 },
722 };
723 
724 buf.push(meta(1), vec![1, 9]);
725 assert!(buf.pop().is_some());
726 
727 // Sequence 2 is lost while another PT remains active.
728 for seq in 3..=1_002 {
729 buf.push_padding(meta(seq));
730 assert!(buf.pop().is_none());
731 }
732 
733 assert!(buf.queue.len() <= buf.hold_back);
734 
735 // A discontinuity waits for the normal hold-back before it is emitted.
736 buf.push(meta(1_003), vec![1, 9]);
737 assert!(buf.pop().is_none());
738 buf.push(meta(1_004), vec![1, 9]);
739 assert!(buf.pop().is_none());
740 buf.push(meta(1_005), vec![1, 9]);
741 let dep = buf.pop().expect("packet emitted").expect("valid packet");
742 assert!(!dep.contiguous);
743 }
744 
745 #[test]
746 fn padding_after_incomplete_vp8_frame_does_not_grow_without_bound() {
747 let depack = CodecDepacketizer::Vp8(Default::default());
748 let mut buf = DepacketizingBuffer::new(depack, 3);
749 let meta = |seq: u64, marker: bool| RtpMeta {
750 received: Instant::now(),
751 seq_no: seq.into(),
752 time: MediaTime::from_90khz(seq),
753 last_sender_info: None,
754 header: RtpHeader {
755 marker,
756 sequence_number: seq as u16,
757 timestamp: seq as u32,
758 ..Default::default()
759 },
760 };
761 
762 // Emit one complete VP8 frame for this PT.
763 buf.push(meta(1, true), [0x10, 0x00]);
764 assert!(buf.pop().is_some());
765 
766 // The next frame's head is lost, leaving an S=0 fragment at the front.
767 buf.push(meta(2, false), [0x00]);
768 assert!(buf.pop().is_none());
769 
770 // The sender switches PT. Every packet becomes synthetic padding here.
771 for seq in 3..=1_002 {
772 buf.push_padding(meta(seq, false));
773 assert!(buf.pop().is_none());
774 }
775 
776 assert!(
777 buf.queue.len() <= buf.hold_back + 1,
778 "incomplete VP8 frame retained {} entries after a 1,000-packet PT switch",
779 buf.queue.len()
780 );
781 }
782 
783 #[test]
784 fn padding_after_waiting_vp8_frame_does_not_grow_without_bound() {
785 let depack = CodecDepacketizer::Vp8(Default::default());
786 let mut buf = DepacketizingBuffer::new(depack, 3);
787 let meta = |seq: u64, marker: bool| RtpMeta {
788 received: Instant::now(),
789 seq_no: seq.into(),
790 time: MediaTime::from_90khz(seq),
791 last_sender_info: None,
792 header: RtpHeader {
793 marker,
794 sequence_number: seq as u16,
795 timestamp: seq as u32,
796 ..Default::default()
797 },
798 };
799 
800 buf.push(meta(1, true), [0x10, 0x00]);
801 assert!(buf.pop().is_some());
802 
803 // Lose a frame head, then switch away long enough to compact its padding.
804 buf.push(meta(2, false), [0x00]);
805 assert!(buf.pop().is_none());
806 for seq in 3..=1_002 {
807 buf.push_padding(meta(seq, false));
808 assert!(buf.pop().is_none());
809 }
810 
811 // One complete frame returns, but waits behind the orphan for hold-back.
812 buf.push(meta(1_003, true), [0x10, 0x00]);
813 assert!(buf.pop().is_none());
814 
815 // Switching away again must remain bounded while that frame waits.
816 for seq in 1_004..=2_003 {
817 buf.push_padding(meta(seq, false));
818 assert!(buf.pop().is_none());
819 }
820 
821 assert!(
822 buf.queue.len() <= buf.hold_back + 2,
823 "waiting VP8 frame retained {} entries after a second 1,000-packet PT switch",
824 buf.queue.len()
825 );
826 }
827 
828 #[test]
829 fn padding_after_initial_incomplete_vp8_frame_does_not_grow_without_bound() {
830 let depack = CodecDepacketizer::Vp8(Default::default());
831 let mut buf = DepacketizingBuffer::new(depack, 3);
832 let meta = |seq: u64| RtpMeta {
833 received: Instant::now(),
834 seq_no: seq.into(),
835 time: MediaTime::from_90khz(seq),
836 last_sender_info: None,
837 header: RtpHeader {
838 sequence_number: seq as u16,
839 timestamp: seq as u32,
840 ..Default::default()
841 },
842 };
843 
844 // Start observing this PT after the head of a VP8 frame was lost.
845 buf.push(meta(1), [0x00]);
846 assert!(buf.pop().is_none());
847 
848 // The sender switches PT before this depayloader has emitted anything.
849 for seq in 2..=1_001 {
850 buf.push_padding(meta(seq));
851 assert!(buf.pop().is_none());
852 }
853 
854 assert!(
855 buf.queue.len() <= buf.hold_back + 1,
856 "initial incomplete VP8 frame retained {} entries after a 1,000-packet PT switch",
857 buf.queue.len()
858 );
859 }
860 
861 fn test(
862 v: &[(
863 u64, // seq
864 u64, // time
865 &[u8], // data
866 &[(
867 u64, // time
868 &[u8], // depacketized data
869 )],
870 )],
871 ) {
872 test_n(3, v)
873 }
874 
875 fn test0(
876 v: &[(
877 u64, // seq
878 u64, // time
879 &[u8], // data
880 &[(
881 u64, // time
882 &[u8], // depacketized data
883 )],
884 )],
885 ) {
886 test_n(0, v)
887 }
888 
889 fn test_n(
890 hold_back: usize,
891 v: &[(
892 u64, // seq
893 u64, // time
894 &[u8], // data
895 &[(
896 u64, // time
897 &[u8], // depacketized data
898 )],
899 )],
900 ) {
901 let depack = CodecDepacketizer::Boxed(Box::new(TestDepack));
902 let mut buf = DepacketizingBuffer::new(depack, hold_back);
903 
904 let mut step = 1;
905 
906 for (seq, time, data, checks) in v {
907 let meta = RtpMeta {
908 received: Instant::now(),
909 seq_no: (*seq).into(),
910 time: MediaTime::from_90khz(*time),
911 last_sender_info: None,
912 header: RtpHeader {
913 sequence_number: *seq as u16,
914 timestamp: *time as u32,
915 ..Default::default()
916 },
917 };
918 
919 buf.push(meta, data.to_vec());
920 
921 let mut depacks = vec![];
922 while let Some(res) = buf.pop() {
923 let d = res.unwrap();
924 depacks.push(d);
925 }
926 
927 assert_eq!(
928 depacks.len(),
929 checks.len(),
930 "Step {}: check count not matching {} != {}",
931 step,
932 depacks.len(),
933 checks.len()
934 );
935 
936 let iter = depacks.into_iter().zip(checks.iter());
937 
938 for (depack, (dtime, ddata)) in iter {
939 assert_eq!(
940 depack.time.numer(),
941 *dtime,
942 "Step {}: Time not matching {} != {}",
943 step,
944 depack.time.numer(),
945 *dtime
946 );
947 
948 assert_eq!(
949 depack.data, *ddata,
950 "Step {}: Data not correct {:?} != {:?}",
951 step, depack.data, *ddata
952 );
953 }
954 
955 step += 1;
956 }
957 }
958 
959 #[derive(Debug)]
960 struct TestDepack;
961 
962 impl Depacketizer for TestDepack {
963 fn out_size_hint(&self, packets_size: usize) -> Option<usize> {
964 Some(packets_size)
965 }
966 
967 fn depacketize(
968 &mut self,
969 packet: &[u8],
970 out: &mut Vec<u8>,
971 _: &mut CodecExtra,
972 ) -> Result<(), PacketError> {
973 out.extend_from_slice(packet);
974 Ok(())
975 }
976 
977 fn is_partition_head(&self, packet: &[u8]) -> bool {
978 !packet.is_empty() && packet[0] == 1
979 }
980 
981 fn is_partition_tail(&self, _marker: bool, packet: &[u8]) -> bool {
982 !packet.is_empty() && packet.contains(&9)
983 }
984 }
985 
986 #[test]
987 fn rtp_out_of_order() {
988 let construct_input =
989 |(time, seq, marker, cc, data): (u32, u16, bool, u16, Vec<u8>)| -> (RtpMeta, Vec<u8>) {
990 (
991 RtpMeta {
992 received: Instant::now(),
993 time: MediaTime::new(time.into(), Frequency::new(90000).unwrap()),
994 seq_no: SeqNo::from(seq as u64),
995 header: RtpHeader {
996 version: 2,
997 has_padding: false,
998 has_extension: true,
999 csrc_count: 0,
1000 marker,
1001 payload_type: Pt::new_with_value(98),
1002 sequence_number: seq,
1003 timestamp: time,
1004 ssrc: Ssrc::from(2930203832),
1005 csrc: [0; 15],
1006 ext_vals: ExtensionValues {
1007 transport_cc: Some(cc),
1008 ..Default::default()
1009 },
1010 header_len: 28,
1011 },
1012 last_sender_info: None,
1013 },
1014 data,
1015 )
1016 };
1017 
1018 let inputs = [
1019 // PID: 23860
1020 (
1021 821395241, // Timestamp
1022 8685, // SeqN
1023 false, // Marker
1024 56, // Transport CC
1025 // VP9 header--------+
1026 // |
1027 // +--------------------+
1028 Vec::from([236, 221, 52, 80, 26, 10, 1, 1, 1, 1, 1, 1, 1, 1]), // Data
1029 ),
1030 // PID: 23860
1031 (
1032 821395241, // Timestamp
1033 8686, // SeqN
1034 true, // Marker
1035 57, // Transport CC
1036 // VP9 header--------+
1037 // |
1038 // +--------------------+
1039 Vec::from([237, 221, 52, 83, 26, 10, 2, 2, 2, 2, 2, 2, 2, 2]), // Data
1040 ),
1041 // PID: 23861
1042 (
1043 821398481, // Timestamp
1044 8687, // SeqN
1045 false, // Marker
1046 60, // Transport CC
1047 // VP9 header--------+
1048 // |
1049 // +------------------------+
1050 Vec::from([170, 221, 53, 16, 27, 56, 20, 0, 0, 0, 0, 0, 0, 0, 0]), // Data
1051 ),
1052 // PID: 23861
1053 (
1054 821398481, // Timestamp
1055 8688, // SeqN
1056 false, // Marker
1057 61, // Transport CC
1058 // VP9 header--------+
1059 // |
1060 // +--------------------+
1061 Vec::from([160, 221, 53, 16, 27, 20, 1, 1, 1, 1, 1, 1, 1, 1]), // Data
1062 ),
1063 // PID: 23861
1064 (
1065 821398481,
1066 8689,
1067 false,
1068 62,
1069 // VP9 header--------+
1070 // |
1071 // +--------------------+
1072 Vec::from([164, 221, 53, 16, 27, 20, 2, 2, 2, 2, 2, 2, 2, 2]), // Data
1073 ),
1074 // PID: 23861
1075 (
1076 821398481, // Timestamp
1077 8690, // SeqN
1078 false, // Marker
1079 63, // Transport CC
1080 // VP9 header--------+
1081 // |
1082 // +--------------------+
1083 Vec::from([169, 221, 53, 19, 27, 20, 3, 3, 3, 3, 3, 3, 3, 3]), // Data
1084 ),
1085 // PID: 23861
1086 (
1087 821398481, // Timestamp
1088 8691, // SeqN
1089 false, // Marker
1090 64, // Transport CC
1091 // VP9 header--------+
1092 // |
1093 // +--------------------+
1094 Vec::from([161, 221, 53, 19, 27, 20, 4, 4, 4, 4, 4, 4, 4, 4]), // Data
1095 ),
1096 // PID: 23861
1097 (
1098 821398481, // Timestamp
1099 8692, // SeqN
1100 false, // Marker
1101 65, // Transport CC
1102 // VP9 header--------+
1103 // |
1104 // +--------------------+
1105 Vec::from([161, 221, 53, 19, 27, 20, 5, 5, 5, 5, 5, 5, 5, 5]), // Data
1106 ),
1107 // PID: 23861
1108 (
1109 821398481, // Timestamp
1110 8693, // SeqN
1111 false, // Marker
1112 66, // Transport CC
1113 // VP9 header--------+
1114 // |
1115 // +--------------------+
1116 Vec::from([161, 221, 53, 19, 27, 20, 6, 6, 6, 6, 6, 6, 6, 6]), // Data
1117 ),
1118 // PID: 23861
1119 (
1120 821398481, // Timestamp
1121 8694, // SeqN
1122 true, // Marker
1123 67, // Transport CC
1124 // VP9 header--------+
1125 // |
1126 // +--------------------+
1127 Vec::from([165, 221, 53, 19, 27, 20, 7, 7, 7, 7, 7, 7, 7, 7]), // Data
1128 ),
1129 ];
1130 
1131 let mut buffer =
1132 DepacketizingBuffer::new(CodecDepacketizer::Vp9(Vp9Depacketizer::default()), 30);
1133 
1134 for input in &inputs {
1135 let (meta, data) = construct_input(input.clone());
1136 buffer.push(meta, data);
1137 }
1138 
1139 let res0before = buffer.pop().unwrap().unwrap(); // Pop PID: 23860, `contiguous_seq == true`.
1140 let res1before = buffer.pop().unwrap().unwrap(); // Pop PID: 23861, `contiguous_seq == true`.
1141 
1142 let mut buffer =
1143 DepacketizingBuffer::new(CodecDepacketizer::Vp9(Vp9Depacketizer::default()), 30);
1144 
1145 for input in &inputs {
1146 let (meta, data) = construct_input(input.clone());
1147 if meta.seq_no == SeqNo::from(8689) {
1148 continue; // Skip RTP packet with seq_num=8689 vp9_payload=[20, 2, 2, 2, 2, 2, 2, 2, 2].
1149 }
1150 buffer.push(meta.clone(), data.clone());
1151 }
1152 
1153 // Pop PID: 23860, `contiguous_seq == true`.
1154 let res0after = buffer.pop().unwrap().unwrap();
1155 // Try to pop PID: 23861. `None` because `contiguous_seq == false` -- no seq_num=8689.
1156 assert!(buffer.pop().is_none());
1157 assert!(buffer.pop().is_none()); // Ensure once again.
1158 
1159 for input in &inputs {
1160 let (meta, data) = construct_input(input.clone());
1161 if meta.seq_no == SeqNo::from(8689) {
1162 // Send RTP packet with seq_num=8689 vp9_payload=[20, 2, 2, 2, 2, 2, 2, 2, 2].
1163 buffer.push(meta.clone(), data.clone());
1164 break;
1165 }
1166 }
1167 
1168 let res1after = buffer.pop().unwrap().unwrap();
1169 
1170 assert_eq!(res0before.data, res0after.data);
1171 assert_eq!(res1before.data, res1after.data);
1172 }
1173}