Skip to content
File

Blob: firmware/vendor/str0m/src/rtp/rtcp/twcc.rs

rust2328 lines
1use std::collections::VecDeque;
2use std::collections::vec_deque;
3use std::fmt;
4use std::ops::RangeInclusive;
5use std::time::{Duration, Instant};
6 
7use super::{FeedbackMessageType, RtcpHeader, RtcpPacket, extend_u16};
8use super::{RtcpType, Ssrc, TransportType};
9 
10use crate::rtp_::{TwccClusterId, TwccSeq};
11 
12/// Transport Wide Congestion Control.
13///
14/// Sent in response to every RTP packet, but does ranges of packets to respond to.
15#[derive(Clone, PartialEq, Eq)]
16pub struct Twcc {
17 /// Sender of this feedback. Mostly irrelevant, but part of RTCP packets.
18 pub sender_ssrc: Ssrc,
19 /// The SSRC this report is for.
20 pub ssrc: Ssrc,
21 /// Start sequence number.
22 pub base_seq: u16,
23 /// Number of reported statuses.
24 pub status_count: u16,
25 /// Clock time this report was produced. Used for RTT measurement.
26 pub reference_time: u32, // 24 bit
27 /// Increasing counter for each TWCC. For deduping.
28 pub feedback_count: u8, // counter for each Twcc
29 /// Ranges received.
30 pub chunks: VecDeque<PacketChunk>,
31 /// Delta times for the ranges received.
32 pub delta: VecDeque<Delta>,
33}
34 
35impl Twcc {
36 fn chunks_byte_len(&self) -> usize {
37 self.chunks.len() * 2
38 }
39 
40 fn delta_byte_len(&self) -> usize {
41 self.delta.iter().map(|d| d.byte_len()).sum()
42 }
43 
44 /// Iterate over the reported sequences.
45 pub fn into_iter(self, time_zero: Instant, extend_from: TwccSeq) -> TwccIter {
46 let millis = self.reference_time as u64 * 64;
47 let time_base = time_zero + Duration::from_millis(millis);
48 let base_seq = extend_u16(Some(*extend_from), self.base_seq);
49 let last_seq = base_seq + self.status_count as u64;
50 
51 TwccIter {
52 base_seq,
53 last_seq,
54 time_base,
55 index: 0,
56 twcc: self,
57 }
58 }
59}
60 
61pub struct TwccIter {
62 base_seq: u64,
63 last_seq: u64,
64 time_base: Instant,
65 index: usize,
66 twcc: Twcc,
67}
68 
69impl Iterator for TwccIter {
70 type Item = (TwccSeq, PacketStatus, Option<Instant>);
71 
72 fn next(&mut self) -> Option<Self::Item> {
73 let seq: TwccSeq = (self.base_seq + self.index as u64).into();
74 
75 if *seq == self.last_seq {
76 return None;
77 }
78 
79 let head = self.twcc.chunks.front()?;
80 
81 let (status, amount) = match head {
82 PacketChunk::Run(s, n) => {
83 use PacketStatus::*;
84 let status = match s {
85 NotReceived | Unknown => NotReceived,
86 ReceivedSmallDelta => ReceivedSmallDelta,
87 PacketStatus::ReceivedLargeOrNegativeDelta => ReceivedLargeOrNegativeDelta,
88 };
89 (status, *n)
90 }
91 PacketChunk::VectorSingle(v, n) => {
92 let status = if 1 << (13 - self.index) & v > 0 {
93 PacketStatus::ReceivedSmallDelta
94 } else {
95 PacketStatus::NotReceived
96 };
97 (status, *n)
98 }
99 PacketChunk::VectorDouble(v, n) => {
100 let e = ((v >> (12 - self.index * 2)) & 0b11) as u8;
101 let status = PacketStatus::from(e);
102 (status, *n)
103 }
104 };
105 
106 let instant = match status {
107 PacketStatus::NotReceived => None,
108 PacketStatus::ReceivedSmallDelta => match self.twcc.delta.pop_front()? {
109 Delta::Small(v) => Some(self.time_base + Duration::from_micros(250 * v as u64)),
110 Delta::Large(_) => panic!("Incorrect large delta size"),
111 },
112 PacketStatus::ReceivedLargeOrNegativeDelta => match self.twcc.delta.pop_front()? {
113 Delta::Small(_) => panic!("Incorrect small delta size"),
114 Delta::Large(v) => {
115 let dur = Duration::from_micros(250 * v.unsigned_abs() as u64);
116 Some(if v < 0 {
117 self.time_base.checked_sub(dur).unwrap_or(self.time_base)
118 } else {
119 self.time_base + dur
120 })
121 }
122 },
123 // A VectorDouble chunk can carry 0b11 for any of its statuses. The
124 // parser maps that to Unknown and consumes no delta, so there is nothing
125 // to report for it either.
126 PacketStatus::Unknown => None,
127 };
128 
129 if let Some(new_timebase) = instant {
130 self.time_base = new_timebase;
131 }
132 
133 self.index += 1;
134 if self.index == amount as usize {
135 self.twcc.chunks.pop_front();
136 self.base_seq = *seq + 1;
137 self.index = 0;
138 }
139 
140 Some((seq, status, instant))
141 }
142}
143 
144impl RtcpPacket for Twcc {
145 fn header(&self) -> RtcpHeader {
146 RtcpHeader {
147 rtcp_type: RtcpType::TransportLayerFeedback,
148 feedback_message_type: FeedbackMessageType::TransportFeedback(
149 TransportType::TransportWide,
150 ),
151 words_less_one: (self.length_words() - 1) as u16,
152 }
153 }
154 
155 fn length_words(&self) -> usize {
156 // header: 1
157 // sender ssrc: 1
158 // ssrc: 1
159 // base seq + packet status: 1
160 // ref time + feedback count: 1
161 // chunks byte len + delta byte len + padding
162 
163 let mut total = self.chunks_byte_len() + self.delta_byte_len();
164 
165 let pad = 4 - total % 4;
166 if pad < 4 {
167 total += pad;
168 }
169 
170 assert!(total % 4 == 0);
171 
172 let total_words = total / 4;
173 
174 5 + total_words
175 }
176 
177 fn write_to(&self, buf: &mut [u8]) -> usize {
178 let len_start = buf.len();
179 
180 let mut total = {
181 let buf = &mut buf[..];
182 
183 self.header().write_to(buf);
184 buf[4..8].copy_from_slice(&self.sender_ssrc.to_be_bytes());
185 buf[8..12].copy_from_slice(&self.ssrc.to_be_bytes());
186 
187 buf[12..14].copy_from_slice(&self.base_seq.to_be_bytes());
188 buf[14..16].copy_from_slice(&self.status_count.to_be_bytes());
189 
190 let ref_time = self.reference_time.to_be_bytes();
191 buf[16..19].copy_from_slice(&ref_time[1..4]);
192 buf[19] = self.feedback_count;
193 
194 let mut buf = &mut buf[20..];
195 for p in &self.chunks {
196 p.write_to(buf);
197 buf = &mut buf[2..];
198 }
199 
200 for d in &self.delta {
201 let n = d.write_to(buf);
202 buf = &mut buf[n..];
203 }
204 
205 len_start - buf.len()
206 };
207 
208 let pad = 4 - total % 4;
209 if pad < 4 {
210 for i in 0..pad {
211 buf[total + i] = 0;
212 }
213 buf[total + pad - 1] = pad as u8;
214 
215 total += pad;
216 // Toggle padding bit
217 buf[0] |= 0b00_1_00000;
218 }
219 
220 total
221 }
222}
223 
224#[derive(Debug)]
225pub struct TwccRecvRegister {
226 // How many packets to keep when they are reported. This is to handle packets arriving out
227 // of order and where two consecutive calls to `build_report` needs to go "backwards" in
228 // base_seq.
229 keep_reported: usize,
230 
231 /// Queue of packets to form Twcc reports of.
232 ///
233 /// Once the queue has some content, we will always keep at least one entry to "remember" for the
234 /// next report.
235 queue: VecDeque<Receipt>,
236 
237 /// Index into queue from where we start reporting on next build_report().
238 report_from: usize,
239 
240 /// Interims built in this for every build_report.
241 interims: VecDeque<ChunkInterim>,
242 
243 /// The point in time we consider 0. All reported values are offset from this. Set to first
244 /// unreported packet in first `build_reported`.
245 ///
246 // https://datatracker.ietf.org/doc/html/draft-holmer-rmcat-transport-wide-cc-extensions-01#page-5
247 // reference time: 24 bits Signed integer indicating an absolute
248 // reference time in some (unknown) time base chosen by the
249 // sender of the feedback packets.
250 time_start: Option<Instant>,
251 
252 /// Counter that increases by one for each report generated.
253 generated_reports: u64,
254 
255 /// Data to calculate received loss.
256 receive_window: ReceiveWindow,
257}
258 
259#[derive(Debug, Clone, Copy, PartialEq, Eq)]
260struct Receipt {
261 seq: TwccSeq,
262 time: Instant,
263}
264 
265impl TwccRecvRegister {
266 pub fn new(keep_reported: usize) -> Self {
267 TwccRecvRegister {
268 keep_reported,
269 queue: VecDeque::new(),
270 report_from: 0,
271 interims: VecDeque::new(),
272 time_start: None,
273 generated_reports: 0,
274 receive_window: ReceiveWindow::default(),
275 }
276 }
277 
278 /// Find the last sequence number, or `None` if the queue is empty.
279 pub fn max_seq(&self) -> Option<TwccSeq> {
280 // The highest seq must be the last since update_seq inserts values
281 // using a binary search.
282 self.queue.back().map(|r| r.seq)
283 }
284 
285 pub fn update_seq(&mut self, seq: TwccSeq, time: Instant) {
286 self.receive_window.record_seq(seq);
287 
288 match self.queue.binary_search_by_key(&seq, |r| r.seq) {
289 Ok(_) => {
290 // Exact same SeqNo found. This is an error where the sender potentially
291 // used the same twcc sequence number for two packets. Let's ignore it.
292 }
293 Err(idx) => {
294 if let Some(time_start) = self.time_start {
295 // If time goes back more than 8192 millis from the time point we've
296 // chosen as our 0, we can't represent that in the report. Let's just
297 // forget about it and hope for the best.
298 if time_start - time >= Duration::from_millis(8192) {
299 return;
300 }
301 }
302 
303 self.queue.insert(idx, Receipt { seq, time });
304 
305 if idx < self.report_from {
306 self.report_from = idx;
307 }
308 }
309 }
310 }
311 
312 pub fn build_report(&mut self, max_byte_size: usize) -> Option<Twcc> {
313 if max_byte_size > 10_000 {
314 warn!("Refuse to build too large Twcc report");
315 return None;
316 }
317 
318 // First unreported is the self.time_start relative offset of the next Twcc.
319 let first = self.queue.get(self.report_from);
320 let first = first?;
321 
322 // Set once on first ever built report.
323 if self.time_start.is_none() {
324 self.time_start = Some(first.time);
325 }
326 
327 let (base_seq, first_time) = (first.seq, first.time);
328 let time_start = self.time_start.expect("a start time");
329 
330 // The difference between our Twcc reference time and the first ever report start time.
331 let first_time_rel = first_time - time_start;
332 
333 // The value is to be interpreted in multiples of 64ms.
334 let reference_time = (first_time_rel.as_micros() as u64 / 64_000) as u32;
335 
336 let mut twcc = Twcc {
337 sender_ssrc: 0.into(),
338 ssrc: 0.into(),
339 feedback_count: self.generated_reports as u8,
340 base_seq: *base_seq as u16,
341 reference_time,
342 status_count: 0,
343 chunks: VecDeque::new(),
344 delta: VecDeque::new(),
345 };
346 
347 // Because reference time is in steps of 64ms, the first reported packet might have an
348 // offset (packet time resolution is 250us). This base_time is calculated backwards from
349 // reference time so that we can offset all packets from the "truncated" 64ms steps.
350 // The RFC says:
351 // The first recv delta in this packet is relative to the reference time.
352 let base_time = time_start + Duration::from_micros(reference_time as u64 * 64_000);
353 
354 // The ChunkInterim are helpers structures that hold the deltas between
355 // the registered receptions.
356 build_interims(
357 &self.queue,
358 self.report_from,
359 base_seq,
360 base_time,
361 &mut self.interims,
362 );
363 let interims = &mut self.interims;
364 
365 // 20 bytes is the size of the fixed fields in Twcc.
366 let mut bytes_left = max_byte_size - 20;
367 
368 while !interims.is_empty() {
369 // 2 chunk + 2 large delta + 3 padding
370 const MIN_RUN_SIZE: usize = 2 + 2 + 3;
371 
372 if bytes_left < MIN_RUN_SIZE {
373 break;
374 }
375 
376 // Chose the packet chunk type that can fit the most interims.
377 let (mut chunk, max) = {
378 let first_status = interims.front().expect("at least one interim").status();
379 
380 let c_run = PacketChunk::Run(first_status, 0);
381 let c_single = PacketChunk::VectorSingle(0, 0);
382 let c_double = PacketChunk::VectorDouble(0, 0);
383 
384 let max_run = c_run.append_max(interims.iter());
385 let max_single = c_single.append_max(interims.iter());
386 let max_double = c_double.append_max(interims.iter());
387 
388 let max = max_run.max(max_single).max(max_double);
389 
390 // 2 chunk + 14 small delta + 3 padding
391 const MAX_SINGLE_SIZE: usize = 2 + 14 + 3;
392 // 2 chunk + 7 large delta + 3 padding
393 const MAX_DOUBLE_SIZE: usize = 2 + 14 + 3;
394 
395 if max == max_run {
396 (c_run, max_run)
397 } else if max == max_single && bytes_left >= MAX_SINGLE_SIZE {
398 (c_single, max_single)
399 } else if max == max_double && bytes_left >= MAX_DOUBLE_SIZE {
400 (c_double, max_double)
401 } else {
402 // fallback, since we can always do runs.
403 (c_run, max_run)
404 }
405 };
406 
407 // we should _definitely_ be able to fit this many reported.
408 let mut todo = max;
409 
410 loop {
411 if bytes_left < MIN_RUN_SIZE {
412 break;
413 }
414 
415 if todo == 0 {
416 break;
417 }
418 
419 let i = match interims.front_mut() {
420 Some(v) => v,
421 None => break,
422 };
423 
424 let appended = chunk.append(i);
425 assert!(appended > 0);
426 todo -= appended;
427 twcc.status_count += appended;
428 
429 if i.consume(appended) {
430 // it was fully consumed.
431 if matches!(i, ChunkInterim::Received(_, _)) {
432 self.report_from += 1;
433 }
434 
435 if let Some(delta) = i.delta() {
436 twcc.delta.push_back(delta);
437 bytes_left -= delta.byte_len();
438 }
439 
440 // move on to next interim
441 interims.pop_front();
442 } else {
443 // not fully consumed, then we must have run out of space in the chunk.
444 assert!(todo == 0);
445 }
446 }
447 
448 let free = chunk.free();
449 if chunk.must_be_full() && free > 0 {
450 // this must be at the end where we can shift in missing
451 assert!(interims.is_empty());
452 chunk.append(&ChunkInterim::Missing(free));
453 }
454 
455 twcc.chunks.push_back(chunk);
456 bytes_left -= 2;
457 }
458 
459 // libWebRTC demands at least one chunk, or it will warn with
460 // "Buffer too small (16 bytes) to fit a FeedbackPacket. Minimum size = 18"
461 // (18 bytes here is not including the RTCP header).
462 if twcc.chunks.is_empty() {
463 return None;
464 }
465 
466 self.generated_reports += 1;
467 
468 // clean up
469 if self.report_from > self.keep_reported {
470 let to_remove = self.report_from - self.keep_reported;
471 self.queue.drain(..to_remove);
472 self.report_from -= to_remove;
473 }
474 
475 Some(twcc)
476 }
477 
478 pub fn has_unreported(&self) -> bool {
479 self.queue.len() > self.report_from
480 }
481 
482 /// Calculate the fraction of lost packets since the last call.
483 ///
484 /// To get periodic stats call this method at fixed intervals.
485 pub fn loss(&mut self) -> Option<f32> {
486 // Based on the algorithm described in
487 // [RFC 3550 Appendix A](https://www.rfc-editor.org/rfc/rfc3550#appendix-A.3), but instead
488 // of applying it to an individual RTP stream it's applied to the whole session using TWCC
489 // sequence numbers rather than RTP sequence numbers.
490 let max_seq = self.receive_window.max_seq?;
491 let base_seq = self.receive_window.base_seq?;
492 let expected = *max_seq - *base_seq + 1_u64;
493 
494 let expected_interval = expected - self.receive_window.expected_prior;
495 self.receive_window.expected_prior = expected;
496 
497 let received_interval = self.receive_window.received - self.receive_window.received_prior;
498 self.receive_window.received_prior = self.receive_window.received;
499 
500 let lost_interval = expected_interval.saturating_sub(received_interval);
501 
502 (expected_interval != 0).then_some(lost_interval as f32 / expected_interval as f32)
503 }
504}
505 
506/// Interims are deltas between `Receiption` which is an intermediary format before
507/// we populate the Twcc report.
508fn build_interims(
509 queue: &VecDeque<Receipt>,
510 report_from: usize,
511 base_seq: TwccSeq,
512 base_time: Instant,
513 interims: &mut VecDeque<ChunkInterim>,
514) {
515 interims.clear();
516 let report_from = queue.iter().skip(report_from);
517 
518 let mut prev = (base_seq, base_time);
519 
520 for r in report_from {
521 let diff_seq = *r.seq - *prev.0;
522 
523 if diff_seq > 1 {
524 let mut todo = diff_seq - 1;
525 while todo > 0 {
526 // max 2^13 run length in each missing chunk
527 let n = todo.min(8192);
528 interims.push_back(ChunkInterim::Missing(n as u16));
529 todo -= n;
530 }
531 }
532 
533 let diff_time = if r.time < prev.1 {
534 // negative
535 let dur = prev.1 - r.time;
536 -(dur.as_micros() as i64)
537 } else {
538 let dur = r.time - prev.1;
539 dur.as_micros() as i64
540 };
541 
542 let (status, time) = if diff_time < -8_192_000 || diff_time > 8_191_750 {
543 // This is too large to be representable in deltas.
544 // Abort, make a report of what we got, and start anew.
545 break;
546 } else if diff_time < 0 || diff_time > 63_750 {
547 let t = diff_time / 250;
548 assert!(t >= -32_768 && t <= 32_767);
549 (PacketStatus::ReceivedLargeOrNegativeDelta, t as i16)
550 } else {
551 let t = diff_time / 250;
552 assert!(t >= 0 && t <= 255);
553 (PacketStatus::ReceivedSmallDelta, t as i16)
554 };
555 
556 interims.push_back(ChunkInterim::Received(status, time));
557 prev = (r.seq, r.time);
558 }
559}
560 
561#[derive(Debug, Clone, Copy)]
562enum ChunkInterim {
563 Missing(u16), // max 2^13 (one run length)
564 Received(PacketStatus, i16),
565}
566 
567#[derive(Debug, Default)]
568struct ReceiveWindow {
569 /// The base seq num, set on the first receive.
570 base_seq: Option<TwccSeq>,
571 
572 /// The large seq num received.
573 max_seq: Option<TwccSeq>,
574 
575 /// The previous number of packets expected, used to calculate a delta for loss.
576 expected_prior: u64,
577 
578 /// The total number of packets received
579 received: u64,
580 
581 /// The previous number of packets received, used to calculate a delta for loss.
582 received_prior: u64,
583}
584 
585impl ReceiveWindow {
586 fn record_seq(&mut self, seq: TwccSeq) {
587 if self.base_seq.is_none() {
588 self.base_seq = Some(seq);
589 }
590 
591 self.received += 1;
592 self.max_seq = self.max_seq.max(Some(seq));
593 }
594}
595 
596impl ChunkInterim {
597 fn status(&self) -> PacketStatus {
598 match self {
599 ChunkInterim::Missing(_) => PacketStatus::NotReceived,
600 ChunkInterim::Received(s, _) => *s,
601 }
602 }
603 
604 fn delta(&self) -> Option<Delta> {
605 match self {
606 ChunkInterim::Missing(_) => None,
607 ChunkInterim::Received(s, d) => match *s {
608 PacketStatus::ReceivedSmallDelta => Some(Delta::Small(*d as u8)),
609 PacketStatus::ReceivedLargeOrNegativeDelta => Some(Delta::Large(*d)),
610 _ => unreachable!(),
611 },
612 }
613 }
614 
615 fn consume(&mut self, n: u16) -> bool {
616 match self {
617 ChunkInterim::Missing(c) => {
618 *c -= n;
619 *c == 0
620 }
621 ChunkInterim::Received(_, _) => {
622 assert!(n <= 1);
623 n == 1
624 }
625 }
626 }
627}
628 
629#[derive(Debug, Clone, Copy, PartialEq, Eq)]
630pub enum PacketChunk {
631 Run(PacketStatus, u16), // 13 bit repeat
632 VectorSingle(u16, u16),
633 VectorDouble(u16, u16),
634}
635 
636impl PacketChunk {
637 fn append_max<'a>(&self, iter: impl Iterator<Item = &'a ChunkInterim>) -> u16 {
638 let mut to_fill = *self;
639 
640 let mut reached_end = true;
641 
642 for i in iter {
643 if to_fill.free() == 0 {
644 reached_end = false;
645 break;
646 }
647 
648 // The status is not possible to add in this chunk. This could be
649 // a large delta in a single, or a mismatching run.
650 if !to_fill.can_append_status(i.status()) {
651 reached_end = false;
652 break;
653 }
654 
655 to_fill.append(i);
656 }
657 
658 // As a special case, single/double must be completely filled. However if
659 // we reached the end of the interims, we can shift in "missing" to make
660 // them full.
661 if to_fill.must_be_full() && to_fill.free() > 0 && !reached_end {
662 return 0;
663 }
664 
665 self.free() - to_fill.free()
666 }
667 
668 fn append(&mut self, i: &ChunkInterim) -> u16 {
669 use ChunkInterim::*;
670 use PacketChunk::*;
671 let free = self.free();
672 match (self, i) {
673 (Run(s, n), Missing(c)) => {
674 if *s != PacketStatus::NotReceived {
675 return 0;
676 }
677 let max = free.min(*c);
678 *n += max;
679 max
680 }
681 (Run(s, n), Received(s2, _)) => {
682 if *s != *s2 {
683 return 0;
684 }
685 let max = free.min(1);
686 *n += max;
687 max
688 }
689 (VectorSingle(n, f), Missing(c)) => {
690 let max = free.min(*c);
691 *n <<= max;
692 *f += max;
693 max
694 }
695 (VectorSingle(n, f), Received(s2, _)) => {
696 if *s2 == PacketStatus::ReceivedLargeOrNegativeDelta {
697 return 0;
698 }
699 let max = free.min(1);
700 if max == 1 {
701 *n <<= 1;
702 *n |= 1;
703 *f += 1;
704 }
705 max
706 }
707 (VectorDouble(n, f), Missing(c)) => {
708 let max = free.min(*c);
709 *n <<= max * 2;
710 *f += max;
711 max
712 }
713 (VectorDouble(n, f), Received(s2, _)) => {
714 let max = free.min(1);
715 if max == 1 {
716 *n <<= 2;
717 *n |= *s2 as u16;
718 *f += 1;
719 }
720 max
721 }
722 }
723 }
724 
725 fn must_be_full(&self) -> bool {
726 match self {
727 PacketChunk::Run(_, _) => false,
728 PacketChunk::VectorSingle(_, _) => true,
729 PacketChunk::VectorDouble(_, _) => true,
730 }
731 }
732 
733 fn free(&self) -> u16 {
734 match self {
735 PacketChunk::Run(_, n) => 8192 - *n,
736 PacketChunk::VectorSingle(_, filled) => 14 - *filled,
737 PacketChunk::VectorDouble(_, filled) => 7 - *filled,
738 }
739 }
740 
741 fn write_to(&self, buf: &mut [u8]) {
742 let x = match self {
743 // 0 1
744 // 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5
745 // +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
746 // |T| S | Run Length |
747 // +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
748 // chunk type (T): 1 bit A zero identifies this as a run length chunk.
749 // packet status symbol (S): 2 bits The symbol repeated in this run.
750 // See above.
751 // run length (L): 13 bits An unsigned integer denoting the run length.
752 PacketChunk::Run(s, n) => {
753 let mut x = 0_u16;
754 x |= (*s as u16) << 13;
755 assert!(*n <= 8192);
756 x |= n;
757 x
758 }
759 
760 // Corrected according to email exchange at the bottom..
761 //
762 // 0 1
763 // 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5
764 // +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
765 // |T|S| symbol list |
766 // +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
767 // chunk type (T): 1 bit A one identifies this as a status vector
768 // chunk.
769 // symbol size (S): 1 bit A zero means this vector contains only
770 // "packet received" (1) and "packet not received" (0)
771 // symbols. This means we can compress each symbol to just
772 // one bit, 14 in total. A one means this vector contains
773 // the normal 2-bit symbols, 7 in total.
774 // symbol list: 14 bits A list of packet status symbols, 7 or 14 in
775 // total.
776 PacketChunk::VectorSingle(n, fill) => {
777 assert!(*fill == 14);
778 let mut x: u16 = 1 << 15;
779 assert!(*n <= 16384);
780 x |= *n;
781 x
782 }
783 PacketChunk::VectorDouble(n, fill) => {
784 assert!(*fill == 7);
785 let mut x: u16 = 1 << 15;
786 assert!(*n <= 16384);
787 x |= 1 << 14;
788 x |= *n;
789 x
790 }
791 };
792 buf[..2].copy_from_slice(&x.to_be_bytes());
793 }
794 
795 fn max_possible_status_count(&self) -> usize {
796 match self {
797 PacketChunk::Run(_, n) => *n as usize,
798 PacketChunk::VectorSingle(_, _) => 14,
799 PacketChunk::VectorDouble(_, _) => 7,
800 }
801 }
802 
803 fn can_append_status(&self, status: PacketStatus) -> bool {
804 match self {
805 PacketChunk::Run(s, _) => *s == status,
806 PacketChunk::VectorSingle(_, _) => status != PacketStatus::ReceivedLargeOrNegativeDelta,
807 PacketChunk::VectorDouble(_, _) => true,
808 }
809 }
810}
811 
812impl Delta {
813 fn write_to(&self, buf: &mut [u8]) -> usize {
814 match self {
815 Delta::Small(v) => {
816 buf[0] = *v;
817 1
818 }
819 Delta::Large(v) => {
820 buf[..2].copy_from_slice(&v.to_be_bytes());
821 2
822 }
823 }
824 }
825 
826 fn byte_len(&self) -> usize {
827 match self {
828 Delta::Small(_) => 1,
829 Delta::Large(_) => 2,
830 }
831 }
832}
833 
834#[derive(Debug, Clone, Copy, PartialEq, Eq)]
835pub enum PacketStatus {
836 NotReceived = 0b00,
837 ReceivedSmallDelta = 0b01,
838 ReceivedLargeOrNegativeDelta = 0b10,
839 Unknown = 0b11,
840}
841 
842#[derive(Debug, Clone, Copy, PartialEq, Eq)]
843pub enum Delta {
844 Small(u8),
845 Large(i16),
846}
847 
848impl From<PacketStatus> for u8 {
849 fn from(val: PacketStatus) -> Self {
850 val as usize as u8
851 }
852}
853 
854impl From<u8> for PacketStatus {
855 fn from(v: u8) -> Self {
856 match v {
857 0b00 => Self::NotReceived,
858 0b01 => Self::ReceivedSmallDelta,
859 0b10 => Self::ReceivedLargeOrNegativeDelta,
860 _ => Self::Unknown,
861 }
862 }
863}
864 
865impl<'a> TryFrom<&'a [u8]> for Twcc {
866 type Error = &'static str;
867 
868 fn try_from(buf: &'a [u8]) -> Result<Self, Self::Error> {
869 if buf.len() < 16 {
870 return Err("Less than 16 bytes for start of Twcc");
871 }
872 
873 let sender_ssrc = u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]).into();
874 let ssrc = u32::from_be_bytes([buf[4], buf[5], buf[6], buf[7]]).into();
875 let base_seq = u16::from_be_bytes([buf[8], buf[9]]);
876 let status_count = u16::from_be_bytes([buf[10], buf[11]]);
877 let reference_time = u32::from_be_bytes([0, buf[12], buf[13], buf[14]]);
878 let feedback_count = buf[15];
879 
880 let mut twcc = Twcc {
881 sender_ssrc,
882 ssrc,
883 base_seq,
884 status_count,
885 reference_time,
886 feedback_count,
887 chunks: VecDeque::new(),
888 delta: VecDeque::new(),
889 };
890 
891 let mut todo = status_count as isize;
892 let mut buf = &buf[16..];
893 loop {
894 if todo <= 0 {
895 break;
896 }
897 
898 let chunk: PacketChunk = buf.try_into()?;
899 
900 todo -= chunk.max_possible_status_count() as isize;
901 
902 twcc.chunks.push_back(chunk);
903 buf = &buf[2..];
904 }
905 
906 if twcc.chunks.is_empty() {
907 return Ok(twcc);
908 }
909 
910 fn read_delta_small(
911 buf: &[u8],
912 n: usize,
913 ) -> Result<impl Iterator<Item = Delta> + '_, &'static str> {
914 if buf.len() < n {
915 return Err("Not enough buf for small deltas");
916 }
917 Ok((0..n).map(|i| Delta::Small(buf[i])))
918 }
919 
920 fn read_delta_large(
921 buf: &[u8],
922 n: usize,
923 ) -> Result<impl Iterator<Item = Delta> + '_, &'static str> {
924 if buf.len() < n * 2 {
925 return Err("Not enough buf for large deltas");
926 }
927 Ok((0..(n * 2))
928 .step_by(2)
929 .map(|i| Delta::Large(i16::from_be_bytes([buf[i], buf[i + 1]]))))
930 }
931 
932 for c in &twcc.chunks {
933 match c {
934 PacketChunk::Run(PacketStatus::ReceivedSmallDelta, n) => {
935 let n = *n as usize;
936 twcc.delta.extend(read_delta_small(buf, n)?);
937 buf = &buf[n..];
938 }
939 PacketChunk::Run(PacketStatus::ReceivedLargeOrNegativeDelta, n) => {
940 let n = *n as usize;
941 twcc.delta.extend(read_delta_large(buf, n)?);
942 buf = &buf[n..];
943 }
944 PacketChunk::VectorSingle(v, _) => {
945 let n = v.count_ones() as usize;
946 twcc.delta.extend(read_delta_small(buf, n)?);
947 buf = &buf[n..];
948 }
949 PacketChunk::VectorDouble(v, _) => {
950 for n in (0..=12).step_by(2) {
951 let x = (*v >> (12 - n)) & 0b11;
952 match PacketStatus::from(x as u8) {
953 PacketStatus::ReceivedSmallDelta => {
954 twcc.delta.extend(read_delta_small(buf, 1)?);
955 buf = &buf[1..];
956 }
957 PacketStatus::ReceivedLargeOrNegativeDelta => {
958 twcc.delta.extend(read_delta_large(buf, 1)?);
959 buf = &buf[2..];
960 }
961 _ => {}
962 }
963 }
964 }
965 _ => {}
966 }
967 }
968 
969 Ok(twcc)
970 }
971}
972 
973impl<'a> TryFrom<&'a [u8]> for PacketChunk {
974 type Error = &'static str;
975 
976 fn try_from(buf: &'a [u8]) -> Result<Self, Self::Error> {
977 if buf.len() < 2 {
978 return Err("Less than 2 bytes for PacketChunk");
979 }
980 
981 let x = u16::from_be_bytes([buf[0], buf[1]]);
982 
983 let is_vec = (x & 0b1000_0000_0000_0000) > 0;
984 
985 let p = if is_vec {
986 let is_double = (x & 0b0100_0000_0000_0000) > 0;
987 let n = x & 0b0011_1111_1111_1111;
988 if is_double {
989 PacketChunk::VectorDouble(n, 7)
990 } else {
991 PacketChunk::VectorSingle(n, 14)
992 }
993 } else {
994 let s: PacketStatus = ((x >> 13) as u8).into();
995 let n = x & 0b0001_1111_1111_1111;
996 PacketChunk::Run(s, n)
997 };
998 
999 Ok(p)
1000 }
1001}
1002 
1003#[derive(Debug)]
1004pub struct TwccSendRegister {
1005 /// How many send records to keep.
1006 keep: usize,
1007 
1008 /// Circular buffer of send records.
1009 queue: VecDeque<TwccSendRecord>,
1010 
1011 /// 0 offset for remote time in Twcc structs.
1012 time_zero: Option<Instant>,
1013 
1014 /// Counter of invocations of apply_report. Used to identify
1015 /// which TwccSendRecord resulted from each invocation.
1016 apply_report_counter: u64,
1017 
1018 /// Last registered Twcc number.
1019 last_registered: TwccSeq,
1020}
1021 
1022impl<'a> IntoIterator for &'a TwccSendRegister {
1023 type Item = &'a TwccSendRecord;
1024 type IntoIter = vec_deque::Iter<'a, TwccSendRecord>;
1025 
1026 fn into_iter(self) -> Self::IntoIter {
1027 self.queue.iter()
1028 }
1029}
1030 
1031/// Packet identification for TWCC tracking.
1032///
1033/// Groups together the TWCC sequence number and optional probe cluster context.
1034#[derive(Debug, Clone, Copy)]
1035pub struct TwccPacketId {
1036 /// TWCC sequence number
1037 seq: TwccSeq,
1038 /// Probe cluster this packet belongs to, if any
1039 cluster: Option<TwccClusterId>,
1040}
1041 
1042impl TwccPacketId {
1043 /// Create a packet ID for a regular media packet (no probe cluster).
1044 pub fn new(seq: impl Into<TwccSeq>) -> Self {
1045 Self {
1046 seq: seq.into(),
1047 cluster: None,
1048 }
1049 }
1050 
1051 /// Create a packet ID for a probe packet.
1052 pub fn with_cluster(seq: impl Into<TwccSeq>, cluster: impl Into<TwccClusterId>) -> Self {
1053 Self {
1054 seq: seq.into(),
1055 cluster: Some(cluster.into()),
1056 }
1057 }
1058 
1059 /// Get the TWCC sequence number.
1060 pub fn seq(&self) -> TwccSeq {
1061 self.seq
1062 }
1063 
1064 /// Get the probe cluster, if any.
1065 pub fn cluster(&self) -> Option<TwccClusterId> {
1066 self.cluster
1067 }
1068}
1069 
1070/// Record for a send entry in twcc.
1071#[derive(Debug)]
1072pub struct TwccSendRecord {
1073 /// Packet identification (sequence + optional probe cluster)
1074 packet_id: TwccPacketId,
1075 
1076 /// The (local) time we sent the packet represented by seq.
1077 local_send_time: Instant,
1078 
1079 /// Size in bytes of the payload sent.
1080 size: u16,
1081 
1082 recv_report: Option<TwccRecvReport>,
1083}
1084 
1085impl TwccSendRecord {
1086 /// The twcc sequence number of the packet we sent.
1087 pub fn seq(&self) -> TwccSeq {
1088 self.packet_id.seq()
1089 }
1090 
1091 /// The probe cluster this packet belongs to, if it's a probe packet.
1092 pub fn cluster(&self) -> Option<TwccClusterId> {
1093 self.packet_id.cluster()
1094 }
1095 
1096 /// The time we sent the packet.
1097 pub fn local_send_time(&self) -> Instant {
1098 self.local_send_time
1099 }
1100 
1101 /// The time we received this TWCC record. [`None`] if no feedback has been received yet.
1102 pub fn local_recv_time(&self) -> Option<Instant> {
1103 self.recv_report.as_ref().map(|r| r.local_recv_time)
1104 }
1105 
1106 pub fn size(&self) -> usize {
1107 self.size as usize
1108 }
1109 
1110 /// The time indicated by the remote side for when they received the packet.
1111 pub fn remote_recv_time(&self) -> Option<Instant> {
1112 self.recv_report.as_ref().and_then(|r| r.remote_recv_time)
1113 }
1114 
1115 /// The rtt time between sending the packet and receiving the twcc report response.
1116 pub fn rtt(&self) -> Option<Duration> {
1117 let recv_report = self.recv_report.as_ref()?;
1118 Some(recv_report.local_recv_time - self.local_send_time)
1119 }
1120}
1121 
1122#[cfg(test)]
1123impl TwccSendRecord {
1124 /// Test-only constructor to build TWCC send records with arbitrary receive status.
1125 ///
1126 /// This is used by unit tests for BWE/probing, allowing them to model received vs lost packets
1127 /// without constructing full RTCP TWCC reports.
1128 pub(crate) fn test_new(
1129 packet_id: TwccPacketId,
1130 local_send_time: Instant,
1131 size: usize,
1132 local_recv_time: Instant,
1133 remote_recv_time: Option<Instant>,
1134 ) -> Self {
1135 Self {
1136 packet_id,
1137 local_send_time,
1138 size: size as u16,
1139 recv_report: Some(TwccRecvReport {
1140 local_recv_time,
1141 remote_recv_time,
1142 apply_report_counter: 0,
1143 }),
1144 }
1145 }
1146}
1147 
1148#[derive(Debug, Copy, Clone)]
1149pub struct TwccRecvReport {
1150 /// The (local) time we received confirmation the other side received the seq.
1151 local_recv_time: Instant,
1152 
1153 /// The remote time the other side received the seq.
1154 remote_recv_time: Option<Instant>,
1155 
1156 /// The invocation count of apply_report(). Used for filtering.
1157 apply_report_counter: u64,
1158}
1159 
1160impl TwccSendRegister {
1161 pub fn new(keep: usize) -> Self {
1162 TwccSendRegister {
1163 keep,
1164 queue: VecDeque::new(),
1165 time_zero: None,
1166 apply_report_counter: 0,
1167 last_registered: 0.into(),
1168 }
1169 }
1170 
1171 pub fn register_seq(&mut self, packet_id: TwccPacketId, now: Instant, size: usize) {
1172 self.last_registered = packet_id.seq();
1173 self.queue.push_back(TwccSendRecord {
1174 packet_id,
1175 local_send_time: now,
1176 // In practice the max sizes is constrained by the MTU and will max out around 1200
1177 // bytes, hence this cast is fine.
1178 size: size as u16,
1179 // The recv report, derived from TWCC feedback later.
1180 recv_report: None,
1181 });
1182 while self.queue.len() > self.keep {
1183 self.queue.pop_front();
1184 }
1185 }
1186 
1187 /// Apply a TWCC RTCP report.
1188 ///
1189 /// Returns iterator over [`TwccSendRecord`]s included in the given [`Twcc`]
1190 /// except for ones that was already acked and returned before.
1191 pub fn apply_report(
1192 &mut self,
1193 twcc: Twcc,
1194 now: Instant,
1195 ) -> Option<impl Iterator<Item = &TwccSendRecord>> {
1196 if self.time_zero.is_none() {
1197 self.time_zero = Some(now);
1198 }
1199 
1200 self.apply_report_counter += 1;
1201 let apply_report_counter = self.apply_report_counter;
1202 
1203 let time_zero = self.time_zero.unwrap();
1204 let head_seq = self.queue.front().map(|r| r.seq())?;
1205 
1206 let mut iter = twcc
1207 .into_iter(time_zero, self.last_registered)
1208 .skip_while(|(seq, _, _)| seq < &head_seq);
1209 let (first_seq_no, _, first_instant) = iter.next()?;
1210 
1211 let mut iter2 = self
1212 .queue
1213 .iter_mut()
1214 .skip_while(|r| *r.seq() < *first_seq_no);
1215 let first_record = iter2.next()?;
1216 
1217 fn update(
1218 now: Instant,
1219 r: &mut TwccSendRecord,
1220 seq: TwccSeq,
1221 remote_recv_time: Option<Instant>,
1222 apply_report_counter: u64,
1223 ) -> bool {
1224 if r.seq() != seq {
1225 return false;
1226 }
1227 
1228 let apply_report_counter = if let Some(rr) = r.recv_report {
1229 // This packed was already acked and handled before so carry
1230 // over previous apply_report_counter, so it won't be included
1231 // in the current apply_report() call result.
1232 rr.remote_recv_time
1233 .map(|_| rr.apply_report_counter)
1234 .unwrap_or_else(|| apply_report_counter)
1235 } else {
1236 apply_report_counter
1237 };
1238 
1239 // Carry over remote recv time if this packet was acked before.
1240 let remote_recv_time = r.remote_recv_time().or(remote_recv_time);
1241 let recv_report = TwccRecvReport {
1242 local_recv_time: now,
1243 remote_recv_time,
1244 apply_report_counter,
1245 };
1246 r.recv_report = Some(recv_report);
1247 
1248 true
1249 }
1250 
1251 if first_record.seq() != first_seq_no {
1252 // Old report for which we no longer have any send records.
1253 return None;
1254 }
1255 
1256 let mut problematic_seq = None;
1257 
1258 if !update(
1259 now,
1260 first_record,
1261 first_seq_no,
1262 first_instant,
1263 apply_report_counter,
1264 ) {
1265 problematic_seq = Some((first_record.seq(), first_seq_no));
1266 }
1267 
1268 let mut last_seq_no = first_seq_no;
1269 for ((seq, _, instant), record) in iter.zip(iter2) {
1270 if problematic_seq.is_some() {
1271 break;
1272 }
1273 
1274 if !update(now, record, seq, instant, apply_report_counter) {
1275 problematic_seq = Some((record.seq(), seq));
1276 }
1277 last_seq_no = seq;
1278 }
1279 
1280 if let Some((record_seq, report_seq)) = problematic_seq {
1281 let queue_tail: Vec<_> = self.queue.iter().rev().take(100).collect();
1282 panic!(
1283 "Unexpected TWCC sequence when applying TWCC report. \
1284 Send Record Seq({record_seq}) does not match Report Seq({report_seq}).\
1285 \nLast 100 entires in queue: {queue_tail:?}."
1286 );
1287 }
1288 
1289 let first_index = self
1290 .queue
1291 .binary_search_by_key(&first_seq_no, |r| r.seq())
1292 .expect("first_seq_no to be registered");
1293 
1294 let range = first_seq_no..=last_seq_no;
1295 
1296 Some(
1297 TwccSendRecordsIter {
1298 range: range.clone(),
1299 index: first_index,
1300 current: first_seq_no,
1301 queue: &self.queue,
1302 }
1303 // We only want the records that were registered in this invocation of
1304 // apply_report_counter(). This is to not double count in the BWE,
1305 // which is the consumer of this returned iterator.
1306 .filter(move |s| {
1307 s.recv_report
1308 .map(|r| r.apply_report_counter == apply_report_counter)
1309 .unwrap_or_default()
1310 }),
1311 )
1312 }
1313 
1314 /// Calculate the egress loss for given time window.
1315 ///
1316 /// **Note:** The register only keeps a limited number of records and using `duration` values
1317 /// larger than ~1-2 seconds is liable to be inaccurate since some packets sent might have already
1318 /// been evicted from the register.
1319 pub fn loss(&self, duration: Duration, now: Instant) -> Option<f32> {
1320 // Consider only packets in the span specified by the caller
1321 let lower_bound = now - duration;
1322 
1323 let packets = self
1324 .queue
1325 .iter()
1326 .rev()
1327 // If there's ingress loss but no egress loss, there's a chance the TWCC reports
1328 // themselves are lost. In this case considering packets that haven't been reported as
1329 // lost will incorrectly conclude that there is in fact egress loss.
1330 .filter(|s| s.recv_report.is_some())
1331 .take_while(|s| s.local_send_time >= lower_bound);
1332 
1333 let (total, lost) = packets.fold((0, 0), |(total, lost), s| {
1334 let was_lost = s
1335 .recv_report
1336 .as_ref()
1337 .map(|rr| rr.remote_recv_time.is_none())
1338 .unwrap_or(true);
1339 
1340 (total + 1, lost + u64::from(was_lost))
1341 });
1342 
1343 if total == 0 {
1344 return None;
1345 }
1346 
1347 Some((lost as f32) / (total as f32))
1348 }
1349 
1350 /// Calculate the RTT for the most recently reported packet.
1351 pub fn rtt(&self) -> Option<Duration> {
1352 self.queue.iter().rev().find_map(|s| s.rtt())
1353 }
1354}
1355 
1356#[derive()]
1357struct TwccSendRecordsIter<'a> {
1358 range: RangeInclusive<TwccSeq>,
1359 current: TwccSeq,
1360 index: usize,
1361 queue: &'a VecDeque<TwccSendRecord>,
1362}
1363 
1364impl<'a> Iterator for TwccSendRecordsIter<'a> {
1365 type Item = &'a TwccSendRecord;
1366 
1367 fn next(&mut self) -> Option<Self::Item> {
1368 if self.current > *self.range.end() || self.current < *self.range.start() {
1369 return None;
1370 }
1371 
1372 let item = &self.queue[self.index];
1373 assert!(self.current == item.seq());
1374 self.current.inc();
1375 self.index += 1;
1376 
1377 Some(item)
1378 }
1379}
1380 
1381// Below is a clarification of the RFC draft from an email exchange with Erik Sprรฅng (one of the authors).
1382//
1383// > I'm trying to implement the draft spec
1384// > https://datatracker.ietf.org/doc/html/draft-holmer-rmcat-transport-wide-cc-extensions-01
1385// > I found a number of errors/inconsistencies in the RFC, and wonder who I should address this to.
1386// > I think the RFC could benefit from another revision.
1387// >
1388// > There are three problems listed below.
1389// >
1390// > 1. There's a contradiction between section 3.1.1 and the example in 3.1.3. First the RFC tells
1391// > me 11 is Reserved, later it shows an example using 11 saying it is a run of packets received w/o
1392// > recv delta. Which one is right?
1393// >
1394// > Section 3.1.1
1395// > ...
1396// > The status of a packet is described using a 2-bit symbol:
1397// > ...
1398// > 11 [Reserved]
1399// >
1400// > Section 3.1.3
1401// > ...
1402// > Example 2:
1403// >
1404// > 0 1
1405// > 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5
1406// > +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
1407// > |0|1 1|0 0 0 0 0 0 0 0 1 1 0 0 0|
1408// > +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
1409// >
1410// >
1411// > This is a run of the "packet received, w/o recv delta" status of
1412// > length 24.
1413//
1414// I believe this example is in error. Packets without receive deltas was a proposal for the v1
1415// protocol but was dropped iirc. Note that there is a newer header extension available, which
1416// when negotiated allows send-side control over when feedback is generated - and that provides
1417// the option to omit all receive deltas from the feedback.
1418//
1419// > 2. In section 3.1.4 when using a 1-bit vector to indicate packet received or not received,
1420// > there's a contradiction between the the definition and the example. The definition says
1421// > "packet received" (0) and "packet not received" (1), while the example is the opposite way
1422// > around: 0 is packet not received. Which way around is it?
1423// >
1424// > symbol size (S): 1 bit A zero means this vector contains only
1425// > "packet received" (0) and "packet not received" (1)
1426// > symbols. This means we can compress each symbol to just
1427// > one bit, 14 in total. A one means this vector contains
1428// > the normal 2-bit symbols, 7 in total.
1429// > ...
1430// > Example 1:
1431// >
1432// > 0 1
1433// > 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5
1434// > +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
1435// > |1|0|0 1 1 1 1 1 0 0 0 1 1 1 0 0|
1436// > +-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
1437// >
1438// > This chunk contains, in order:
1439// >
1440// > 1x "packet not received"
1441// >
1442// > 5x "packet received"
1443//
1444// I believe the definition is wrong in this case. The intent is to just truncate the 2-bit values:
1445//
1446// 3.1.1. Packet Status Symbols
1447//
1448// The status of a packet is described using a 2-bit symbol:
1449//
1450// 00 Packet not received
1451//
1452// 01 Packet received, small delta
1453//
1454// So (0) for not received and (1) for received, small delta.
1455// This also matches what the libwebrtc source code does.
1456//
1457// > 3. In section 3.1.4 when using a 1-bit vector, the RFC doesn't say what a "packet received" in that
1458// > vector should be accompanied by in receive delta size. Is it an 8 bit delta or 16 bit
1459// > delta per "packet received"?
1460//
1461// Same as the question above, this is a truncation to (0) for not received and (1) for
1462// received, small delta.
1463 
1464impl fmt::Debug for Twcc {
1465 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1466 f.debug_struct("Twcc")
1467 .field("sender_ssrc", &self.sender_ssrc)
1468 .field("ssrc", &self.ssrc)
1469 .field("base_seq", &self.base_seq)
1470 .field("status_count", &self.status_count)
1471 .field("reference_time", &self.reference_time)
1472 .field("feedback_count", &self.feedback_count)
1473 .field("chunks", &self.chunks)
1474 .field("delta", &self.delta.len())
1475 .finish()
1476 }
1477}
1478 
1479#[allow(clippy::assign_op_pattern)]
1480#[cfg(test)]
1481mod test {
1482 use std::time::Duration;
1483 
1484 use super::*;
1485 
1486 use Delta::*;
1487 use PacketChunk::*;
1488 use PacketStatus::*;
1489 
1490 #[test]
1491 fn register_write_parse_small_delta() {
1492 let mut reg = TwccRecvRegister::new(100);
1493 
1494 let now = Instant::now();
1495 
1496 reg.update_seq(10.into(), now + Duration::from_millis(0));
1497 reg.update_seq(11.into(), now + Duration::from_millis(12));
1498 reg.update_seq(12.into(), now + Duration::from_millis(23));
1499 reg.update_seq(13.into(), now + Duration::from_millis(43));
1500 
1501 let report = reg.build_report(1000).unwrap();
1502 let mut buf = vec![0_u8; 1500];
1503 let n = report.write_to(&mut buf[..]);
1504 buf.truncate(n);
1505 
1506 let header: RtcpHeader = (&buf[..]).try_into().unwrap();
1507 let parsed: Twcc = (&buf[4..]).try_into().unwrap();
1508 
1509 assert_eq!(header, report.header());
1510 assert_eq!(parsed, report);
1511 }
1512 
1513 #[test]
1514 fn register_write_parse_small_delta_missing() {
1515 let mut reg = TwccRecvRegister::new(100);
1516 
1517 let now = Instant::now();
1518 
1519 reg.update_seq(10.into(), now + Duration::from_millis(0));
1520 reg.update_seq(11.into(), now + Duration::from_millis(12));
1521 reg.update_seq(12.into(), now + Duration::from_millis(23));
1522 // 13 is not there
1523 reg.update_seq(14.into(), now + Duration::from_millis(43));
1524 
1525 let report = reg.build_report(1000).unwrap();
1526 let mut buf = vec![0_u8; 1500];
1527 let n = report.write_to(&mut buf[..]);
1528 buf.truncate(n);
1529 
1530 let header: RtcpHeader = (&buf[..]).try_into().unwrap();
1531 let parsed: Twcc = (&buf[4..]).try_into().unwrap();
1532 
1533 assert_eq!(header, report.header());
1534 assert_eq!(parsed, report);
1535 }
1536 
1537 #[test]
1538 fn register_write_parse_large_delta() {
1539 let mut reg = TwccRecvRegister::new(100);
1540 
1541 let now = Instant::now();
1542 
1543 reg.update_seq(10.into(), now + Duration::from_millis(0));
1544 reg.update_seq(11.into(), now + Duration::from_millis(70));
1545 reg.update_seq(12.into(), now + Duration::from_millis(140));
1546 reg.update_seq(13.into(), now + Duration::from_millis(210));
1547 
1548 let report = reg.build_report(1000).unwrap();
1549 let mut buf = vec![0_u8; 1500];
1550 let n = report.write_to(&mut buf[..]);
1551 buf.truncate(n);
1552 
1553 let header: RtcpHeader = (&buf[..]).try_into().unwrap();
1554 let parsed: Twcc = (&buf[4..]).try_into().unwrap();
1555 
1556 assert_eq!(header, report.header());
1557 assert_eq!(parsed, report);
1558 }
1559 
1560 #[test]
1561 fn register_write_parse_mixed_delta() {
1562 let mut reg = TwccRecvRegister::new(100);
1563 
1564 let now = Instant::now();
1565 
1566 reg.update_seq(10.into(), now + Duration::from_millis(0));
1567 reg.update_seq(11.into(), now + Duration::from_millis(12));
1568 reg.update_seq(12.into(), now + Duration::from_millis(140));
1569 reg.update_seq(13.into(), now + Duration::from_millis(152));
1570 
1571 let report = reg.build_report(1000).unwrap();
1572 let mut buf = vec![0_u8; 1500];
1573 let n = report.write_to(&mut buf[..]);
1574 buf.truncate(n);
1575 
1576 let header: RtcpHeader = (&buf[..]).try_into().unwrap();
1577 let parsed: Twcc = (&buf[4..]).try_into().unwrap();
1578 
1579 assert_eq!(header, report.header());
1580 assert_eq!(parsed, report);
1581 }
1582 
1583 #[test]
1584 fn too_big_time_gap_requires_two_reports() {
1585 let mut reg = TwccRecvRegister::new(100);
1586 
1587 let now = Instant::now();
1588 
1589 reg.update_seq(10.into(), now + Duration::from_millis(0));
1590 reg.update_seq(11.into(), now + Duration::from_millis(12));
1591 reg.update_seq(12.into(), now + Duration::from_millis(9000));
1592 
1593 let _ = reg.build_report(1000).unwrap();
1594 let report2 = reg.build_report(1000).unwrap();
1595 
1596 // 9000 milliseconds is not possible to set as exact reference time which
1597 // is in multiples of 64ms. 9000/64 = 140.625.
1598 assert_eq!(report2.reference_time, 140);
1599 
1600 // 140 * 64 = 8960
1601 // So the first offset must be 40ms, i.e. 40_000us / 250us = 160
1602 assert_eq!(report2.delta[0], Small(160));
1603 }
1604 
1605 #[test]
1606 fn report_padded_to_even_word() {
1607 let mut reg = TwccRecvRegister::new(100);
1608 
1609 let now = Instant::now();
1610 
1611 reg.update_seq(10.into(), now + Duration::from_millis(0));
1612 
1613 let report = reg.build_report(1000).unwrap();
1614 let mut buf = vec![0_u8; 1500];
1615 let n = report.write_to(&mut buf[..]);
1616 
1617 assert!(n % 4 == 0);
1618 }
1619 
1620 #[test]
1621 fn report_truncated_to_max_byte_size() {
1622 let mut reg = TwccRecvRegister::new(100);
1623 
1624 let now = Instant::now();
1625 
1626 reg.update_seq(10.into(), now + Duration::from_millis(0));
1627 reg.update_seq(11.into(), now + Duration::from_millis(12));
1628 reg.update_seq(12.into(), now + Duration::from_millis(140));
1629 reg.update_seq(13.into(), now + Duration::from_millis(152));
1630 
1631 let report = reg.build_report(28).unwrap();
1632 
1633 assert_eq!(report.status_count, 2);
1634 assert_eq!(report.chunks, vec![Run(ReceivedSmallDelta, 2)]);
1635 assert_eq!(report.delta, vec![Small(0), Small(48)]);
1636 
1637 let report = reg.build_report(28).unwrap();
1638 
1639 assert_eq!(report.status_count, 2);
1640 assert_eq!(report.chunks, vec![Run(ReceivedSmallDelta, 2)]);
1641 assert_eq!(report.delta, vec![Small(48), Small(48)]);
1642 }
1643 
1644 #[test]
1645 fn truncated_counts_gaps_correctly() {
1646 let mut reg = TwccRecvRegister::new(100);
1647 
1648 let now = Instant::now();
1649 
1650 reg.update_seq(10.into(), now + Duration::from_millis(0));
1651 // gap
1652 reg.update_seq(13.into(), now + Duration::from_millis(12));
1653 reg.update_seq(14.into(), now + Duration::from_millis(140));
1654 reg.update_seq(15.into(), now + Duration::from_millis(152));
1655 
1656 let report = reg.build_report(32).unwrap();
1657 
1658 assert_eq!(report.status_count, 4);
1659 assert_eq!(
1660 report.chunks,
1661 vec![
1662 Run(ReceivedSmallDelta, 1),
1663 Run(NotReceived, 2),
1664 Run(ReceivedSmallDelta, 1)
1665 ]
1666 );
1667 assert_eq!(report.delta, vec![Small(0), Small(48)]);
1668 }
1669 
1670 #[test]
1671 fn run_max_is_8192() {
1672 let mut reg = TwccRecvRegister::new(100);
1673 
1674 let now = Instant::now();
1675 
1676 reg.update_seq(0.into(), now + Duration::from_millis(0));
1677 reg.update_seq(8194.into(), now + Duration::from_millis(10));
1678 
1679 let report = reg.build_report(1000).unwrap();
1680 
1681 assert_eq!(report.status_count, 8195);
1682 assert_eq!(
1683 report.chunks,
1684 vec![
1685 VectorSingle(8192, 14),
1686 Run(NotReceived, 8180),
1687 Run(ReceivedSmallDelta, 1)
1688 ]
1689 );
1690 }
1691 
1692 #[test]
1693 fn single_followed_by_missing() {
1694 let mut reg = TwccRecvRegister::new(100);
1695 
1696 let now = Instant::now();
1697 
1698 reg.update_seq(10.into(), now + Duration::from_millis(0));
1699 reg.update_seq(12.into(), now + Duration::from_millis(10));
1700 reg.update_seq(100.into(), now + Duration::from_millis(20));
1701 
1702 let report = reg.build_report(2016).unwrap();
1703 
1704 assert_eq!(report.status_count, 91);
1705 assert_eq!(
1706 report.chunks,
1707 vec![
1708 VectorSingle(10240, 14),
1709 Run(NotReceived, 76),
1710 Run(ReceivedSmallDelta, 1)
1711 ]
1712 );
1713 assert_eq!(report.delta, vec![Small(0), Small(40), Small(40)]);
1714 }
1715 
1716 #[test]
1717 fn time_jump_small_back_for_second_report() {
1718 let mut reg = TwccRecvRegister::new(100);
1719 
1720 let now = Instant::now();
1721 
1722 reg.update_seq(10.into(), now + Duration::from_millis(8000));
1723 let _ = reg.build_report(2016).unwrap();
1724 
1725 reg.update_seq(9.into(), now + Duration::from_millis(0));
1726 let report = reg.build_report(2016).unwrap();
1727 
1728 assert_eq!(report.status_count, 2);
1729 assert_eq!(report.chunks, vec![Run(ReceivedLargeOrNegativeDelta, 2)]);
1730 assert_eq!(report.delta, vec![Large(-32000), Large(32000)]);
1731 }
1732 
1733 #[test]
1734 fn time_jump_large_back_for_second_report() {
1735 let mut reg = TwccRecvRegister::new(100);
1736 
1737 let now = Instant::now();
1738 
1739 reg.update_seq(10.into(), now + Duration::from_millis(9000));
1740 let _ = reg.build_report(2016).unwrap();
1741 
1742 reg.update_seq(9.into(), now + Duration::from_millis(0));
1743 assert!(reg.build_report(2016).is_none());
1744 
1745 assert_eq!(reg.queue.len(), 1);
1746 }
1747 
1748 #[test]
1749 fn empty_twcc() {
1750 let twcc = Twcc {
1751 sender_ssrc: 0.into(),
1752 ssrc: 0.into(),
1753 base_seq: 0,
1754 status_count: 0,
1755 reference_time: 0,
1756 feedback_count: 0,
1757 chunks: VecDeque::new(),
1758 delta: VecDeque::new(),
1759 };
1760 
1761 let mut buf = vec![0_u8; 1500];
1762 let n = twcc.write_to(&mut buf[..]);
1763 buf.truncate(n);
1764 
1765 let header: RtcpHeader = (&buf[..]).try_into().unwrap();
1766 let parsed: Twcc = (&buf[4..]).try_into().unwrap();
1767 
1768 assert_eq!(header, twcc.header());
1769 assert_eq!(parsed, twcc);
1770 }
1771 
1772 #[test]
1773 fn negative_deltas() {
1774 let mut reg = TwccRecvRegister::new(100);
1775 
1776 let now = Instant::now();
1777 
1778 reg.update_seq(10.into(), now + Duration::from_millis(12));
1779 reg.update_seq(11.into(), now + Duration::from_millis(0));
1780 reg.update_seq(12.into(), now + Duration::from_millis(23));
1781 
1782 let report = reg.build_report(1000).unwrap();
1783 
1784 assert_eq!(report.status_count, 3);
1785 assert_eq!(report.base_seq, 10);
1786 assert_eq!(report.reference_time, 0);
1787 assert_eq!(report.chunks, vec![VectorDouble(6400, 7)]);
1788 assert_eq!(report.delta, vec![Small(0), Large(-48), Small(92)]);
1789 
1790 let base = reg.time_start.unwrap();
1791 
1792 let mut iter = report.into_iter(base, 10.into());
1793 assert_eq!(
1794 iter.next(),
1795 Some((
1796 10.into(),
1797 PacketStatus::ReceivedSmallDelta,
1798 Some(base + Duration::from_millis(0))
1799 ))
1800 );
1801 assert_eq!(
1802 iter.next(),
1803 Some((
1804 11.into(),
1805 PacketStatus::ReceivedLargeOrNegativeDelta,
1806 Some(base.checked_sub(Duration::from_millis(12)).unwrap())
1807 ))
1808 );
1809 assert_eq!(
1810 iter.next(),
1811 Some((
1812 12.into(),
1813 PacketStatus::ReceivedSmallDelta,
1814 Some(base + Duration::from_millis(11))
1815 ))
1816 );
1817 }
1818 
1819 #[test]
1820 fn twcc_fuzz_fail() {
1821 let mut reg = TwccRecvRegister::new(100);
1822 
1823 let now = Instant::now();
1824 
1825 // [Register(, ), Register(, ), Register(, ), BuildReport(43)]
1826 
1827 reg.update_seq(4542.into(), now + Duration::from_millis(2373281424));
1828 reg.update_seq(15918.into(), now + Duration::from_millis(2373862820));
1829 reg.update_seq(8405.into(), now + Duration::from_millis(2379074367));
1830 
1831 let report = reg.build_report(43).unwrap();
1832 
1833 let mut buf = vec![0_u8; 1500];
1834 let n = report.write_to(&mut buf[..]);
1835 buf.truncate(n);
1836 
1837 let header: RtcpHeader = match (&buf[..]).try_into() {
1838 Ok(v) => v,
1839 Err(_) => return,
1840 };
1841 let parsed: Twcc = match (&buf[4..]).try_into() {
1842 Ok(v) => v,
1843 Err(_) => return,
1844 };
1845 
1846 assert_eq!(header, report.header());
1847 assert_eq!(parsed, report);
1848 }
1849 
1850 #[test]
1851 fn twcc_large_time_delta_edges() {
1852 let mut reg = TwccRecvRegister::new(100);
1853 
1854 // Stretch the boundaries of the i16 type used for deltas
1855 let now = Instant::now();
1856 reg.update_seq(0.into(), now + Duration::from_micros(8_192_000));
1857 reg.update_seq(1.into(), now + Duration::from_micros(0));
1858 reg.update_seq(2.into(), now + Duration::from_micros(8_191_750));
1859 
1860 let report = reg.build_report(1000).unwrap();
1861 
1862 assert_eq!(report.status_count, 3);
1863 assert_eq!(report.delta, vec![Small(0), Large(-32768), Large(32767)]);
1864 }
1865 
1866 #[test]
1867 fn twcc_crazy_negative_time_delta() {
1868 let mut reg = TwccRecvRegister::new(100);
1869 
1870 // Deltas so big they wrap around the bounds of an i32 to become small again.
1871 // These constants are chosen carefully to look normal when wrapped
1872 let now = Instant::now();
1873 reg.update_seq(0.into(), now + Duration::from_micros(4_294_967_547)); // Wraps to -251
1874 reg.update_seq(1.into(), now + Duration::from_micros(0));
1875 
1876 // The bogus value should be ignored
1877 let report = reg.build_report(1000).unwrap();
1878 assert_eq!(report.status_count, 1);
1879 assert_eq!(report.delta, vec![Small(0)]);
1880 }
1881 
1882 #[test]
1883 fn twcc_crazy_positive_time_delta() {
1884 let mut reg = TwccRecvRegister::new(100);
1885 
1886 // Deltas so big they wrap around the bounds of an i32 to become small again.
1887 // These constants are chosen carefully to look normal when wrapped
1888 let now = Instant::now();
1889 reg.update_seq(0.into(), now + Duration::from_micros(0));
1890 reg.update_seq(1.into(), now + Duration::from_micros(4_294_967_547)); // Wraps to 251
1891 
1892 // The bogus value should be ignored
1893 let report = reg.build_report(1000).unwrap();
1894 assert_eq!(report.status_count, 1);
1895 assert_eq!(report.delta, vec![Small(0)]);
1896 }
1897 
1898 #[test]
1899 fn test_send_register_apply_report_for_old_seq_numbers() {
1900 let mut reg = TwccSendRegister::new(25);
1901 let mut now = Instant::now();
1902 
1903 for i in 0..50 {
1904 reg.register_seq(TwccPacketId::new(i), now, 0);
1905 now = now + Duration::from_micros(15);
1906 }
1907 
1908 // At this point the front of the internal queue should be seq no 25.
1909 //
1910 // Set time zero base with empty packet
1911 reg.apply_report(
1912 Twcc {
1913 sender_ssrc: Ssrc::new(),
1914 ssrc: Ssrc::new(),
1915 base_seq: 0,
1916 status_count: 0,
1917 reference_time: 0,
1918 feedback_count: 0,
1919 chunks: [].into(),
1920 delta: [].into(),
1921 },
1922 now,
1923 );
1924 now = now + Duration::from_millis(35);
1925 
1926 let iter = reg.apply_report(
1927 Twcc {
1928 sender_ssrc: Ssrc::new(),
1929 ssrc: Ssrc::new(),
1930 base_seq: 20,
1931 status_count: 8,
1932 reference_time: 35,
1933 feedback_count: 0,
1934 chunks: [PacketChunk::Run(PacketStatus::ReceivedSmallDelta, 8)].into(),
1935 delta: [
1936 Delta::Small(10),
1937 Delta::Small(10),
1938 Delta::Small(10),
1939 Delta::Small(10),
1940 Delta::Small(10),
1941 Delta::Small(10),
1942 Delta::Small(10),
1943 Delta::Small(10),
1944 ]
1945 .into(),
1946 },
1947 now,
1948 );
1949 let iter = iter.unwrap();
1950 
1951 for record in iter {
1952 assert!(
1953 record.recv_report.is_some(),
1954 "Report should have recorded recv_report"
1955 );
1956 }
1957 }
1958 
1959 #[test]
1960 fn test_twcc_iter_correct_deltas() {
1961 let twcc = Twcc {
1962 sender_ssrc: 0.into(),
1963 ssrc: 0.into(),
1964 base_seq: 1,
1965 status_count: 12,
1966 reference_time: 1337,
1967 feedback_count: 0,
1968 chunks: [
1969 PacketChunk::Run(PacketStatus::ReceivedSmallDelta, 2),
1970 PacketChunk::VectorDouble(0b1101_0010_0101_0000, 7),
1971 PacketChunk::Run(PacketStatus::ReceivedLargeOrNegativeDelta, 3),
1972 ]
1973 .into(),
1974 delta: [
1975 // Run of length 2
1976 Delta::Small(10),
1977 Delta::Small(15),
1978 // Double status vector with 4 deltas
1979 Delta::Small(7),
1980 Delta::Large(280),
1981 Delta::Small(3),
1982 Delta::Small(13),
1983 // Run of length 3
1984 Delta::Large(-37),
1985 Delta::Large(32),
1986 Delta::Large(89),
1987 ]
1988 .into(),
1989 };
1990 
1991 let now = Instant::now();
1992 let base = now + Duration::from_millis(1337 * 64);
1993 let expected = vec![
1994 (
1995 1.into(),
1996 PacketStatus::ReceivedSmallDelta,
1997 Some(base + Duration::from_micros(10 * 250)),
1998 ),
1999 (
2000 2.into(),
2001 PacketStatus::ReceivedSmallDelta,
2002 Some(base + Duration::from_micros(25 * 250)),
2003 ),
2004 (
2005 3.into(),
2006 PacketStatus::ReceivedSmallDelta,
2007 Some(base + Duration::from_micros(32 * 250)),
2008 ),
2009 (4.into(), PacketStatus::NotReceived, None),
2010 (
2011 5.into(),
2012 PacketStatus::ReceivedLargeOrNegativeDelta,
2013 Some(base + Duration::from_micros(312 * 250)),
2014 ),
2015 (
2016 6.into(),
2017 PacketStatus::ReceivedSmallDelta,
2018 Some(base + Duration::from_micros(315 * 250)),
2019 ),
2020 (
2021 7.into(),
2022 PacketStatus::ReceivedSmallDelta,
2023 Some(base + Duration::from_micros(328 * 250)),
2024 ),
2025 (8.into(), PacketStatus::NotReceived, None),
2026 (9.into(), PacketStatus::NotReceived, None),
2027 (
2028 10.into(),
2029 PacketStatus::ReceivedLargeOrNegativeDelta,
2030 Some(base + Duration::from_micros(291 * 250)),
2031 ),
2032 (
2033 11.into(),
2034 PacketStatus::ReceivedLargeOrNegativeDelta,
2035 Some(base + Duration::from_micros(323 * 250)),
2036 ),
2037 (
2038 12.into(),
2039 PacketStatus::ReceivedLargeOrNegativeDelta,
2040 Some(base + Duration::from_micros(412 * 250)),
2041 ),
2042 ];
2043 
2044 let result: Vec<_> = twcc.into_iter(now, 1.into()).collect();
2045 
2046 assert_eq!(result, expected);
2047 }
2048 
2049 #[test]
2050 fn test_twcc_iter_limited_with_status_count() {
2051 let status_count = 3;
2052 
2053 // [(1, NotReceived), (2, ReceivedSmallDelta), (3, ReceivedSmallDelta)]
2054 let twcc_iter_count = Twcc {
2055 sender_ssrc: 1.into(),
2056 ssrc: 2.into(),
2057 base_seq: 1,
2058 status_count,
2059 reference_time: 406753,
2060 feedback_count: 1,
2061 chunks: VecDeque::from(vec![VectorDouble(0b00_00_01_01_00_00_00_00, 7)]),
2062 delta: VecDeque::from(vec![Small(236), Small(1)]),
2063 }
2064 .into_iter(Instant::now(), 1.into())
2065 .count();
2066 
2067 assert_eq!(twcc_iter_count, status_count as usize);
2068 }
2069 
2070 #[test]
2071 fn test_twcc_register_send_records() {
2072 let mut reg = TwccSendRegister::new(25);
2073 let mut now = Instant::now();
2074 for i in 0..25 {
2075 reg.register_seq(TwccPacketId::new(i), now, 0);
2076 now = now + Duration::from_micros(15);
2077 }
2078 
2079 let iter = reg
2080 .apply_report(
2081 Twcc {
2082 sender_ssrc: Ssrc::new(),
2083 ssrc: Ssrc::new(),
2084 base_seq: 0,
2085 status_count: 8,
2086 reference_time: 35,
2087 feedback_count: 0,
2088 chunks: [PacketChunk::Run(PacketStatus::ReceivedSmallDelta, 8)].into(),
2089 delta: [
2090 Delta::Small(10),
2091 Delta::Small(10),
2092 Delta::Small(10),
2093 Delta::Small(10),
2094 Delta::Small(10),
2095 Delta::Small(10),
2096 Delta::Small(10),
2097 Delta::Small(10),
2098 ]
2099 .into(),
2100 },
2101 now,
2102 )
2103 .expect("apply_report to return Some(_)");
2104 
2105 assert_eq!(
2106 iter.map(|r| *r.seq()).collect::<Vec<_>>(),
2107 vec![0, 1, 2, 3, 4, 5, 6, 7]
2108 );
2109 }
2110 
2111 #[test]
2112 fn test_twcc_send_register_loss() {
2113 let mut reg = TwccSendRegister::new(25);
2114 let mut now = Instant::now();
2115 for i in 0..9 {
2116 reg.register_seq(TwccPacketId::new(i), now, 0);
2117 now = now + Duration::from_millis(15);
2118 }
2119 
2120 now = now + Duration::from_millis(5);
2121 #[allow(unused_must_use)]
2122 reg.apply_report(
2123 Twcc {
2124 sender_ssrc: Ssrc::new(),
2125 ssrc: Ssrc::new(),
2126 base_seq: 0,
2127 status_count: 9,
2128 reference_time: 35,
2129 feedback_count: 0,
2130 chunks: [
2131 PacketChunk::VectorDouble(0b11_01_01_01_00_01_00_01, 7),
2132 PacketChunk::Run(PacketStatus::ReceivedSmallDelta, 2),
2133 ]
2134 .into(),
2135 delta: [
2136 Delta::Small(10),
2137 Delta::Small(10),
2138 Delta::Small(10),
2139 Delta::Small(10),
2140 Delta::Small(10),
2141 Delta::Small(10),
2142 Delta::Small(10),
2143 ]
2144 .into(),
2145 },
2146 now,
2147 )
2148 .expect("apply_report to return Some(_)");
2149 
2150 now = now + Duration::from_millis(20);
2151 let loss = reg
2152 .loss(Duration::from_millis(150), now)
2153 .expect("Should be able to calcualte loss");
2154 
2155 let pct = (loss * 100.0).floor() as u32;
2156 
2157 assert_eq!(
2158 pct, 25,
2159 "The loss percentage should be 25 as 2 out of 8 packets are lost"
2160 );
2161 }
2162 
2163 #[test]
2164 fn test_twcc_recv_register_loss() {
2165 let mut reg = TwccRecvRegister::new(25);
2166 let mut now = Instant::now();
2167 
2168 for i in 0..10 {
2169 if i == 3 || i == 7 {
2170 // simulate loss
2171 continue;
2172 }
2173 reg.update_seq(i.into(), now);
2174 now = now + Duration::from_millis(50);
2175 }
2176 
2177 assert_eq!(reg.loss(), Some(2.0 / 10.0));
2178 
2179 for i in 10..20 {
2180 if i == 11 || i == 13 || i == 15 || i == 17 {
2181 // simulate loss
2182 continue;
2183 }
2184 reg.update_seq(i.into(), now);
2185 now = now + Duration::from_millis(50);
2186 }
2187 
2188 assert_eq!(reg.loss(), Some(4.0 / 10.0));
2189 }
2190 
2191 #[test]
2192 fn no_acked_duplicates_when_reordered() {
2193 let now = Instant::now();
2194 let mut twcc_gen = TwccRecvRegister::new(1000);
2195 let mut twcc_handler = TwccSendRegister::new(1000);
2196 
2197 twcc_handler.register_seq(TwccPacketId::new(1), now + Duration::from_millis(1), 0);
2198 twcc_handler.register_seq(TwccPacketId::new(2), now + Duration::from_millis(2), 0);
2199 twcc_handler.register_seq(TwccPacketId::new(3), now + Duration::from_millis(3), 0);
2200 twcc_handler.register_seq(TwccPacketId::new(4), now + Duration::from_millis(4), 0);
2201 
2202 {
2203 let acked_packets = twcc_handler
2204 .apply_report(
2205 {
2206 // 3rd packet is delayed
2207 twcc_gen.update_seq(1.into(), now + Duration::from_millis(5));
2208 twcc_gen.update_seq(2.into(), now + Duration::from_millis(6));
2209 twcc_gen.update_seq(4.into(), now + Duration::from_millis(7));
2210 twcc_gen.build_report(10_000).unwrap()
2211 },
2212 now + Duration::from_millis(8),
2213 )
2214 .unwrap()
2215 .filter_map(|sr| sr.remote_recv_time().map(|_| sr.seq().as_u16()))
2216 .collect::<Vec<_>>();
2217 
2218 assert_eq!(acked_packets, [1, 2, 4]);
2219 }
2220 
2221 twcc_handler.register_seq(TwccPacketId::new(5), now + Duration::from_millis(9), 0);
2222 twcc_handler.register_seq(TwccPacketId::new(6), now + Duration::from_millis(10), 0);
2223 twcc_handler.register_seq(TwccPacketId::new(7), now + Duration::from_millis(11), 0);
2224 
2225 {
2226 let acked_packets = twcc_handler
2227 .apply_report(
2228 {
2229 // So the receipt order is 1, 2, 4, 3, 7
2230 twcc_gen.update_seq(3.into(), now + Duration::from_millis(12));
2231 twcc_gen.update_seq(7.into(), now + Duration::from_millis(13));
2232 twcc_gen.build_report(10_000).unwrap()
2233 },
2234 now + Duration::from_millis(14),
2235 )
2236 .unwrap()
2237 .filter_map(|sr| sr.remote_recv_time().map(|_| sr.seq().as_u16()))
2238 .collect::<Vec<_>>();
2239 
2240 // 3 was delayed before and is acked now
2241 // 4 is excluded since it was already returned from the previous call
2242 // [5, 6] are delayed/lost
2243 // 7 is acked in the last report
2244 assert_eq!(acked_packets, [3, 7]);
2245 }
2246 }
2247 
2248 #[test]
2249 fn test_twcc_recv_register_high_initial_sequence() {
2250 // Test TWCC sequence numbers starting at high values (>= 32768)
2251 let test_cases = vec![
2252 (32767, "just below half"),
2253 (32768, "exactly half"),
2254 (65434, "just below u16"),
2255 (65535, "maximum u16"),
2256 ];
2257 
2258 for (transport_cc, description) in test_cases {
2259 let mut twcc_rx_register = TwccRecvRegister::new(100);
2260 let now = Instant::now();
2261 
2262 // Empty register should return `None`
2263 assert_eq!(
2264 twcc_rx_register.max_seq(),
2265 None,
2266 "Empty register should return None for {}",
2267 description
2268 );
2269 
2270 // This is what `session.rs` does:
2271 let prev = twcc_rx_register.max_seq();
2272 let extended = extend_u16(prev.map(|s| *s), transport_cc);
2273 twcc_rx_register.update_seq(extended.into(), now);
2274 
2275 assert_eq!(
2276 twcc_rx_register.max_seq(),
2277 Some(extended.into()),
2278 "First packet for {}",
2279 description
2280 );
2281 
2282 // Add a few more packets, simulating production code path
2283 for i in 1..5 {
2284 let raw_seq = ((transport_cc as u32 + i) % 65536) as u16;
2285 let prev = twcc_rx_register.max_seq();
2286 let extended = extend_u16(prev.map(|s| *s), raw_seq);
2287 twcc_rx_register
2288 .update_seq(extended.into(), now + Duration::from_millis(i as u64 * 10));
2289 }
2290 
2291 // Verify wrap-around handling if we crossed the boundary
2292 if transport_cc >= 65533 {
2293 // Should have wrapped to extended values > 65536
2294 let max = twcc_rx_register.max_seq().unwrap();
2295 assert!(
2296 *max > 65536,
2297 "Should wrap correctly for {}: got {}",
2298 description,
2299 *max
2300 );
2301 }
2302 
2303 // build_report should NOT panic or allocate huge memory (leading to OOM)
2304 let report = twcc_rx_register.build_report(1000);
2305 assert!(report.is_some(), "Should build report for {}", description);
2306 
2307 let report = report.unwrap();
2308 assert!(
2309 report.status_count > 0,
2310 "Should have packets for {}",
2311 description
2312 );
2313 assert!(
2314 !report.chunks.is_empty(),
2315 "Should have chunks for {}",
2316 description
2317 );
2318 
2319 // Sanity check: should never have an absurdly large number of chunks
2320 assert!(
2321 report.chunks.len() < 10,
2322 "Should not create excessive chunks for {}",
2323 description
2324 );
2325 }
2326 }
2327}