use crate::rtpfb::acknowledgement::PacketReport;
use std::time::{Duration, Instant};
pub const DEFAULT_BURST_INTERVAL: Duration = Duration::from_millis(5);
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct ArrivalGroup {
pub first_departure: Instant,
pub departure: Instant,
pub arrival: Duration,
pub packets: usize,
pub size: usize,
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct InterGroupDelay {
pub delta_ms: f64,
pub at: Instant,
pub size: usize,
}
#[derive(Debug, Clone)]
pub struct ArrivalGroupAccumulator {
burst_interval: Duration,
current: Option<ArrivalGroup>,
previous: Option<ArrivalGroup>,
}
impl Default for ArrivalGroupAccumulator {
fn default() -> Self {
Self::new(DEFAULT_BURST_INTERVAL)
}
}
impl ArrivalGroupAccumulator {
pub fn new(burst_interval: Duration) -> Self {
Self {
burst_interval,
current: None,
previous: None,
}
}
pub fn accumulate(&mut self, report: &PacketReport) -> Option<InterGroupDelay> {
let arrival = report.arrival?;
if !report.arrived {
return None;
}
let Some(current) = self.current.as_mut() else {
self.current = Some(ArrivalGroup {
first_departure: report.departure,
departure: report.departure,
arrival,
packets: 1,
size: report.size,
});
return None;
};
if report
.departure
.saturating_duration_since(current.first_departure)
<= self.burst_interval
{
current.departure = current.departure.max(report.departure);
current.arrival = current.arrival.max(arrival);
current.packets += 1;
current.size += report.size;
return None;
}
let closed = *current;
self.current = Some(ArrivalGroup {
first_departure: report.departure,
departure: report.departure,
arrival,
packets: 1,
size: report.size,
});
let measurement = self.previous.map(|previous| gradient(&previous, &closed));
self.previous = Some(closed);
measurement
}
pub fn flush(&mut self) -> Option<InterGroupDelay> {
let closed = self.current.take()?;
let measurement = self.previous.map(|previous| gradient(&previous, &closed));
self.previous = Some(closed);
measurement
}
}
fn gradient(previous: &ArrivalGroup, current: &ArrivalGroup) -> InterGroupDelay {
let arrival_delta = current.arrival.as_secs_f64() - previous.arrival.as_secs_f64();
let departure_delta = current
.departure
.saturating_duration_since(previous.departure)
.as_secs_f64();
InterGroupDelay {
delta_ms: (arrival_delta - departure_delta) * 1_000.0,
at: current.departure,
size: current.size,
}
}
#[cfg(test)]
mod tests {
use super::*;
use rtcp::transport_feedbacks::cc_feedback_report::Ecn;
fn report(departure: Instant, arrival_ms: u64, size: usize) -> PacketReport {
PacketReport {
ssrc: 1,
id: 0,
rtp_sequence_number: 0,
is_twcc: true,
twcc_sequence_number: 0,
size,
arrived: true,
departure,
arrival: Some(Duration::from_millis(arrival_ms)),
ecn: Ecn::default(),
}
}
#[test]
fn packets_sent_together_form_one_group() {
let epoch = Instant::now();
let mut accumulator = ArrivalGroupAccumulator::default();
for offset in [0, 1, 2, 3, 4] {
assert_eq!(
None,
accumulator.accumulate(&report(epoch + Duration::from_millis(offset), 100, 1200)),
"nothing is emitted until a group closes"
);
}
assert_eq!(
None,
accumulator.accumulate(&report(epoch + Duration::from_millis(20), 120, 1200))
);
let group = accumulator
.flush()
.expect("the second group closes against the first");
assert_eq!(1200, group.size, "the second group holds one packet");
}
#[test]
fn a_path_that_does_not_queue_has_a_zero_gradient() {
let epoch = Instant::now();
let mut accumulator = ArrivalGroupAccumulator::default();
let mut gradients = Vec::new();
for burst in 0..5u64 {
let departure = epoch + Duration::from_millis(burst * 20);
if let Some(delay) = accumulate_and_flush(&mut accumulator, departure, 100 + burst * 20)
{
gradients.push(delay.delta_ms);
}
}
assert!(
gradients.iter().all(|delta| delta.abs() < 1e-6),
"a non-queueing path must measure zero delay gradient: {gradients:?}"
);
}
#[test]
fn a_growing_queue_has_a_positive_gradient() {
let epoch = Instant::now();
let mut accumulator = ArrivalGroupAccumulator::default();
let mut gradients = Vec::new();
for burst in 0..5u64 {
let departure = epoch + Duration::from_millis(burst * 20);
if let Some(delay) = accumulate_and_flush(&mut accumulator, departure, 100 + burst * 30)
{
gradients.push(delay.delta_ms);
}
}
assert!(!gradients.is_empty(), "no gradients were produced");
assert!(
gradients.iter().all(|delta| (*delta - 10.0).abs() < 1e-6),
"each group should measure 10 ms of added delay: {gradients:?}"
);
}
#[test]
fn a_draining_queue_has_a_negative_gradient() {
let epoch = Instant::now();
let mut accumulator = ArrivalGroupAccumulator::default();
let mut gradients = Vec::new();
for burst in 0..5u64 {
let departure = epoch + Duration::from_millis(burst * 20);
if let Some(delay) = accumulate_and_flush(&mut accumulator, departure, 200 + burst * 15)
{
gradients.push(delay.delta_ms);
}
}
assert!(
gradients.iter().all(|delta| *delta < 0.0),
"a draining queue must measure negative: {gradients:?}"
);
}
#[test]
fn lost_packets_are_not_measured() {
let epoch = Instant::now();
let mut accumulator = ArrivalGroupAccumulator::default();
let mut lost = report(epoch, 0, 1200);
lost.arrived = false;
lost.arrival = None;
assert_eq!(None, accumulator.accumulate(&lost));
assert_eq!(
None,
accumulator.flush(),
"a lost packet must not open a group"
);
}
fn accumulate_and_flush(
accumulator: &mut ArrivalGroupAccumulator,
departure: Instant,
arrival_ms: u64,
) -> Option<InterGroupDelay> {
accumulator.accumulate(&report(departure, arrival_ms, 1200))
}
}