use std::time::{Duration, Instant};
use hang::moq_net::Timestamp;
pub(super) struct Presentation {
pacer: moq_mux::Pacer,
delay: Duration,
speaker: bool,
restart: bool,
}
impl Presentation {
pub(super) fn new(delay: Duration) -> Self {
Self {
pacer: moq_mux::Pacer::default().with_delay(delay),
delay,
speaker: false,
restart: false,
}
}
pub(super) fn video(&mut self, timestamp: Timestamp, now: Instant) -> bool {
if self.speaker {
return false;
}
let before = self.pacer.due(timestamp);
before != Some(self.pacer.pace(timestamp, now))
}
pub(super) fn video_restarted(&mut self) {
if !self.speaker {
self.pacer = moq_mux::Pacer::default().with_delay(self.delay);
}
}
pub(super) fn audio(&mut self, end: Duration, buffered: Duration, now: Instant) -> bool {
let timestamp = media(end);
let edge = now
.checked_add(buffered)
.and_then(|at| at.checked_sub(self.delay))
.unwrap_or(now);
let before = self.pacer.due(timestamp);
let after = if !self.speaker || std::mem::take(&mut self.restart) {
self.pacer.hurry(timestamp, edge)
} else {
self.pacer.pace(timestamp, edge)
};
self.speaker = true;
before != Some(after)
}
pub(super) fn restarted(&mut self) {
self.restart = true;
}
pub(super) fn stopped(&mut self) {
self.speaker = false;
self.restart = false;
}
pub(super) fn due(&self, timestamp: Timestamp) -> Option<Instant> {
self.pacer.due(timestamp)
}
}
pub(super) fn timestamp(timestamp: Timestamp) -> Duration {
Duration::from_micros(timestamp.as_micros().min(u64::MAX as u128) as u64)
}
fn media(duration: Duration) -> Timestamp {
const MAX: u128 = (1 << 62) - 1;
Timestamp::from_micros(duration.as_micros().min(MAX) as u64).expect("clamped to the wire maximum")
}
#[derive(Default)]
pub(super) struct AudioTimeline {
origin: Option<Duration>,
end: Option<Duration>,
written: u64,
}
pub(super) struct AudioTiming {
pub(super) end: Duration,
pub(super) silence: u64,
pub(super) reset_sink: bool,
}
impl AudioTimeline {
pub(super) fn push(&mut self, start: Duration, samples: usize, sample_rate: u32, fill_max: u64) -> AudioTiming {
let duration = Duration::from_secs_f64(samples as f64 / sample_rate as f64);
let end = start.saturating_add(duration);
let tolerance = Duration::from_millis(1).saturating_add(Duration::from_secs_f64(2.0 / sample_rate as f64));
let rewound = self
.end
.is_some_and(|previous| start.saturating_add(tolerance) < previous);
if rewound {
self.origin = None;
self.written = 0;
}
let origin = *self.origin.get_or_insert(start);
let expected = (start.saturating_sub(origin).as_secs_f64() * sample_rate as f64).round() as u64;
let hole = expected.saturating_sub(self.written);
let skipped = hole > fill_max;
let silence = if skipped { 0 } else { hole };
let reset_sink = rewound || skipped;
self.written = self
.written
.max(expected)
.saturating_add(u64::try_from(samples).unwrap_or(u64::MAX));
self.end = Some(end);
AudioTiming {
end,
silence,
reset_sink,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
const DELAY: Duration = Duration::from_millis(100);
fn ms(millis: u64) -> Timestamp {
Timestamp::from_millis(millis).unwrap()
}
#[test]
fn a_late_first_frame_catches_up_to_live() {
let start = Instant::now();
let mut presentation = Presentation::new(DELAY);
assert!(presentation.video(ms(0), start), "the first anchor");
assert_eq!(presentation.due(ms(0)), Some(start + DELAY));
let now = start + Duration::from_millis(20);
assert!(presentation.video(ms(40), now), "an early frame left the anchor alone");
assert_eq!(presentation.due(ms(40)), Some(now + DELAY));
let late = now + Duration::from_millis(70);
assert!(!presentation.video(ms(80), late), "a late arrival moved the anchor");
assert_eq!(presentation.due(ms(80)), Some(now + Duration::from_millis(140)));
}
#[test]
fn the_speaker_owns_the_anchor_while_it_plays() {
let start = Instant::now();
let mut presentation = Presentation::new(DELAY);
presentation.audio(Duration::from_secs(1), DELAY, start);
let anchored = presentation.due(ms(1_000));
presentation.video(ms(1_040), start);
presentation.video(ms(4_000), start);
assert_eq!(
presentation.due(ms(1_000)),
anchored,
"video moved the speaker's anchor"
);
presentation.stopped();
presentation.video(ms(4_040), start);
assert_eq!(presentation.due(ms(4_040)), Some(start + DELAY));
}
#[test]
fn the_speaker_anchors_at_the_live_edge() {
let start = Instant::now();
let mut presentation = Presentation::new(DELAY);
presentation.audio(Duration::from_secs(1), DELAY, start);
assert_eq!(presentation.due(ms(1_000)), Some(start + DELAY));
let now = start + Duration::from_millis(500);
presentation.audio(Duration::from_millis(1_500), Duration::from_millis(40), now);
assert_eq!(presentation.due(ms(1_500)), Some(now + Duration::from_millis(40)));
}
#[test]
fn the_speaker_re_anchors_when_it_takes_over() {
let start = Instant::now();
let mut presentation = Presentation::new(DELAY);
presentation.video(ms(4_000), start);
presentation.audio(Duration::from_secs(1), DELAY, start);
assert_eq!(presentation.due(ms(1_000)), Some(start + DELAY));
presentation.stopped();
let now = start + Duration::from_millis(20);
presentation.audio(Duration::from_millis(200), DELAY, now);
assert_eq!(presentation.due(ms(200)), Some(now + DELAY));
}
#[test]
fn a_restarted_speaker_re_anchors() {
let start = Instant::now();
let mut presentation = Presentation::new(DELAY);
presentation.audio(Duration::from_secs(10), DELAY, start);
let now = start + Duration::from_millis(20);
presentation.restarted();
presentation.audio(Duration::from_secs(1), DELAY, now);
assert_eq!(presentation.due(ms(1_000)), Some(now + DELAY));
}
#[test]
fn a_stopped_speaker_drops_a_pending_restart() {
let start = Instant::now();
let mut presentation = Presentation::new(DELAY);
presentation.audio(Duration::from_secs(10), DELAY, start);
presentation.restarted();
presentation.stopped();
presentation.audio(Duration::from_secs(1), DELAY, start);
let now = start + Duration::from_millis(40);
presentation.audio(Duration::from_millis(1_020), DELAY, now);
assert_eq!(
presentation.due(ms(1_020)),
Some(start + Duration::from_millis(20) + DELAY)
);
}
#[test]
fn a_moved_anchor_is_reported() {
let start = Instant::now();
let mut presentation = Presentation::new(DELAY);
assert!(
presentation.audio(Duration::from_secs(1), DELAY, start),
"the first anchor"
);
let now = start + Duration::from_millis(40);
assert!(
!presentation.audio(Duration::from_millis(1_040), DELAY, now),
"a frame on the anchor moved it"
);
assert!(
presentation.audio(Duration::from_millis(1_100), DELAY, now),
"an early frame left the anchor alone"
);
}
#[test]
fn audio_timeline_restarts_when_media_time_rewinds() {
let mut timeline = AudioTimeline::default();
let first = timeline.push(Duration::from_secs(10), 960, 48_000, 24_000);
assert!(!first.reset_sink);
let rewound = timeline.push(Duration::from_secs(5), 960, 48_000, 24_000);
assert!(rewound.reset_sink);
assert_eq!(rewound.silence, 0);
let next = timeline.push(Duration::from_millis(5_020), 960, 48_000, 24_000);
assert!(!next.reset_sink);
assert_eq!(next.silence, 0);
}
#[test]
fn audio_timeline_tolerates_millisecond_stamp_rounding() {
let mut timeline = AudioTimeline::default();
let first = timeline.push(Duration::ZERO, 1024, 44_100, 22_050);
assert!(!first.reset_sink);
let rounded = timeline.push(Duration::from_millis(23), 1024, 44_100, 22_050);
assert!(!rounded.reset_sink);
}
#[test]
fn audio_timeline_resets_sink_when_forward_hole_exceeds_fill_cap() {
let mut timeline = AudioTimeline::default();
timeline.push(Duration::ZERO, 960, 48_000, 4_800);
let filled = timeline.push(Duration::from_millis(100), 960, 48_000, 4_800);
assert!(!filled.reset_sink);
assert_eq!(filled.silence, 3_840);
let skipped = timeline.push(Duration::from_secs(1), 960, 48_000, 4_800);
assert!(skipped.reset_sink);
assert_eq!(skipped.silence, 0);
let next = timeline.push(Duration::from_millis(1_020), 960, 48_000, 4_800);
assert!(!next.reset_sink);
assert_eq!(next.silence, 0);
}
}