use crate::jitterbuffer::sequence::SequenceExtender;
use rtcp::transport_feedbacks::cc_feedback_report::{
CcFeedbackMetricBlock, CcFeedbackReportBlock, Ecn,
};
use std::collections::HashMap;
use std::time::Instant;
pub(crate) const MAX_REPORTS_PER_BLOCK: usize = 16384;
const OFFSET_AFTER_REPORT: u16 = 0x1FFF;
const OFFSET_TOO_OLD: u16 = 0x1FFE;
const OFFSET_MAX: u16 = 0x1FFD;
#[derive(Debug, Clone, Copy)]
struct PacketReport {
arrival: Instant,
ecn: Ecn,
}
#[derive(Debug)]
pub(crate) struct StreamLog {
ssrc: u32,
extender: SequenceExtender,
next_to_report: u64,
highest_received: u64,
initialised: bool,
log: HashMap<u64, PacketReport>,
}
impl StreamLog {
pub(crate) fn new(ssrc: u32) -> Self {
Self {
ssrc,
extender: SequenceExtender::new(),
next_to_report: 0,
highest_received: 0,
initialised: false,
log: HashMap::new(),
}
}
pub(crate) fn add(&mut self, arrival: Instant, sequence_number: u16, ecn: Ecn) {
let extended = self.extender.extend(sequence_number);
if !self.initialised {
self.initialised = true;
self.next_to_report = extended;
}
if extended < self.next_to_report {
return;
}
self.log.insert(extended, PacketReport { arrival, ecn });
self.highest_received = self.highest_received.max(extended);
}
pub(crate) fn metrics_after(
&mut self,
reference: Instant,
max_blocks: usize,
) -> CcFeedbackReportBlock {
if self.log.is_empty() {
return CcFeedbackReportBlock {
media_ssrc: self.ssrc,
begin_sequence: self.next_to_report as u16,
metric_blocks: Vec::new(),
};
}
let mut count = self.highest_received - self.next_to_report + 1;
if count > max_blocks as u64 {
count = max_blocks as u64;
let new_next = self.highest_received + 1 - count;
self.log
.retain(|&sequence_number, _| sequence_number >= new_next);
self.next_to_report = new_next;
}
if count == 0 {
return CcFeedbackReportBlock {
media_ssrc: self.ssrc,
begin_sequence: self.next_to_report as u16,
metric_blocks: Vec::new(),
};
}
let begin = self.next_to_report;
let mut metric_blocks = Vec::with_capacity(count as usize);
for extended in begin..=self.highest_received {
let report = self.log.get(&extended).copied();
metric_blocks.push(match report {
Some(report) => CcFeedbackMetricBlock {
received: true,
ecn: report.ecn,
arrival_time_offset: arrival_time_offset(reference, report.arrival),
},
None => CcFeedbackMetricBlock::default(),
});
if report.is_some() && extended == self.next_to_report {
self.log.remove(&extended);
self.next_to_report += 1;
}
}
CcFeedbackReportBlock {
media_ssrc: self.ssrc,
begin_sequence: begin as u16,
metric_blocks,
}
}
pub(crate) fn is_empty(&self) -> bool {
self.log.is_empty()
}
}
fn arrival_time_offset(reference: Instant, arrival: Instant) -> u16 {
if arrival > reference {
return OFFSET_AFTER_REPORT;
}
let offset = reference.duration_since(arrival).as_secs_f64() * 1024.0;
if offset > f64::from(OFFSET_MAX) {
return OFFSET_TOO_OLD;
}
offset as u16
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
fn received(block: &CcFeedbackReportBlock) -> Vec<bool> {
block
.metric_blocks
.iter()
.map(|metric| metric.received)
.collect()
}
#[test]
fn an_empty_log_reports_nothing() {
let mut log = StreamLog::new(1);
let block = log.metrics_after(Instant::now(), MAX_REPORTS_PER_BLOCK);
assert_eq!(1, block.media_ssrc);
assert!(block.metric_blocks.is_empty());
}
#[test]
fn a_contiguous_run_is_reported_once_and_then_forgotten() {
let now = Instant::now();
let mut log = StreamLog::new(1);
for offset in 0..3u16 {
log.add(now, 100 + offset, Ecn::NotEct);
}
let block = log.metrics_after(now, MAX_REPORTS_PER_BLOCK);
assert_eq!(100, block.begin_sequence);
assert_eq!(vec![true, true, true], received(&block));
assert!(
log.is_empty(),
"everything reported was contiguous, so nothing is held"
);
let second = log.metrics_after(now, MAX_REPORTS_PER_BLOCK);
assert!(
second.metric_blocks.is_empty(),
"and it is not reported a second time"
);
}
#[test]
fn reporting_stops_advancing_at_a_gap() {
let now = Instant::now();
let mut log = StreamLog::new(1);
log.add(now, 100, Ecn::NotEct);
log.add(now, 102, Ecn::NotEct);
let block = log.metrics_after(now, MAX_REPORTS_PER_BLOCK);
assert_eq!(100, block.begin_sequence);
assert_eq!(vec![true, false, true], received(&block));
log.add(now, 101, Ecn::NotEct);
let second = log.metrics_after(now, MAX_REPORTS_PER_BLOCK);
assert_eq!(101, second.begin_sequence, "resumes where it stopped");
assert_eq!(vec![true, true], received(&second));
}
#[test]
fn a_packet_older_than_the_window_is_not_reported_again() {
let now = Instant::now();
let mut log = StreamLog::new(1);
log.add(now, 100, Ecn::NotEct);
log.metrics_after(now, MAX_REPORTS_PER_BLOCK);
log.add(now, 100, Ecn::NotEct);
assert!(
log.is_empty(),
"already reported: saying so again would claim it arrived twice"
);
}
#[test]
fn the_window_gives_up_its_oldest_end_rather_than_overflowing_the_block() {
let now = Instant::now();
let mut log = StreamLog::new(1);
for offset in 0..10u16 {
log.add(now, 100 + offset, Ecn::NotEct);
}
let block = log.metrics_after(now, 4);
assert_eq!(4, block.metric_blocks.len(), "capped at the limit");
assert_eq!(
106, block.begin_sequence,
"the newest four, since the oldest are the least useful"
);
}
#[test]
fn ecn_markings_are_carried_through() {
let now = Instant::now();
let mut log = StreamLog::new(1);
log.add(now, 100, Ecn::Ce);
log.add(now, 101, Ecn::Ect1);
let block = log.metrics_after(now, MAX_REPORTS_PER_BLOCK);
assert_eq!(Ecn::Ce, block.metric_blocks[0].ecn);
assert_eq!(Ecn::Ect1, block.metric_blocks[1].ecn);
}
#[test]
fn a_sequence_number_wrap_does_not_reorder_the_window() {
let now = Instant::now();
let mut log = StreamLog::new(1);
for sequence_number in [65534u16, 65535, 0, 1] {
log.add(now, sequence_number, Ecn::NotEct);
}
let block = log.metrics_after(now, MAX_REPORTS_PER_BLOCK);
assert_eq!(65534, block.begin_sequence);
assert_eq!(
vec![true, true, true, true],
received(&block),
"0 follows 65535 rather than opening a 65534-packet gap"
);
}
#[test]
fn an_offset_is_measured_back_from_the_report_timestamp() {
let now = Instant::now();
assert_eq!(
512,
arrival_time_offset(now, now - Duration::from_millis(500))
);
assert_eq!(0, arrival_time_offset(now, now));
}
#[test]
fn offsets_clamp_at_both_edges() {
let now = Instant::now();
assert_eq!(
OFFSET_AFTER_REPORT,
arrival_time_offset(now, now + Duration::from_millis(1)),
"arrived after the report was stamped"
);
let far_past = now - Duration::from_secs(60);
assert_eq!(OFFSET_TOO_OLD, arrival_time_offset(now, far_past));
let at_limit = now - Duration::from_secs_f64(f64::from(OFFSET_MAX) / 1024.0);
assert!(
arrival_time_offset(now, at_limit) <= OFFSET_MAX,
"the largest expressible offset is not mistaken for a reserved value"
);
}
}