1use 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
8const REPORT_OVERHEAD: usize = 12;
10
11const BLOCK_OVERHEAD: usize = 8;
13
14const METRIC_BLOCK_SIZE: usize = 2;
16
17#[derive(Debug, Default)]
19pub struct CcFeedbackRecorder {
20 streams: HashMap<u32, StreamLog>,
21}
22
23impl CcFeedbackRecorder {
24 pub fn new() -> Self {
26 Self::default()
27 }
28
29 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 pub fn is_empty(&self) -> bool {
39 self.streams.values().all(StreamLog::is_empty)
40 }
41
42 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 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 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 #[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 #[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 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 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 #[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}