moq-cli 0.12.8

Media over QUIC
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
//! Media time: the playout clock, and the audio track's own timeline.

use std::time::{Duration, Instant};

use hang::moq_net::Timestamp;

/// The playout clock: one anchor the window and the speaker both present
/// against.
///
/// Every arriving frame is folded into a single [`Pacer`](moq_mux::Pacer), which
/// maps media time onto the wall clock and holds each result back by the
/// configured delay. That is what pulls playback toward live: a frame arriving
/// earlier than the anchor predicted re-pins it, and the delay is measured from
/// the edge it just set rather than from wherever the first frame happened to
/// land.
///
/// While audio is playing it is the only half allowed to move the anchor. The
/// speaker drains on its own clock and cannot skip forward with a re-anchor, so
/// a video tune-in burst re-anchoring would leave the picture ahead of the
/// sound. Video reads the anchor and follows it.
pub(super) struct Presentation {
	pacer: moq_mux::Pacer,
	/// How far behind the live edge playback runs.
	delay: Duration,
	/// Whether the speaker owns the anchor.
	speaker: bool,
	/// Set when the speaker restarts on a new timeline mid-track, so the next
	/// audio frame pins the anchor outright instead of pacing against media it
	/// will never play.
	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,
		}
	}

	/// Fold an arriving video frame into the anchor, unless the speaker owns it,
	/// reporting whether that moved the schedule.
	///
	/// A move makes every queued frame due earlier, so the window has to recompute
	/// the deadline it is already asleep on.
	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))
	}

	/// Start a new video timeline without taking ownership from a speaker.
	pub(super) fn video_restarted(&mut self) {
		if !self.speaker {
			self.pacer = moq_mux::Pacer::default().with_delay(self.delay);
		}
	}

	/// Fold the speaker's position into the anchor, reporting whether that moved
	/// the schedule.
	///
	/// `end` is the media time of the last sample written and `buffered` is how
	/// much of the write is still queued, so `end` does not sound until
	/// `buffered` from now. The anchor is pinned at the instant that leaves
	/// exactly the delay before it does, which is where the live edge sits on a
	/// clock trailing by that delay. Pinning the sample sounding *now* instead
	/// would schedule video a whole delay behind the speaker.
	///
	/// A move makes every queued video frame due earlier, so the window has to
	/// recompute the deadline it is already asleep on.
	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);

		// Taking ownership re-pins outright rather than pacing. Whatever set the
		// anchor before is either video, which the speaker cannot follow, or a
		// track that has since ended, and a timestamp that maps into the past under
		// that anchor would never re-anchor on its own: the sink would play now
		// while the picture stayed on a timeline nothing is feeding.
		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)
	}

	/// The speaker restarted on a new timeline, so the media either side of the
	/// break is unrelated: re-pin the anchor on the next audio frame rather than
	/// pacing across a jump the speaker never plays.
	pub(super) fn restarted(&mut self) {
		self.restart = true;
	}

	/// The speaker went quiet, so nothing holds playback to its cadence any
	/// more and video takes the anchor back.
	///
	/// A replacement rendition takes ownership again on its first frame, and
	/// re-pins there: it is a track boundary, so its timestamps are not
	/// necessarily continuous with what just ended. That covers a pending
	/// restart, which would otherwise re-pin the frame after it too.
	pub(super) fn stopped(&mut self) {
		self.speaker = false;
		self.restart = false;
	}

	/// When `timestamp` should be presented, or `None` before any frame has
	/// anchored the clock. `None` means as soon as possible: there is nothing
	/// left to wait for.
	pub(super) fn due(&self, timestamp: Timestamp) -> Option<Instant> {
		self.pacer.due(timestamp)
	}
}

/// A wire timestamp as a duration from the start of the track.
pub(super) fn timestamp(timestamp: Timestamp) -> Duration {
	Duration::from_micros(timestamp.as_micros().min(u64::MAX as u128) as u64)
}

/// A duration from the start of the track as a wire timestamp, the inverse of
/// [`timestamp`].
fn media(duration: Duration) -> Timestamp {
	// A wire timestamp caps at 2^62 - 1, which is ~146,000 years of microseconds.
	const MAX: u128 = (1 << 62) - 1;
	Timestamp::from_micros(duration.as_micros().min(MAX) as u64).expect("clamped to the wire maximum")
}

/// Where the audio track has reached, measured from its own origin so
/// timestamp rounding can't accumulate into drift.
#[derive(Default)]
pub(super) struct AudioTimeline {
	origin: Option<Duration>,
	end: Option<Duration>,
	written: u64,
}

/// What the speaker owes before the frame just pushed: silence to play a hole
/// through, or a fresh sink when the timeline jumped too far to fill.
pub(super) struct AudioTiming {
	/// Media time the pushed frame ends at.
	pub(super) end: Duration,
	/// Samples of silence to write first.
	pub(super) silence: u64,
	/// Whether the buffered sink has to be replaced.
	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);
		// Millisecond-stamped input can put adjacent frames on either side of their
		// exact boundary. Two output samples cover the conversions on top of that.
		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;
		}

		// Measure every hole from the track origin so timestamp rounding cannot
		// accumulate into drift. Advancing to `expected` even when the hole is skipped
		// keeps the next frame contiguous with the new timeline position.
		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()
	}

	/// Earlier arrivals pull the anchor forward so playback converges on the live
	/// edge without giving up the configured delay.
	#[test]
	fn a_late_first_frame_catches_up_to_live() {
		let start = Instant::now();
		let mut presentation = Presentation::new(DELAY);

		// Produced 500ms ago, though nothing here knows that yet. The window is
		// asleep on the deadline the anchor last gave it, so every move is reported.
		assert!(presentation.video(ms(0), start), "the first anchor");
		assert_eq!(presentation.due(ms(0)), Some(start + DELAY));

		// 40ms of media later, but only 20ms of wall clock: the anchor was 480ms
		// behind live, so it moves onto this frame.
		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));

		// A frame that merely arrives late keeps its media instant. The delay is
		// the room it has to be late in.
		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)));
	}

	/// The speaker cannot skip forward with a re-anchor, so while it is playing a
	/// video burst must not move the anchor out from under it.
	#[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));

		// A tune-in burst: seconds of video arriving at once.
		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"
		);

		// Once audio stops, video anchors again.
		presentation.stopped();
		presentation.video(ms(4_040), start);
		assert_eq!(presentation.due(ms(4_040)), Some(start + DELAY));
	}

	/// The speaker reports the sample sounding now, which is a delay behind the
	/// edge. Pinning that instant directly would schedule video a second delay
	/// behind the sound.
	#[test]
	fn the_speaker_anchors_at_the_live_edge() {
		let start = Instant::now();
		let mut presentation = Presentation::new(DELAY);

		// The last sample written is a full delay from sounding, which is exactly
		// where a settled sink sits: media time and the picture agree.
		presentation.audio(Duration::from_secs(1), DELAY, start);
		assert_eq!(presentation.due(ms(1_000)), Some(start + DELAY));

		// Filling up, with only 40ms queued: that sample sounds in 40ms, so the
		// picture that goes with it is due then too.
		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)));
	}

	/// Video usually anchors first, and audio's first frame can sit behind that
	/// anchor. Pacing it would map it into the past, which never re-anchors, so
	/// the sink would play now while the picture stayed on video's timeline.
	#[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));

		// A replacement rendition is a track boundary too: its timestamps need not
		// continue the one that just ended.
		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));
	}

	/// A hole too large to play through drops the buffered audio and starts a new
	/// sink, and a rewind restarts the timeline outright. Pacing across either
	/// would schedule video against samples the speaker never plays.
	#[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);

		// The publisher rewound: without the restart this pins media 9 seconds
		// behind the anchor, leaving the picture 9 seconds in the past.
		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));
	}

	/// A retired rendition marks a restart for its replacement, which may never
	/// come. The next track re-pins on its first frame for taking ownership, so a
	/// restart left over would re-pin the frame after it onto that frame's jitter.
	#[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);
		// 20ms of media, 40ms of wall clock: a late frame paces and leaves the
		// anchor where it was.
		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)
		);
	}

	/// The window sleeps on the deadline it last computed, so an anchor the
	/// speaker moved has to say so or the queue presents late by that much.
	#[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"
		);

		// On schedule: 40ms of media, 40ms of wall clock, same sink depth.
		let now = start + Duration::from_millis(40);
		assert!(
			!presentation.audio(Duration::from_millis(1_040), DELAY, now),
			"a frame on the anchor moved it"
		);

		// Arriving early re-anchors, and the window has to hear about 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);

		// 1024 frames end at 23.22 ms, but an FLV timestamp carries 23 ms.
		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);
	}
}