Skip to main content

rtc_interceptor/rfc8888/
recorder.rs

1//! Turns observed arrivals into an RFC 8888 feedback report.
2
3use super::stream_log::{MAX_REPORTS_PER_BLOCK, StreamLog};
4use rtcp::transport_feedbacks::cc_feedback_report::{CcFeedbackReport, Ecn};
5use std::collections::HashMap;
6use std::time::Instant;
7
8/// Header, sender SSRC and report timestamp — the fixed cost of any report.
9const REPORT_OVERHEAD: usize = 12;
10
11/// Media SSRC, base sequence number and count — the fixed cost of each report block.
12const BLOCK_OVERHEAD: usize = 8;
13
14/// Bytes per reported packet.
15const METRIC_BLOCK_SIZE: usize = 2;
16
17/// Records packet arrivals per stream and builds feedback reports from them.
18#[derive(Debug, Default)]
19pub struct CcFeedbackRecorder {
20    streams: HashMap<u32, StreamLog>,
21}
22
23impl CcFeedbackRecorder {
24    /// A recorder with nothing observed yet.
25    pub fn new() -> Self {
26        Self::default()
27    }
28
29    /// Note that a packet of `ssrc` arrived at `arrival`.
30    pub fn add_packet(&mut self, arrival: Instant, ssrc: u32, sequence_number: u16, ecn: Ecn) {
31        self.streams
32            .entry(ssrc)
33            .or_insert_with(|| StreamLog::new(ssrc))
34            .add(arrival, sequence_number, ecn);
35    }
36
37    /// Whether anything is waiting to be reported.
38    pub fn is_empty(&self) -> bool {
39        self.streams.values().all(StreamLog::is_empty)
40    }
41
42    /// Build a report of everything observed since the last one.
43    ///
44    /// `max_size` is the byte budget for the whole report, shared evenly between streams: a
45    /// report that does not fit the path's MTU would be fragmented or dropped, and feedback that
46    /// does not arrive is worse than feedback that describes fewer packets.
47    ///
48    /// `report_timestamp` is the middle 32 bits of an NTP timestamp. It is supplied rather than
49    /// read from a clock because a sans-I/O interceptor has none — `shared::time::SystemInstant`
50    /// is how the caller converts the monotonic `now` into wall-clock time.
51    pub fn build_report(
52        &mut self,
53        now: Instant,
54        sender_ssrc: u32,
55        report_timestamp: u32,
56        max_size: usize,
57    ) -> CcFeedbackReport {
58        let mut report = CcFeedbackReport {
59            sender_ssrc,
60            report_blocks: Vec::new(),
61            report_timestamp,
62        };
63
64        let stream_count = self.streams.len();
65        if stream_count == 0 {
66            return report;
67        }
68
69        let budget = max_size
70            .saturating_sub(REPORT_OVERHEAD)
71            .saturating_sub(BLOCK_OVERHEAD * stream_count)
72            / METRIC_BLOCK_SIZE;
73        let per_stream = (budget / stream_count).min(MAX_REPORTS_PER_BLOCK);
74
75        // Ordered so a report is reproducible; a `HashMap` would otherwise vary run to run.
76        let mut ssrcs: Vec<u32> = self.streams.keys().copied().collect();
77        ssrcs.sort_unstable();
78
79        for ssrc in ssrcs {
80            let Some(stream) = self.streams.get_mut(&ssrc) else {
81                continue;
82            };
83            report
84                .report_blocks
85                .push(stream.metrics_after(now, per_stream));
86        }
87
88        report
89    }
90
91    /// Forget a stream entirely.
92    pub fn remove_stream(&mut self, ssrc: u32) {
93        self.streams.remove(&ssrc);
94    }
95}
96
97#[cfg(test)]
98mod tests {
99    use super::*;
100
101    trait ReportSsrcsForTest {
102        fn destination_ssrcs_for_test(&self) -> Vec<u32>;
103    }
104
105    impl ReportSsrcsForTest for CcFeedbackReport {
106        fn destination_ssrcs_for_test(&self) -> Vec<u32> {
107            self.report_blocks
108                .iter()
109                .map(|block| block.media_ssrc)
110                .collect()
111        }
112    }
113
114    const MTU: usize = 1200;
115
116    #[test]
117    fn a_recorder_with_nothing_observed_reports_no_blocks() {
118        let mut recorder = CcFeedbackRecorder::new();
119        let report = recorder.build_report(Instant::now(), 7, 42, MTU);
120
121        assert_eq!(7, report.sender_ssrc);
122        assert_eq!(42, report.report_timestamp);
123        assert!(report.report_blocks.is_empty());
124        assert!(recorder.is_empty());
125    }
126
127    #[test]
128    fn each_stream_gets_its_own_report_block() {
129        let now = Instant::now();
130        let mut recorder = CcFeedbackRecorder::new();
131        recorder.add_packet(now, 1, 100, Ecn::NotEct);
132        recorder.add_packet(now, 2, 500, Ecn::NotEct);
133        recorder.add_packet(now, 1, 101, Ecn::NotEct);
134
135        let report = recorder.build_report(now, 0, 0, MTU);
136
137        assert_eq!(2, report.report_blocks.len());
138        assert_eq!(vec![1, 2], report.destination_ssrcs_for_test());
139        assert_eq!(
140            2,
141            report.report_blocks[0].metric_blocks.len(),
142            "two on ssrc 1"
143        );
144        assert_eq!(
145            1,
146            report.report_blocks[1].metric_blocks.len(),
147            "one on ssrc 2"
148        );
149    }
150
151    /// Report blocks come out in a stable order, so two runs over the same input produce the same
152    /// bytes — a `HashMap` iteration order would not.
153    #[test]
154    fn report_blocks_are_ordered_by_ssrc() {
155        let now = Instant::now();
156        let mut recorder = CcFeedbackRecorder::new();
157        for ssrc in [900, 100, 500, 300] {
158            recorder.add_packet(now, ssrc, 1, Ecn::NotEct);
159        }
160
161        let report = recorder.build_report(now, 0, 0, MTU);
162        assert_eq!(
163            vec![100, 300, 500, 900],
164            report.destination_ssrcs_for_test()
165        );
166    }
167
168    /// A report has to fit the path. Describing fewer packets is better than a report that gets
169    /// fragmented or dropped, since feedback that does not arrive is worth nothing.
170    #[test]
171    fn the_byte_budget_bounds_what_a_report_describes() {
172        let now = Instant::now();
173        let mut recorder = CcFeedbackRecorder::new();
174        for sequence_number in 0..100u16 {
175            recorder.add_packet(now, 1, sequence_number, Ecn::NotEct);
176        }
177
178        // 12 report overhead + 8 block overhead leaves 20 bytes, i.e. 10 metric blocks.
179        let report = recorder.build_report(now, 0, 0, 40);
180        assert_eq!(10, report.report_blocks[0].metric_blocks.len());
181    }
182
183    #[test]
184    fn a_budget_too_small_for_any_packet_reports_none() {
185        let now = Instant::now();
186        let mut recorder = CcFeedbackRecorder::new();
187        recorder.add_packet(now, 1, 1, Ecn::NotEct);
188
189        let report = recorder.build_report(now, 0, 0, 4);
190        assert!(report.report_blocks[0].metric_blocks.is_empty());
191    }
192
193    #[test]
194    fn the_budget_is_shared_between_streams() {
195        let now = Instant::now();
196        let mut recorder = CcFeedbackRecorder::new();
197        for sequence_number in 0..50u16 {
198            recorder.add_packet(now, 1, sequence_number, Ecn::NotEct);
199            recorder.add_packet(now, 2, sequence_number, Ecn::NotEct);
200        }
201
202        // 12 + 8*2 = 28 overhead; 72 bytes left is 36 metric blocks, 18 per stream.
203        let report = recorder.build_report(now, 0, 0, 100);
204        assert_eq!(18, report.report_blocks[0].metric_blocks.len());
205        assert_eq!(18, report.report_blocks[1].metric_blocks.len());
206    }
207
208    #[test]
209    fn removing_a_stream_stops_it_being_reported() {
210        let now = Instant::now();
211        let mut recorder = CcFeedbackRecorder::new();
212        recorder.add_packet(now, 1, 1, Ecn::NotEct);
213        recorder.add_packet(now, 2, 1, Ecn::NotEct);
214
215        recorder.remove_stream(1);
216        let report = recorder.build_report(now, 0, 0, MTU);
217
218        assert_eq!(vec![2], report.destination_ssrcs_for_test());
219    }
220
221    /// A report describing packets already reported would tell the sender they arrived twice.
222    #[test]
223    fn a_second_report_describes_only_what_arrived_since() {
224        let now = Instant::now();
225        let mut recorder = CcFeedbackRecorder::new();
226        recorder.add_packet(now, 1, 100, Ecn::NotEct);
227
228        let first = recorder.build_report(now, 0, 0, MTU);
229        assert_eq!(1, first.report_blocks[0].metric_blocks.len());
230
231        let second = recorder.build_report(now, 0, 0, MTU);
232        assert!(second.report_blocks[0].metric_blocks.is_empty());
233
234        recorder.add_packet(now, 1, 101, Ecn::NotEct);
235        let third = recorder.build_report(now, 0, 0, MTU);
236        assert_eq!(1, third.report_blocks[0].metric_blocks.len());
237        assert_eq!(101, third.report_blocks[0].begin_sequence);
238    }
239}