1use 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
15const TWCC_REFERENCE_TICK: Duration = Duration::from_millis(64);
17
18const TWCC_DELTA_UNIT: Duration = Duration::from_micros(1);
20
21const CCFB_OFFSET_DENOMINATOR: u32 = 1024;
23
24const CCFB_OFFSET_AFTER_REPORT: u16 = 0x1FFF;
27
28pub fn convert_twcc(feedback: &TransportLayerCc) -> Vec<Acknowledgement> {
33 let mut acknowledgements = Vec::new();
34
35 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 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 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
119pub fn convert_ccfb(feedback: &CcFeedbackReport) -> (Duration, HashMap<u32, Vec<Acknowledgement>>) {
126 let mut per_stream = HashMap::new();
127 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 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 #[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 #[test]
230 fn deltas_accumulate_from_the_reference_time() {
231 let feedback = TransportLayerCc {
232 base_sequence_number: 10,
233 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 #[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 #[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 #[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 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 #[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 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 #[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}