Skip to main content

rtc_interceptor/rtpfb/
convert.rs

1//! Turning the two feedback formats into a common [`Acknowledgement`].
2//!
3//! TWCC and RFC 8888 answer the same question — which packets arrived, and when — in different
4//! shapes. Converting both to one type is what lets congestion control read either without
5//! knowing which is negotiated.
6
7use super::acknowledgement::Acknowledgement;
8use rtcp::transport_feedbacks::cc_feedback_report::{CcFeedbackReport, Ecn};
9use rtcp::transport_feedbacks::transport_layer_cc::{
10    PacketStatusChunk, SymbolTypeTcc, TransportLayerCc,
11};
12use std::collections::HashMap;
13use std::time::Duration;
14
15/// TWCC's reference time counts in multiples of 64 ms.
16const TWCC_REFERENCE_TICK: Duration = Duration::from_millis(64);
17
18/// TWCC receive deltas are in microseconds.
19const TWCC_DELTA_UNIT: Duration = Duration::from_micros(1);
20
21/// RFC 8888 arrival offsets are in units of 1/1024 s.
22const CCFB_OFFSET_DENOMINATOR: u32 = 1024;
23
24/// The RFC 8888 offset meaning "arrived after the report timestamp", which carries no usable
25/// arrival time.
26const CCFB_OFFSET_AFTER_REPORT: u16 = 0x1FFF;
27
28/// Convert a TWCC feedback packet into one acknowledgement per reported packet.
29///
30/// Arrival times are relative to the feedback's own 24-bit reference time, so they are comparable
31/// across successive TWCC reports but not with any other clock.
32pub fn convert_twcc(feedback: &TransportLayerCc) -> Vec<Acknowledgement> {
33    let mut acknowledgements = Vec::new();
34
35    // The reference time is a count of 64 ms ticks on the receiver's clock; deltas accumulate
36    // from there, so each arrival depends on every arrival before it in the same report.
37    let mut arrival = TWCC_REFERENCE_TICK * feedback.reference_time;
38    let mut delta_index = 0usize;
39    let mut offset = 0u16;
40
41    let push = |symbol: SymbolTypeTcc,
42                offset: &mut u16,
43                arrival: &mut Duration,
44                delta_index: &mut usize,
45                acknowledgements: &mut Vec<Acknowledgement>| {
46        let sequence_number = feedback.base_sequence_number.wrapping_add(*offset);
47        *offset = offset.wrapping_add(1);
48
49        match symbol {
50            SymbolTypeTcc::PacketNotReceived => {
51                acknowledgements.push(Acknowledgement::lost(sequence_number));
52            }
53            SymbolTypeTcc::PacketReceivedSmallDelta | SymbolTypeTcc::PacketReceivedLargeDelta => {
54                // A delta may be negative — packets can be reported out of order — so the running
55                // arrival is adjusted in whichever direction the report says.
56                if let Some(delta) = feedback.recv_deltas.get(*delta_index) {
57                    *delta_index += 1;
58                    let magnitude = TWCC_DELTA_UNIT * delta.delta.unsigned_abs() as u32;
59                    *arrival = if delta.delta < 0 {
60                        arrival.saturating_sub(magnitude)
61                    } else {
62                        *arrival + magnitude
63                    };
64                    acknowledgements.push(Acknowledgement::received(
65                        sequence_number,
66                        Some(*arrival),
67                        Ecn::NotEct,
68                    ));
69                } else {
70                    // The chunks claim more received packets than there are deltas. Report the
71                    // packet as arrived without a time rather than inventing one.
72                    acknowledgements.push(Acknowledgement::received(
73                        sequence_number,
74                        None,
75                        Ecn::NotEct,
76                    ));
77                }
78            }
79            SymbolTypeTcc::PacketReceivedWithoutDelta => {
80                acknowledgements.push(Acknowledgement::received(
81                    sequence_number,
82                    None,
83                    Ecn::NotEct,
84                ));
85            }
86        }
87    };
88
89    for chunk in &feedback.packet_chunks {
90        match chunk {
91            PacketStatusChunk::RunLengthChunk(run) => {
92                for _ in 0..run.run_length {
93                    push(
94                        run.packet_status_symbol,
95                        &mut offset,
96                        &mut arrival,
97                        &mut delta_index,
98                        &mut acknowledgements,
99                    );
100                }
101            }
102            PacketStatusChunk::StatusVectorChunk(vector) => {
103                for &symbol in &vector.symbol_list {
104                    push(
105                        symbol,
106                        &mut offset,
107                        &mut arrival,
108                        &mut delta_index,
109                        &mut acknowledgements,
110                    );
111                }
112            }
113        }
114    }
115
116    acknowledgements
117}
118
119/// Convert an RFC 8888 report into acknowledgements per media stream, with the delay the receiver
120/// added before sending it.
121///
122/// The returned delay is the gap between the newest arrival the report describes and the instant
123/// it was stamped. A round trip measured without subtracting it counts the receiver's own
124/// reporting interval as network time.
125pub fn convert_ccfb(feedback: &CcFeedbackReport) -> (Duration, HashMap<u32, Vec<Acknowledgement>>) {
126    let mut per_stream = HashMap::new();
127    // Arrivals are offsets *back* from the report timestamp, so the newest arrival is the
128    // smallest offset.
129    let mut newest_arrival: Option<Duration> = None;
130
131    for block in &feedback.report_blocks {
132        let mut acknowledgements = Vec::with_capacity(block.metric_blocks.len());
133
134        for (index, metric) in block.metric_blocks.iter().enumerate() {
135            let sequence_number = block.begin_sequence.wrapping_add(index as u16);
136
137            if !metric.received {
138                acknowledgements.push(Acknowledgement::lost(sequence_number));
139                continue;
140            }
141
142            // The reserved offset says the packet arrived after the report was stamped, which
143            // leaves no usable time.
144            let arrival = if metric.arrival_time_offset == CCFB_OFFSET_AFTER_REPORT {
145                None
146            } else {
147                let offset = Duration::from_secs_f64(
148                    f64::from(metric.arrival_time_offset) / f64::from(CCFB_OFFSET_DENOMINATOR),
149                );
150                newest_arrival = Some(match newest_arrival {
151                    Some(newest) => newest.min(offset),
152                    None => offset,
153                });
154                Some(offset)
155            };
156
157            acknowledgements.push(Acknowledgement::received(
158                sequence_number,
159                arrival,
160                metric.ecn,
161            ));
162        }
163
164        per_stream.insert(block.media_ssrc, acknowledgements);
165    }
166
167    (newest_arrival.unwrap_or_default(), per_stream)
168}
169
170#[cfg(test)]
171mod tests {
172    use super::*;
173    use rtcp::transport_feedbacks::cc_feedback_report::{
174        CcFeedbackMetricBlock, CcFeedbackReportBlock,
175    };
176    use rtcp::transport_feedbacks::transport_layer_cc::{
177        RecvDelta, RunLengthChunk, StatusVectorChunk,
178    };
179
180    fn run_length(symbol: SymbolTypeTcc, run_length: u16) -> PacketStatusChunk {
181        PacketStatusChunk::RunLengthChunk(RunLengthChunk {
182            packet_status_symbol: symbol,
183            run_length,
184            ..Default::default()
185        })
186    }
187
188    fn status_vector(symbols: Vec<SymbolTypeTcc>) -> PacketStatusChunk {
189        PacketStatusChunk::StatusVectorChunk(StatusVectorChunk {
190            symbol_list: symbols,
191            ..Default::default()
192        })
193    }
194
195    fn delta(microseconds: i64) -> RecvDelta {
196        RecvDelta {
197            delta: microseconds,
198            ..Default::default()
199        }
200    }
201
202    // ---------------------------------------------------------------------------------------
203    // TWCC
204    // ---------------------------------------------------------------------------------------
205
206    #[test]
207    fn a_run_of_lost_packets_converts_to_losses() {
208        let feedback = TransportLayerCc {
209            base_sequence_number: 100,
210            reference_time: 0,
211            packet_chunks: vec![run_length(SymbolTypeTcc::PacketNotReceived, 3)],
212            recv_deltas: vec![],
213            ..Default::default()
214        };
215
216        let acknowledgements = convert_twcc(&feedback);
217        assert_eq!(
218            vec![
219                Acknowledgement::lost(100),
220                Acknowledgement::lost(101),
221                Acknowledgement::lost(102)
222            ],
223            acknowledgements
224        );
225    }
226
227    /// Deltas accumulate: each arrival is relative to the one before it, so a report describes a
228    /// *sequence* of arrivals rather than independent timestamps.
229    #[test]
230    fn deltas_accumulate_from_the_reference_time() {
231        let feedback = TransportLayerCc {
232            base_sequence_number: 10,
233            // One tick of 64 ms.
234            reference_time: 1,
235            packet_chunks: vec![run_length(SymbolTypeTcc::PacketReceivedSmallDelta, 3)],
236            recv_deltas: vec![delta(1000), delta(2000), delta(500)],
237            ..Default::default()
238        };
239
240        let acknowledgements = convert_twcc(&feedback);
241        let arrivals: Vec<Duration> = acknowledgements
242            .iter()
243            .map(|ack| ack.arrival.expect("arrived with a time"))
244            .collect();
245
246        assert_eq!(
247            vec![
248                Duration::from_millis(64) + Duration::from_micros(1000),
249                Duration::from_millis(64) + Duration::from_micros(3000),
250                Duration::from_millis(64) + Duration::from_micros(3500),
251            ],
252            arrivals
253        );
254        assert!(acknowledgements.iter().all(|ack| ack.arrived));
255    }
256
257    /// Packets can be reported arriving out of order, which is a negative delta. Treating it as
258    /// positive would report the path as *less* delayed than it is.
259    #[test]
260    fn a_negative_delta_moves_the_arrival_backwards() {
261        let feedback = TransportLayerCc {
262            base_sequence_number: 0,
263            reference_time: 1,
264            packet_chunks: vec![run_length(SymbolTypeTcc::PacketReceivedSmallDelta, 2)],
265            recv_deltas: vec![delta(5000), delta(-2000)],
266            ..Default::default()
267        };
268
269        let arrivals: Vec<Duration> = convert_twcc(&feedback)
270            .iter()
271            .map(|ack| ack.arrival.expect("arrived"))
272            .collect();
273
274        assert_eq!(
275            Duration::from_millis(64) + Duration::from_micros(5000),
276            arrivals[0]
277        );
278        assert_eq!(
279            Duration::from_millis(64) + Duration::from_micros(3000),
280            arrivals[1],
281            "the second packet arrived before the first"
282        );
283    }
284
285    #[test]
286    fn a_status_vector_converts_symbol_by_symbol() {
287        let feedback = TransportLayerCc {
288            base_sequence_number: 500,
289            reference_time: 0,
290            packet_chunks: vec![status_vector(vec![
291                SymbolTypeTcc::PacketReceivedSmallDelta,
292                SymbolTypeTcc::PacketNotReceived,
293                SymbolTypeTcc::PacketReceivedLargeDelta,
294            ])],
295            recv_deltas: vec![delta(100), delta(200)],
296            ..Default::default()
297        };
298
299        let acknowledgements = convert_twcc(&feedback);
300        assert_eq!(
301            vec![true, false, true],
302            acknowledgements
303                .iter()
304                .map(|ack| ack.arrived)
305                .collect::<Vec<_>>()
306        );
307        assert_eq!(
308            vec![500, 501, 502],
309            acknowledgements
310                .iter()
311                .map(|ack| ack.sequence_number)
312                .collect::<Vec<_>>()
313        );
314    }
315
316    /// "Received without delta" means the receiver got it but cannot say when. Reporting an
317    /// invented time would feed congestion control a measurement nobody made.
318    #[test]
319    fn a_packet_received_without_a_delta_has_no_arrival_time() {
320        let feedback = TransportLayerCc {
321            base_sequence_number: 0,
322            reference_time: 0,
323            packet_chunks: vec![run_length(SymbolTypeTcc::PacketReceivedWithoutDelta, 1)],
324            recv_deltas: vec![],
325            ..Default::default()
326        };
327
328        let acknowledgements = convert_twcc(&feedback);
329        assert!(acknowledgements[0].arrived);
330        assert_eq!(None, acknowledgements[0].arrival);
331    }
332
333    /// A report claiming more received packets than it carries deltas for is malformed. The
334    /// packets are still reported as arrived, without times, rather than panicking on the index.
335    #[test]
336    fn more_received_packets_than_deltas_does_not_panic() {
337        let feedback = TransportLayerCc {
338            base_sequence_number: 0,
339            reference_time: 0,
340            packet_chunks: vec![run_length(SymbolTypeTcc::PacketReceivedSmallDelta, 4)],
341            recv_deltas: vec![delta(100)],
342            ..Default::default()
343        };
344
345        let acknowledgements = convert_twcc(&feedback);
346        assert_eq!(4, acknowledgements.len());
347        assert!(acknowledgements[0].arrival.is_some(), "the one real delta");
348        assert!(
349            acknowledgements[1..]
350                .iter()
351                .all(|ack| ack.arrival.is_none()),
352            "and no invented times for the rest"
353        );
354    }
355
356    #[test]
357    fn sequence_numbers_wrap_across_a_report() {
358        let feedback = TransportLayerCc {
359            base_sequence_number: 65534,
360            reference_time: 0,
361            packet_chunks: vec![run_length(SymbolTypeTcc::PacketNotReceived, 4)],
362            recv_deltas: vec![],
363            ..Default::default()
364        };
365
366        assert_eq!(
367            vec![65534, 65535, 0, 1],
368            convert_twcc(&feedback)
369                .iter()
370                .map(|ack| ack.sequence_number)
371                .collect::<Vec<_>>()
372        );
373    }
374
375    #[test]
376    fn an_empty_feedback_converts_to_nothing() {
377        let feedback = TransportLayerCc::default();
378        assert!(convert_twcc(&feedback).is_empty());
379    }
380
381    // ---------------------------------------------------------------------------------------
382    // RFC 8888
383    // ---------------------------------------------------------------------------------------
384
385    fn metric(received: bool, offset: u16, ecn: Ecn) -> CcFeedbackMetricBlock {
386        CcFeedbackMetricBlock {
387            received,
388            ecn,
389            arrival_time_offset: offset,
390        }
391    }
392
393    #[test]
394    fn a_ccfb_report_converts_per_stream() {
395        let feedback = CcFeedbackReport {
396            sender_ssrc: 1,
397            report_blocks: vec![
398                CcFeedbackReportBlock {
399                    media_ssrc: 10,
400                    begin_sequence: 100,
401                    metric_blocks: vec![
402                        metric(true, 512, Ecn::NotEct),
403                        metric(false, 0, Ecn::NotEct),
404                    ],
405                },
406                CcFeedbackReportBlock {
407                    media_ssrc: 20,
408                    begin_sequence: 5,
409                    metric_blocks: vec![metric(true, 256, Ecn::Ce)],
410                },
411            ],
412            report_timestamp: 0,
413        };
414
415        let (_, per_stream) = convert_ccfb(&feedback);
416
417        let first = &per_stream[&10];
418        assert_eq!(100, first[0].sequence_number);
419        assert_eq!(Some(Duration::from_millis(500)), first[0].arrival);
420        assert!(!first[1].arrived);
421
422        let second = &per_stream[&20];
423        assert_eq!(5, second[0].sequence_number);
424        assert_eq!(Ecn::Ce, second[0].ecn, "ECN survives conversion");
425    }
426
427    /// The receiver's own reporting delay is not network time. Without subtracting it, a round
428    /// trip includes however long the receiver sat on the report before sending it.
429    #[test]
430    fn the_reporting_delay_is_the_gap_to_the_newest_arrival() {
431        let feedback = CcFeedbackReport {
432            sender_ssrc: 1,
433            report_blocks: vec![CcFeedbackReportBlock {
434                media_ssrc: 10,
435                begin_sequence: 0,
436                // 1024 units = 1 s ago, 102 ≈ 100 ms ago. The newest is the smaller offset.
437                metric_blocks: vec![
438                    metric(true, 1024, Ecn::NotEct),
439                    metric(true, 102, Ecn::NotEct),
440                ],
441            }],
442            report_timestamp: 0,
443        };
444
445        let (delay, _) = convert_ccfb(&feedback);
446        assert!(
447            (Duration::from_millis(99)..=Duration::from_millis(101)).contains(&delay),
448            "the newest arrival was ~100 ms before the report, got {delay:?}"
449        );
450    }
451
452    /// The reserved offset means the packet arrived after the report was stamped, so there is no
453    /// usable time — and it must not be taken as the newest arrival either.
454    #[test]
455    fn the_reserved_offset_yields_no_arrival_time() {
456        let feedback = CcFeedbackReport {
457            sender_ssrc: 1,
458            report_blocks: vec![CcFeedbackReportBlock {
459                media_ssrc: 10,
460                begin_sequence: 0,
461                metric_blocks: vec![
462                    metric(true, CCFB_OFFSET_AFTER_REPORT, Ecn::NotEct),
463                    metric(true, 512, Ecn::NotEct),
464                ],
465            }],
466            report_timestamp: 0,
467        };
468
469        let (delay, per_stream) = convert_ccfb(&feedback);
470        assert!(per_stream[&10][0].arrived);
471        assert_eq!(None, per_stream[&10][0].arrival);
472        assert_eq!(
473            Duration::from_millis(500),
474            delay,
475            "the reserved value did not become the newest arrival"
476        );
477    }
478
479    #[test]
480    fn a_report_with_no_arrivals_has_no_reporting_delay() {
481        let feedback = CcFeedbackReport {
482            sender_ssrc: 1,
483            report_blocks: vec![CcFeedbackReportBlock {
484                media_ssrc: 10,
485                begin_sequence: 0,
486                metric_blocks: vec![metric(false, 0, Ecn::NotEct)],
487            }],
488            report_timestamp: 0,
489        };
490
491        let (delay, per_stream) = convert_ccfb(&feedback);
492        assert_eq!(Duration::ZERO, delay);
493        assert!(!per_stream[&10][0].arrived);
494    }
495
496    #[test]
497    fn ccfb_sequence_numbers_wrap_within_a_block() {
498        let feedback = CcFeedbackReport {
499            sender_ssrc: 1,
500            report_blocks: vec![CcFeedbackReportBlock {
501                media_ssrc: 10,
502                begin_sequence: 65535,
503                metric_blocks: vec![metric(false, 0, Ecn::NotEct), metric(false, 0, Ecn::NotEct)],
504            }],
505            report_timestamp: 0,
506        };
507
508        let (_, per_stream) = convert_ccfb(&feedback);
509        assert_eq!(
510            vec![65535, 0],
511            per_stream[&10]
512                .iter()
513                .map(|ack| ack.sequence_number)
514                .collect::<Vec<_>>()
515        );
516    }
517}