use std::time::Duration;
use moq_net::Timestamp;
const BITRATE_WINDOW: Duration = Duration::from_secs(1);
#[derive(Default)]
pub struct Metrics {
jitter: Jitter,
bitrate: Bitrate,
}
impl Metrics {
pub fn new() -> Self {
Self::default()
}
pub fn record_frame(&mut self, ts: Timestamp, bytes: usize) -> Option<Duration> {
self.bitrate.observe_frame(ts, bytes);
self.jitter.observe(ts)
}
pub fn record_reorder(&mut self, reorder: Timestamp) -> Option<Duration> {
self.jitter.observe_reorder(reorder)
}
pub fn finish_group(&mut self, next: Option<Timestamp>) -> Option<u64> {
self.bitrate.finish_group(next)
}
pub fn jitter(&self) -> Option<Duration> {
self.jitter.current()
}
pub fn bitrate(&self) -> Option<u64> {
self.bitrate.current()
}
}
#[derive(Default)]
struct Bitrate {
group: Option<Group>,
window_bytes: u64,
window_duration: Duration,
max: Option<u64>,
reported: Option<u64>,
}
impl Bitrate {
fn observe_frame(&mut self, ts: Timestamp, bytes: usize) {
let group = self.group.get_or_insert(Group {
start: ts,
max: ts,
bytes: 0,
});
if ts < group.start {
group.start = ts;
}
if ts > group.max {
group.max = ts;
}
group.bytes = group.bytes.saturating_add(bytes as u64);
}
fn finish_group(&mut self, next: Option<Timestamp>) -> Option<u64> {
let group = self.group.take()?;
let duration = next
.and_then(|next| next.checked_sub(group.start).ok())
.filter(|duration| !duration.is_zero())
.or_else(|| {
group
.max
.checked_sub(group.start)
.ok()
.filter(|duration| !duration.is_zero())
})?;
self.window_bytes = self.window_bytes.saturating_add(group.bytes);
self.window_duration += Duration::from(duration);
if self.window_duration < BITRATE_WINDOW {
return None;
}
let bitrate = bits_per_second(self.window_bytes, self.window_duration);
self.window_bytes = 0;
self.window_duration = Duration::ZERO;
if self.max.is_none_or(|max| bitrate > max) {
self.max = Some(bitrate);
}
if self.reported != self.max {
self.reported = self.max;
return self.max;
}
None
}
fn current(&self) -> Option<u64> {
self.max
}
}
struct Group {
start: Timestamp,
max: Timestamp,
bytes: u64,
}
fn bits_per_second(bytes: u64, duration: Duration) -> u64 {
let nanos = duration.as_nanos();
if nanos == 0 {
return 0;
}
let bits_per_second = (bytes as u128).saturating_mul(8).saturating_mul(1_000_000_000) / nanos;
bits_per_second.min(u64::MAX as u128) as u64
}
#[derive(Default)]
pub struct Jitter {
last_timestamp: Option<Timestamp>,
min_duration: Option<Duration>,
max_reorder: Duration,
reported: Option<Duration>,
}
impl Jitter {
pub fn observe(&mut self, ts: Timestamp) -> Option<Duration> {
if let Some(last) = self.last_timestamp.replace(ts)
&& let Ok(duration) = ts.checked_sub(last)
&& !duration.is_zero()
{
let duration = Duration::from(duration);
self.min_duration = Some(match self.min_duration {
Some(min) => min.min(duration),
None => duration,
});
}
self.report()
}
pub fn observe_reorder(&mut self, reorder: Timestamp) -> Option<Duration> {
self.max_reorder = self.max_reorder.max(Duration::from(reorder));
self.report()
}
pub fn current(&self) -> Option<Duration> {
let jitter = self.combined();
(!jitter.is_zero()).then_some(jitter)
}
fn combined(&self) -> Duration {
self.min_duration.unwrap_or(Duration::ZERO).max(self.max_reorder)
}
fn report(&mut self) -> Option<Duration> {
let jitter = self.combined();
if jitter.is_zero() || self.reported == Some(jitter) {
return None;
}
self.reported = Some(jitter);
Some(jitter)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn micros(value: u64) -> Timestamp {
Timestamp::from_micros(value).unwrap()
}
#[test]
fn reports_min_frame_spacing_once_it_changes() {
let mut jitter = Jitter::default();
assert_eq!(jitter.observe(micros(1_000)), None);
assert_eq!(jitter.observe(micros(41_000)), Some(Duration::from_millis(40)));
assert_eq!(jitter.observe(micros(81_000)), None);
assert_eq!(jitter.observe(micros(101_000)), Some(Duration::from_millis(20)));
assert_eq!(jitter.current(), Some(Duration::from_millis(20)));
}
#[test]
fn reorder_delay_wins_over_frame_spacing() {
let mut jitter = Jitter::default();
assert_eq!(jitter.observe(micros(0)), None);
assert_eq!(jitter.observe(micros(16_000)), Some(Duration::from_millis(16)));
assert_eq!(jitter.observe_reorder(micros(48_000)), Some(Duration::from_millis(48)));
assert_eq!(jitter.observe(micros(32_000)), None);
assert_eq!(jitter.current(), Some(Duration::from_millis(48)));
}
#[test]
fn ignores_non_monotonic_presentation_spacing() {
let mut jitter = Jitter::default();
assert_eq!(jitter.observe(micros(100_000)), None);
assert_eq!(jitter.observe(micros(80_000)), None);
assert_eq!(jitter.observe(micros(120_000)), Some(Duration::from_millis(40)));
}
#[test]
fn bitrate_waits_for_group_boundaries_and_reports_max() {
let mut metrics = Metrics::new();
metrics.record_frame(micros(0), 100_000);
metrics.record_frame(micros(500_000), 100_000);
assert_eq!(metrics.finish_group(Some(micros(1_000_000))), Some(1_600_000));
metrics.record_frame(micros(1_000_000), 25_000);
assert_eq!(metrics.finish_group(Some(micros(2_000_000))), None);
assert_eq!(metrics.bitrate(), Some(1_600_000));
metrics.record_frame(micros(2_000_000), 250_000);
assert_eq!(metrics.finish_group(Some(micros(3_000_000))), Some(2_000_000));
}
#[test]
fn metrics_report_jitter_and_bitrate_independently() {
let mut metrics = Metrics::new();
assert_eq!(metrics.record_frame(micros(0), 1), None);
assert_eq!(metrics.record_frame(micros(33_000), 1), Some(Duration::from_millis(33)));
assert_eq!(
metrics.record_reorder(micros(100_000)),
Some(Duration::from_millis(100))
);
assert_eq!(metrics.jitter(), Some(Duration::from_millis(100)));
}
}