Skip to main content

moq_audio/decode/
consumer.rs

1//! Subscribe to an encoded audio track and emit raw PCM.
2
3use std::collections::VecDeque;
4
5use bytes::Bytes;
6
7use super::decoder::{Config, Decoder};
8use crate::resample::{Resampler, remix, validate_remix};
9use crate::{Activity, Error, Format, Frame, Layout};
10
11/// Where a consumer starts on a track that already holds groups.
12#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
13#[non_exhaustive]
14pub enum Start {
15	/// Decode every cached group.
16	#[default]
17	Oldest,
18	/// Start at the newest cached group.
19	Latest,
20}
21
22/// PCM output conversion requested from [`Consumer`].
23#[derive(Clone, Debug, Default)]
24#[non_exhaustive]
25pub struct Output {
26	/// How to pack samples in each emitted frame.
27	pub format: Format,
28	/// Output sample rate, or the codec rate when absent.
29	pub sample_rate: Option<u32>,
30	/// Output layout, or the codec layout when absent.
31	pub layout: Option<Layout>,
32}
33
34/// Subscription, decoder, and output policy for [`Consumer`].
35#[derive(Clone, Debug, Default)]
36#[non_exhaustive]
37pub struct Options {
38	/// Low-level decoder configuration.
39	pub decoder: Config,
40	/// PCM output conversion.
41	pub output: Output,
42	/// Maximum media age accepted from the subscription.
43	pub max_age: std::time::Duration,
44	/// Initial cached-group policy.
45	pub start: Start,
46}
47
48impl Options {
49	/// Build default real-time consumer options.
50	pub fn new() -> Self {
51		Self::default()
52	}
53}
54
55/// Subscribe to a moq-mux audio track and emit decoded PCM in the requested
56/// [`Output`].
57///
58/// The mirror of [`encode::Producer`](crate::encode::Producer): output format /
59/// sample rate / layout are fixed at construction, and
60/// [`read`](Self::read) returns [`Frame`]s carrying the codec activity they
61/// were decoded from.
62pub struct Consumer {
63	decoder: Decoder,
64	track: moq_mux::container::Consumer<moq_mux::catalog::hang::Container>,
65	resampler: Option<Resampler>,
66	options: Options,
67	max_age: std::time::Duration,
68	resolved_sample_rate: u32,
69	resolved_layout: Layout,
70	/// Where the next packet's timestamp should land: the last packet's timestamp
71	/// plus the media it covered, including the codec delay the decoder trimmed off
72	/// the front. A packet that misses it is a hole nobody declared.
73	next_start: Option<moq_net::Timestamp>,
74	/// Frames decoded and not yet handed back, so a gap's tail can be returned
75	/// ahead of the packet that exposed it.
76	ready: VecDeque<Frame>,
77	/// Codec activity spans the resampler's buffered output still covers.
78	spans: VecDeque<ActivitySpan>,
79	/// Activity of the last span the output ran past, for the rounding samples the
80	/// filter leaves beyond the final input boundary.
81	trailing: Activity,
82	/// Timestamp of the first encoded packet in this decoder epoch, used to
83	/// interpret codec delay and a terminal marker.
84	epoch: Option<moq_net::Timestamp>,
85	/// Codec delay trimmed since the current decoder epoch began.
86	delay_trimmed: usize,
87	/// Codec-rate terminal frames emitted since `terminal_start`.
88	frames_decoded: usize,
89	/// Logical endpoint carried by an empty legacy frame before terminal packets.
90	end: Option<moq_net::Timestamp>,
91	/// Presentation time of the first decoded terminal frame.
92	terminal_start: Option<moq_net::Timestamp>,
93	/// Last container playhead generation applied to timeline state.
94	discontinuity: u64,
95}
96
97struct ActivitySpan {
98	end: moq_net::Timestamp,
99	activity: Activity,
100}
101
102impl Consumer {
103	/// Subscribe to `name` in `broadcast`, using the catalog entry to pick the
104	/// codec.
105	pub async fn new(
106		broadcast: &moq_net::broadcast::Consumer,
107		catalog: &hang::catalog::AudioConfig,
108		name: impl Into<String>,
109		options: Options,
110	) -> Result<Self, Error> {
111		let decoder = Decoder::new(catalog, &options.decoder)?;
112		let sample_rate = options.output.sample_rate.unwrap_or_else(|| decoder.sample_rate());
113		let layout = options.output.layout.unwrap_or_else(|| decoder.layout());
114		validate_remix(decoder.layout(), layout)?;
115
116		let resampler = if sample_rate == decoder.sample_rate() {
117			None
118		} else {
119			let chunk_frames = (decoder.sample_rate() as usize * 20) / 1000;
120			Some(Resampler::new(
121				decoder.sample_rate(),
122				sample_rate,
123				decoder.layout().channels(),
124				chunk_frames,
125			)?)
126		};
127
128		let name = name.into();
129		let track = broadcast.track(&name)?;
130		let mut subscriber = track
131			.subscribe(
132				moq_net::track::Subscription::default()
133					.with_priority(hang::catalog::PRIORITY.audio)
134					.with_max_age(options.max_age),
135			)
136			.await?;
137		// A decoder often opens on a track that is already cached: a replacement
138		// decoder subscribes while its predecessor still holds groups, and a
139		// rendition switched away from and back to stays warm on the origin for
140		// `TRACK_IDLE_LINGER` (cached groups, not an upstream subscription). A
141		// caller that asked for `Start::Latest` wants none of that backlog,
142		// because a cursor starting at sequence zero replays every cached group
143		// at decode speed before reaching live media, which on a thirty-second
144		// retention is half a minute of sound raced through.
145		//
146		// This moves the local read cursor and deliberately not
147		// `Subscription::group_start`. That field is a request to the publisher,
148		// aggregated across every live subscriber, so naming a stale cached
149		// sequence there asks the publisher to rewind the track for everyone
150		// reading it. What a player wants is to skip what it already has.
151		if options.start == Start::Latest
152			&& let Some(live_edge) = track.latest()
153		{
154			subscriber.set_groups(live_edge..);
155		}
156		let track = subscriber;
157		let max_age = options.max_age.min(track.info().max_age);
158		// The catalog says how the track is framed, and it is not always the legacy
159		// wire: `moq import fmp4` publishes CMAF. Reading a moof+mdat fragment as a
160		// varint timestamp plus a payload decodes to garbage rather than failing.
161		let container = moq_mux::catalog::hang::Container::try_from(catalog)?;
162		let track = moq_mux::container::Consumer::new(track, container);
163
164		Ok(Self {
165			decoder,
166			track,
167			resampler,
168			options,
169			max_age,
170			resolved_sample_rate: sample_rate,
171			resolved_layout: layout,
172			next_start: None,
173			ready: VecDeque::new(),
174			spans: VecDeque::new(),
175			trailing: Activity::Active,
176			epoch: None,
177			delay_trimmed: 0,
178			frames_decoded: 0,
179			end: None,
180			terminal_start: None,
181			discontinuity: 0,
182		})
183	}
184
185	/// The options this consumer was built with.
186	pub fn options(&self) -> &Options {
187		&self.options
188	}
189
190	/// The effective age budget after clamping to the publisher's retention window.
191	pub fn max_age(&self) -> std::time::Duration {
192		self.max_age
193	}
194
195	/// Sample rate samples are actually delivered at, which is
196	/// [`Output::sample_rate`] resolved against the catalog.
197	pub fn sample_rate(&self) -> u32 {
198		self.resolved_sample_rate
199	}
200
201	/// Layout samples are actually delivered in.
202	pub fn layout(&self) -> Layout {
203		self.resolved_layout
204	}
205
206	/// Read the next decoded PCM frame, or `None` when the track ends.
207	///
208	/// [`Frame::activity`] reports whether the packet these samples came from
209	/// coded audio. It describes where the frame begins, so a resampled
210	/// frame that straddles a change carries the activity its first sample came
211	/// from and the next frame carries the new one.
212	///
213	/// A timestamp that doesn't continue the previous packet is a hole in the
214	/// output, not a splice: nothing is carried across it, and the frames on either
215	/// side stay anchored to their own packet timeline, so the hole is there to
216	/// see. "Doesn't continue" allows for the quantization the stamps carry, which
217	/// on a millisecond-stamped ingest is most of a millisecond.
218	pub async fn read(&mut self) -> Result<Option<Frame>, Error> {
219		loop {
220			if let Some(frame) = self.ready.pop_front() {
221				return Ok(Some(frame));
222			}
223
224			let mux_frame = self.track.read().await?;
225			self.apply_discontinuity()?;
226			let Some(mux_frame) = mux_frame else {
227				return self.flush();
228			};
229
230			if let Some(end) = self.track.end()
231				&& self.end != Some(end)
232			{
233				self.end = Some(end);
234				self.frames_decoded = 0;
235				self.terminal_start = None;
236			}
237
238			// Undeclared holes are routine: a skipped stalled group, a packet the
239			// decoder refused, an ingest that resynced. Drop every stage's state at
240			// the edge, before the packet after it goes anywhere near the decoder.
241			//
242			// Skipped once an end marker arrives, because from there the terminal
243			// phase reconstructs each batch's time from the marker rather than
244			// reading it off the packet, so there is nothing left to compare.
245			if self.end.is_none()
246				&& self
247					.next_start
248					.is_some_and(|next| discontinuous(next, mux_frame.timestamp))
249				&& let Some(frame) = self.gap()?
250			{
251				self.ready.push_back(frame);
252			}
253
254			let rate = self.decoder.sample_rate();
255			let epoch = *self.epoch.get_or_insert(mux_frame.timestamp);
256			let delay = self.decoder.delay_remaining();
257			let decoded = self.decoder.decode(&mux_frame.payload)?;
258			// Codec delay trimmed off the front is media this packet covered even
259			// though no samples came out, so it still moves the packet after it along.
260			let trimmed = delay - self.decoder.delay_remaining();
261			self.delay_trimmed += trimmed;
262			let activity = decoded.activity;
263			let mut decoded = decoded.samples;
264			if let Some(end) = self.end {
265				let terminal_start = *self
266					.terminal_start
267					.get_or_insert(rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch));
268				let total = frames_between(terminal_start, end, rate)?;
269				let remaining = total.saturating_sub(self.frames_decoded);
270				decoded.truncate(remaining.saturating_mul(self.decoder.layout().channels() as usize));
271			}
272
273			let frames = decoded.len() / self.decoder.layout().channels() as usize;
274			let decoded_at = if let Some(terminal_start) = self.terminal_start {
275				advance(terminal_start, self.frames_decoded, rate)?
276			} else {
277				// The codec delay is padding before the epoch, not a hole after the
278				// first short frame. Keep later output contiguous by moving it back over
279				// everything trimmed since this decoder epoch began.
280				rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch)
281			};
282			if self.end.is_some() {
283				self.frames_decoded += frames;
284			}
285			// Packet continuity stays on the encoded timeline. `decoded_at` may be
286			// earlier because codec pre-skip is padding before the decoded epoch.
287			self.next_start = Some(advance(mux_frame.timestamp, frames + trimmed, rate)?);
288			if decoded.is_empty() {
289				continue;
290			}
291
292			let (pcm, timestamp) = match self.resampler.as_mut() {
293				// The resampler works in fixed chunks, so it holds back whatever didn't
294				// fill one. What comes out next starts with those held-back samples, which
295				// arrived before this packet did, so it is stamped where they arrived.
296				// Reading that off this packet instead would place the audio late by up to
297				// a chunk, sawtoothing A/V sync, and drag it the whole way whenever the
298				// source jumps forward without declaring a hole.
299				Some(r) => {
300					let held = if r.pending_frames() == 0 {
301						decoded_at
302					} else {
303						r.held_at().unwrap_or(decoded_at)
304					};
305					let skipped = r.skipped();
306					let pcm = r.process(&decoded, decoded_at)?;
307					(pcm, rewind(held, skipped, self.resolved_sample_rate)?)
308				}
309				None => (decoded, decoded_at),
310			};
311
312			let decoded_end = advance(decoded_at, frames, rate)?;
313
314			// The resampler hands back samples it was holding from earlier packets,
315			// so what comes out starts before the packet that filled its chunk. Track
316			// where each packet's activity ends so the output can be labelled by
317			// where it actually begins, not by the packet just submitted.
318			let resampled = self.resampler.is_some();
319			if resampled {
320				self.spans.push_back(ActivitySpan {
321					end: decoded_end,
322					activity,
323				});
324			}
325
326			// A packet shorter than the resampler's chunk leaves nothing to hand
327			// over yet. Read on rather than returning a frame with no samples, which
328			// a caller would otherwise see as audio arriving.
329			if pcm.is_empty() {
330				continue;
331			}
332
333			let activity = if resampled {
334				self.activity_at(timestamp)
335			} else {
336				activity
337			};
338			// Queued rather than returned, so a tail drained at a gap earlier in this
339			// same iteration still comes out first. The next turn of the loop pops it.
340			let frame = self.frame(pcm, timestamp, activity)?;
341			self.ready.push_back(frame);
342		}
343	}
344
345	/// A playhead event re-applies startup delay and skip. The decoder is not reset:
346	/// the next group already starts on a keyframe, and pre-skip is a play-path concern.
347	fn apply_discontinuity(&mut self) -> Result<(), Error> {
348		let discontinuity = self.track.discontinuity();
349		if discontinuity == self.discontinuity {
350			return Ok(());
351		}
352
353		self.discontinuity = discontinuity;
354		self.next_start = None;
355		self.spans.clear();
356		self.trailing = Activity::Active;
357		self.frames_decoded = 0;
358		self.end = None;
359		self.terminal_start = None;
360		self.epoch = None;
361		self.delay_trimmed = 0;
362		self.decoder.reapply_delay();
363		Ok(())
364	}
365
366	/// Reset codec prediction and resampling state at a hole, returning whatever the
367	/// resampler was still holding from before it.
368	///
369	/// Those samples arrived before the hole and belong before it, so they come
370	/// out as their own frame rather than being filtered together with the audio
371	/// on the far side. The resampler starts over from there, which is what makes
372	/// the next packet's output stamp from the packet itself: nothing is buffered
373	/// to reach back over.
374	fn gap(&mut self) -> Result<Option<Frame>, Error> {
375		self.decoder.reset_prediction()?;
376
377		let mut frame = None;
378		if let Some(resampler) = self.resampler.as_mut() {
379			let held = resampler.held_at();
380			let skipped = resampler.skipped();
381			let pcm = resampler.drain()?;
382			frame = self.tail(pcm, held, skipped)?;
383		}
384
385		self.next_start = None;
386		self.spans.clear();
387		self.trailing = Activity::Active;
388		self.epoch = None;
389		self.delay_trimmed = 0;
390		Ok(frame)
391	}
392
393	/// The tail the resampler is still holding when the track ends, once.
394	///
395	/// Without it the last partial chunk is dropped, which is up to a chunk of
396	/// audio missing from the end of every resampled track. Flushing consumes the
397	/// resampler, which is what makes calling this on every later poll return
398	/// `None` rather than more tails.
399	fn flush(&mut self) -> Result<Option<Frame>, Error> {
400		let Some(resampler) = self.resampler.take() else {
401			return Ok(None);
402		};
403
404		let held = resampler.held_at();
405		let skipped = resampler.skipped();
406		self.tail(resampler.flush()?, held, skipped)
407	}
408
409	/// Stamp and pack a tail the resampler handed back, if it handed back one.
410	///
411	/// `held` is where the samples it was holding arrived, which is where the tail
412	/// begins once the startup frames it dropped of its own are taken off. `None`
413	/// there means the resampler never ran, so there is nothing to place.
414	fn tail(
415		&mut self,
416		pcm: Vec<f32>,
417		held: Option<moq_net::Timestamp>,
418		skipped: usize,
419	) -> Result<Option<Frame>, Error> {
420		let Some(held) = held.filter(|_| !pcm.is_empty()) else {
421			return Ok(None);
422		};
423
424		let timestamp = rewind(held, skipped, self.resolved_sample_rate)?;
425		let activity = self.activity_at(timestamp);
426		Ok(Some(self.frame(pcm, timestamp, activity)?))
427	}
428
429	/// The codec activity covering `timestamp`, dropping the spans it has passed.
430	fn activity_at(&mut self, timestamp: moq_net::Timestamp) -> Activity {
431		while let Some(span) = self.spans.front().filter(|span| span.end <= timestamp) {
432			self.trailing = span.activity;
433			self.spans.pop_front();
434		}
435
436		self.spans.front().map_or(self.trailing, |span| span.activity)
437	}
438
439	/// Remix and pack decoded PCM into an output frame.
440	fn frame(&self, pcm: Vec<f32>, timestamp: moq_net::Timestamp, activity: Activity) -> Result<Frame, Error> {
441		let pcm = if self.decoder.layout() == self.resolved_layout {
442			pcm
443		} else {
444			remix(&pcm, self.decoder.layout(), self.resolved_layout)?
445		};
446
447		let bytes = self
448			.options
449			.output
450			.format
451			.from_interleaved_f32(&pcm, self.resolved_layout.channels())?;
452		Ok(Frame {
453			timestamp,
454			data: Bytes::from(bytes),
455			activity,
456		})
457	}
458}
459
460/// Whether `timestamp` fails to continue `expected`, leaving a hole (or an
461/// overlap) rather than the next packet in line.
462///
463/// Exact contiguity cannot be the test. RTMP stamps in whole milliseconds while a
464/// 1024-sample AAC frame at 44.1 kHz runs 23.22 ms, so on the most common ingest
465/// path every packet lands beside where its predecessor ended.
466///
467/// The slack is the quantization the stamps carry, and nothing else. A frame
468/// duration would be far too much: a single lost packet lands exactly one frame
469/// off, and Opus packets run anywhere from 2.5 ms to 60 ms with no duration
470/// declared in the catalog, so a half-frame rule read off a 20 ms neighbour would
471/// splice straight across a lost 2.5 ms one.
472///
473/// So a packet is discontinuous when it misses `expected` by more than one unit of
474/// the coarsest timescale on the path, plus one unit of the stamp's own scale for
475/// the rounding in the arithmetic that produced `expected`. The coarsest timescale
476/// is the stamp's own scale floored at [`Timescale::default`](moq_net::Timescale):
477/// the legacy hang container re-stamps every frame in microseconds whatever the
478/// source used, and a wire that cannot carry a timescale at all (moq-lite before
479/// 05, IETF moq-transport) falls back to milliseconds, so a millisecond is the
480/// finest quantization a packet can be assumed to have kept. That floor stays under
481/// the shortest packet anything here can send, 2.5 ms of Opus, so it never
482/// swallows a lost one.
483fn discontinuous(expected: moq_net::Timestamp, timestamp: moq_net::Timestamp) -> bool {
484	let scale = expected.scale().max(timestamp.scale());
485	let quantum = scale.min(moq_net::Timescale::default());
486	let tolerance = (scale.as_u64() as u128).div_ceil(quantum.as_u64() as u128) + 1;
487	expected.as_scale(scale).abs_diff(timestamp.as_scale(scale)) > tolerance
488}
489
490/// `timestamp` moved forward by `frames` at `sample_rate`, in its own timescale.
491fn advance(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
492	if frames == 0 {
493		return Ok(timestamp);
494	}
495
496	let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
497	Ok(timestamp.checked_add(offset)?)
498}
499
500/// Codec-rate frames in the interval, rounding a microsecond marker to the nearest frame.
501fn frames_between(start: moq_net::Timestamp, end: moq_net::Timestamp, sample_rate: u32) -> Result<usize, Error> {
502	let duration = end.checked_sub(start)?;
503	let frames = (std::time::Duration::from(duration).as_nanos() * sample_rate as u128 + 500_000_000) / 1_000_000_000;
504	usize::try_from(frames).map_err(|_| Error::Unsupported("audio duration does not fit in memory".into()))
505}
506
507/// `timestamp` moved back by `frames` at `sample_rate`, in its own timescale.
508///
509/// Saturates at zero rather than failing: a publisher whose first timestamps
510/// don't advance is odd, but it isn't a reason to end the track.
511fn rewind(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
512	if frames == 0 {
513		return Ok(timestamp);
514	}
515
516	let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
517	Ok(timestamp
518		.checked_sub(offset)
519		.unwrap_or(moq_net::Timestamp::new(0, timestamp.scale())?))
520}
521
522#[cfg(test)]
523mod tests {
524	use moq_net::Timestamp;
525
526	use super::*;
527	use crate::encode::{Encoder, Input, Options as EncodeOptions, Producer, Settings};
528	use crate::{Format, Layout};
529
530	#[tokio::test]
531	async fn remixes_mono_stream_to_stereo_output() {
532		let mut broadcast = moq_net::broadcast::Info::new().produce();
533		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
534		let subscriber = broadcast.consume();
535		let input = Input {
536			format: Format::F32,
537			sample_rate: 48_000,
538			layout: Layout::Mono,
539		};
540		let options = EncodeOptions {
541			track: Some("audio".to_string()),
542			settings: Settings::new(48_000, Layout::Mono),
543			..EncodeOptions::default()
544		};
545		let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
546		let catalog = Encoder::new(&Settings::new(input.sample_rate, input.layout))
547			.unwrap()
548			.catalog();
549		let mut consumer = Consumer::new(
550			&subscriber,
551			&catalog,
552			"audio",
553			Options {
554				output: Output {
555					layout: Some(Layout::Stereo),
556					..Output::default()
557				},
558				..Options::new()
559			},
560		)
561		.await
562		.unwrap();
563
564		let samples = vec![0.1f32; 960];
565		let mut data = Vec::with_capacity(samples.len() * size_of::<f32>());
566		for sample in samples {
567			data.extend_from_slice(&sample.to_le_bytes());
568		}
569		producer.write(&Frame::new(data.into(), Timestamp::ZERO)).unwrap();
570
571		let frame = consumer.read().await.unwrap().expect("decoded frame");
572		let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
573		assert_eq!(samples.len(), (960 - 312) * 2);
574		for pair in samples.as_chunks::<2>().0.iter() {
575			assert_eq!(pair[0], pair[1]);
576		}
577	}
578
579	/// An imported 44.1 kHz Opus stream decodes on the 48 kHz clock: the pre-skip
580	/// is trimmed once as padding before the first packet, and every later frame
581	/// is stamped where its samples fall.
582	#[tokio::test]
583	async fn opus_timestamps_follow_the_48k_clock() {
584		use crate::decode::decoder::tests::{opus_catalog, opus_packets};
585
586		let broadcast = moq_net::broadcast::Info::new().produce();
587		let track = broadcast
588			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
589			.unwrap();
590		let subscriber = broadcast.consume();
591
592		let catalog = opus_catalog(moq_mux::codec::opus::Config::new(44_100, 1).with_pre_skip(312));
593		let mut producer = moq_mux::container::Producer::new(
594			track,
595			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
596		);
597		let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Options::new())
598			.await
599			.unwrap();
600		assert_eq!(consumer.sample_rate(), 48_000);
601
602		for (packet, payload) in opus_packets(3).into_iter().enumerate() {
603			producer
604				.write(moq_mux::container::Frame {
605					timestamp: Timestamp::from_micros(packet as u64 * 20_000).unwrap(),
606					duration: None,
607					payload,
608					keyframe: packet == 0,
609				})
610				.unwrap();
611		}
612
613		// 312 samples at 48 kHz is 6.5 ms.
614		for (micros, frames) in [(0, 960 - 312), (13_500, 960), (33_500, 960)] {
615			let frame = consumer.read().await.unwrap().expect("decoded frame");
616			assert_eq!(frame.timestamp.as_micros(), micros);
617			assert_eq!(frame.data.len() / size_of::<f32>(), frames);
618		}
619	}
620
621	/// A packet whose sample count isn't a multiple of the resampler's chunk leaves
622	/// samples buffered, and the next output starts with those. Stamping that
623	/// output with the packet that completed the chunk puts it up to a chunk late,
624	/// which is a sawtooth in A/V sync rather than a constant offset. Any codec
625	/// whose frame is not a whole number of chunks reaches it: a 1024-sample frame
626	/// at 44.1 kHz never fills the 882-frame chunk evenly.
627	#[tokio::test]
628	async fn resampled_timestamps_follow_the_samples() {
629		let broadcast = moq_net::broadcast::Info::new().produce();
630		let track = broadcast
631			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
632			.unwrap();
633		let subscriber = broadcast.consume();
634
635		let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
636		let mut producer = moq_mux::container::Producer::new(
637			track,
638			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
639		);
640
641		let mut consumer = Consumer::new(
642			&subscriber,
643			&catalog,
644			"audio",
645			Options {
646				output: Output {
647					sample_rate: Some(48_000),
648					..Output::default()
649				},
650				max_age: std::time::Duration::from_secs(1),
651				..Options::new()
652			},
653		)
654		.await
655		.unwrap();
656
657		// Two 1024-sample packets, back to back at the codec's own rate.
658		const FRAMES: u64 = 1024;
659		let payload: Bytes = vec![0u8; FRAMES as usize * size_of::<f32>()].into();
660		for packet in 0..2 {
661			producer
662				.write(moq_mux::container::Frame {
663					timestamp: moq_net::Timestamp::from_scale(packet * FRAMES, 44_100).unwrap(),
664					duration: None,
665					payload: payload.clone(),
666					keyframe: true,
667				})
668				.unwrap();
669		}
670
671		let first = consumer.read().await.unwrap().expect("decoded frame");
672		assert_eq!(first.timestamp.as_micros(), 0);
673
674		// Continuity, not a fixed number: the second frame starts where the first
675		// one's samples end, whatever they came to. Within a few frames rather than
676		// exactly, because the resampler emits whole frames and its count per chunk
677		// wobbles around the nominal ratio; a real hole (the samples it held back, or
678		// the startup silence it dropped) is twenty times this tolerance.
679		let second = consumer.read().await.unwrap().expect("decoded frame");
680		let first_frames = (first.data.len() / size_of::<f32>()) as u128;
681		let ends_at = first_frames * 1_000_000 / 48_000;
682		let gap = second.timestamp.as_micros().abs_diff(ends_at);
683		assert!(gap < 100, "expected the frames to meet, got a {gap} us gap");
684	}
685
686	/// The resampler only converts whole chunks, so the last partial one has to be
687	/// flushed at end of track or its audio is simply gone. A 1024-sample frame at
688	/// 44.1 kHz guarantees a remainder, never filling the 882-frame chunk evenly.
689	#[tokio::test]
690	async fn resampled_tail_survives_the_end_of_the_track() {
691		let broadcast = moq_net::broadcast::Info::new().produce();
692		let track = broadcast
693			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
694			.unwrap();
695		let subscriber = broadcast.consume();
696
697		let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
698		let mut producer = moq_mux::container::Producer::new(
699			track,
700			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
701		);
702
703		let mut consumer = Consumer::new(
704			&subscriber,
705			&catalog,
706			"audio",
707			Options {
708				output: Output {
709					sample_rate: Some(48_000),
710					..Output::default()
711				},
712				..Options::new()
713			},
714		)
715		.await
716		.unwrap();
717
718		// One 1024-frame packet: 882 fill a chunk, 142 are left holding.
719		const FRAMES: usize = 1024;
720		let payload: Bytes = vec![0u8; FRAMES * size_of::<f32>()].into();
721		producer
722			.write(moq_mux::container::Frame {
723				timestamp: moq_net::Timestamp::ZERO,
724				duration: None,
725				payload,
726				keyframe: true,
727			})
728			.unwrap();
729		producer.finish().unwrap();
730
731		let first = consumer.read().await.unwrap().expect("decoded frame");
732		let first_frames = first.data.len() / size_of::<f32>();
733
734		let tail = consumer.read().await.unwrap().expect("flushed tail");
735		let tail_frames = tail.data.len() / size_of::<f32>();
736
737		// The 142 held-back frames at 44.1 kHz are ~155 at 48 kHz, plus the 69 the
738		// sinc filter still owes: it runs centred, so the end of the track only
739		// emerges once the flush has fed it silence to push it out.
740		assert!((215..=230).contains(&tail_frames), "unexpected tail: {tail_frames}");
741		// It picks up where the first frame's samples ended, within the same few
742		// frames of whole-frame rounding as above.
743		let ends_at = (first_frames as u128) * 1_000_000 / 48_000;
744		let gap = tail.timestamp.as_micros().abs_diff(ends_at);
745		assert!(gap < 100, "expected the tail to meet the body, got a {gap} us gap");
746
747		// Together they cover the packet and no more: 1024 frames at 44.1 kHz is
748		// ~1114 at 48 kHz. The filter's delay does not extend the stream, because
749		// what the drain adds here is what the start dropped off the front.
750		let total = first_frames + tail_frames;
751		assert!((1105..=1120).contains(&total), "unexpected total: {total}");
752		assert!(consumer.read().await.unwrap().is_none());
753	}
754
755	#[tokio::test]
756	async fn resampling_keeps_the_activity_boundary_on_its_source() {
757		let mut encoder = Encoder::new(&Settings {
758			dtx: true,
759			bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
760			frame_duration: std::time::Duration::from_millis(10),
761			..Settings::new(48_000, Layout::Mono)
762		})
763		.unwrap();
764		let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Opus, 48_000, 1);
765
766		let broadcast = moq_net::broadcast::Info::new().produce();
767		let track = broadcast
768			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
769			.unwrap();
770		let subscriber = broadcast.consume();
771		let mut producer = moq_mux::container::Producer::new(
772			track,
773			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
774		);
775		let mut consumer = Consumer::new(
776			&subscriber,
777			&catalog,
778			"audio",
779			Options {
780				output: Output {
781					sample_rate: Some(44_100),
782					..Output::default()
783				},
784				max_age: std::time::Duration::from_secs(1),
785				..Options::new()
786			},
787		)
788		.await
789		.unwrap();
790
791		let active = vec![0.5; encoder.frame_size()];
792		let silence = vec![0.0; encoder.frame_size()];
793		let mut first_dtx = None;
794		for index in 0..40u64 {
795			let packet = encoder.encode(if index == 0 { &active } else { &silence }).unwrap();
796			let timestamp = Timestamp::from_scale(index * encoder.frame_size() as u64, 48_000).unwrap();
797			if first_dtx.is_none() && packet.activity.is_dtx() {
798				first_dtx = Some(timestamp);
799			}
800			producer
801				.write(moq_mux::container::Frame {
802					timestamp,
803					payload: packet.payload,
804					keyframe: true,
805					duration: None,
806				})
807				.unwrap();
808			producer.cut(None).unwrap();
809		}
810		producer.finish().unwrap();
811
812		let expected = first_dtx.expect("silence should enter Opus DTX");
813		let mut actual = None;
814		while let Some(frame) = consumer.read().await.unwrap() {
815			// 10 ms packets do not fill the 20 ms chunk, so the resampler hands back
816			// nothing every other packet. Those must not surface as frames: a frame
817			// with no samples reads as audio arriving, and carries an activity
818			// describing samples that are not there.
819			assert!(!frame.data.is_empty(), "read returned a frame with no samples");
820			if frame.activity.is_dtx() {
821				actual = Some(frame.timestamp);
822				break;
823			}
824		}
825		let actual = actual.expect("consumer should report Opus DTX");
826
827		// Each frame carries the activity its first sample came from, so the label
828		// can lag its source by up to the frame it lands in, but it must never lead
829		// it: leading means samples that are still active got labelled DTX. That is
830		// what labelling by the packet most recently submitted does, since the
831		// resampler is handing back audio from before that packet. It puts the
832		// boundary a chunk early instead of a fraction of a chunk late.
833		let delay = actual.as_micros() as i128 - expected.as_micros() as i128;
834		let chunk_us = 20_000i128;
835		assert!(
836			(0..chunk_us).contains(&delay),
837			"DTX label landed {delay} us from its source, outside [0, {chunk_us})"
838		);
839	}
840
841	/// Publish PCM packets of `frames` samples each at the given stamps, and read
842	/// back every decoded frame as `(microseconds, output frames)`.
843	async fn pcm_gaps(rate: u32, out_rate: u32, frames: usize, stamps: &[Timestamp]) -> Vec<(u128, usize)> {
844		let broadcast = moq_net::broadcast::Info::new().produce();
845		let track = broadcast
846			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
847			.unwrap();
848		let subscriber = broadcast.consume();
849
850		let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, rate, 1);
851		let mut producer = moq_mux::container::Producer::new(
852			track,
853			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
854		);
855		let mut consumer = Consumer::new(
856			&subscriber,
857			&catalog,
858			"audio",
859			// These tests publish a whole track up front and read it back afterwards,
860			// so every packet but the last is already old by the time the consumer
861			// looks. The default budget is zero, which sheds all of them.
862			Options {
863				output: Output {
864					sample_rate: Some(out_rate),
865					..Output::default()
866				},
867				max_age: std::time::Duration::from_secs(1),
868				..Options::new()
869			},
870		)
871		.await
872		.unwrap();
873
874		let payload: Bytes = vec![0u8; frames * size_of::<f32>()].into();
875		for stamp in stamps {
876			producer
877				.write(moq_mux::container::Frame {
878					timestamp: *stamp,
879					duration: None,
880					payload: payload.clone(),
881					keyframe: true,
882				})
883				.unwrap();
884		}
885		producer.finish().unwrap();
886
887		let mut read = Vec::new();
888		while let Some(frame) = consumer.read().await.unwrap() {
889			read.push((frame.timestamp.as_micros(), frame.data.len() / size_of::<f32>()));
890		}
891		read
892	}
893
894	/// A packet that doesn't continue the last one is a hole, not a splice: the
895	/// resampler hands back what it was holding from before the gap as its own
896	/// frame, and the audio after it is stamped from the packet that carried it
897	/// rather than rewound over samples that no longer exist.
898	#[tokio::test]
899	async fn a_missing_packet_leaves_a_hole() {
900		const FRAMES: usize = 1024;
901		// Packets at sample 0 and sample 2048: the one at 1024 never arrived.
902		let stamps = [
903			Timestamp::from_scale(0, 44_100).unwrap(),
904			Timestamp::from_scale(2 * FRAMES as u64, 44_100).unwrap(),
905		];
906		let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
907
908		// The first packet's chunk, then the tail drained at the gap, then the
909		// second packet's chunk. The flush at end of track adds the last tail.
910		assert_eq!(read.len(), 4, "unexpected frames: {read:?}");
911
912		// Everything the first packet carried comes out before the hole: 1024 frames
913		// at 44.1 kHz is ~1114 at 48 kHz, whole-frame rounding aside.
914		let before: usize = read[..2].iter().map(|(_, frames)| frames).sum();
915		assert!((1105..=1120).contains(&before), "unexpected pre-gap audio: {before}");
916
917		// The audio after the hole is stamped by its own packet. Rewinding over the
918		// resampler's buffer instead would put it ~3 ms early, in the middle of the
919		// hole, and splice the two sides together through the filter.
920		assert_eq!(read[2].0, stamps[1].as_micros());
921
922		// And the hole is the packet that never arrived: 1024 frames at 44.1 kHz.
923		let ends_at = read[1].0 + (read[1].1 as u128) * 1_000_000 / 48_000;
924		let hole = read[2].0 - ends_at;
925		assert!((23_100..=23_350).contains(&hole), "unexpected hole: {hole} us");
926	}
927
928	/// A packet can land a hair past where the last one ended without being a hole:
929	/// the stamps are quantized, so `discontinuous` allows a millisecond of slack.
930	/// The jump is still a jump, and the resampler is holding samples from before
931	/// it. Deriving their stamp by counting back from the packet drags them forward
932	/// by the whole jump; reading it off the packet they arrived with does not.
933	#[tokio::test]
934	async fn a_jump_inside_the_slack_leaves_the_held_samples_alone() {
935		// 441 frames at 44.1 kHz is 10 ms, half of the 20 ms chunk, so the first
936		// packet is held whole and the second is what completes the chunk.
937		const FRAMES: usize = 441;
938		// A millisecond past where the first packet ended, which is the slack the
939		// legacy container's microsecond re-stamping is allowed.
940		let stamps = [
941			Timestamp::from_micros(0).unwrap(),
942			Timestamp::from_micros(11_000).unwrap(),
943		];
944		let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
945
946		// The chunk, then the flush at the end of the track: no hole was declared, so
947		// nothing was drained in between.
948		assert_eq!(read.len(), 2, "unexpected frames: {read:?}");
949		// It starts with the first packet's samples, so it is stamped where that
950		// packet was. Rewinding from the second one instead puts it a millisecond late.
951		assert_eq!(read[0].0, 0, "held samples moved with the jump: {read:?}");
952	}
953
954	#[tokio::test]
955	async fn a_jump_after_a_full_chunk_uses_the_new_packet_timestamp() {
956		let stamps = [
957			Timestamp::from_micros(0).unwrap(),
958			Timestamp::from_micros(21_000).unwrap(),
959		];
960		let read = pcm_gaps(44_100, 48_000, 882, &stamps).await;
961		let mut r = crate::resample::Resampler::new(44_100, 48_000, 1, 882).unwrap();
962		r.process(&[0.25; 882], stamps[0]).unwrap();
963		let expected = rewind(stamps[1], r.skipped(), 48_000).unwrap().as_micros();
964		assert_eq!(read[1].0, expected);
965	}
966
967	/// Once an end marker arrives the gap check stops running, because from there
968	/// each batch's time is reconstructed from the marker rather than read off the
969	/// packet. A packet that jumps forward then still moves whatever the resampler
970	/// is holding, and the activity that lands with it: those samples came from
971	/// before the jump and are labelled by the packet they came from.
972	#[tokio::test]
973	async fn a_terminal_jump_leaves_the_held_samples_alone() {
974		let mut encoder = Encoder::new(&Settings {
975			dtx: true,
976			bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
977			..Settings::new(48_000, Layout::Mono)
978		})
979		.unwrap();
980		let catalog = encoder.catalog();
981
982		// One coded packet, then a withheld one to follow it. Taken from the same
983		// encoder rather than published in between, so nothing fills the resampler's
984		// chunk between the two.
985		let active = encoder.encode(&vec![0.5f32; encoder.frame_size()]).unwrap();
986		assert!(active.activity.is_active());
987		let silence = vec![0.0f32; encoder.frame_size()];
988		let dtx = (0..200)
989			.map(|_| encoder.encode(&silence).unwrap())
990			.find(|packet| packet.activity.is_dtx())
991			.expect("silence should enter Opus DTX");
992
993		let broadcast = moq_net::broadcast::Info::new().produce();
994		let track = broadcast
995			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
996			.unwrap();
997		let subscriber = broadcast.consume();
998		let mut producer = moq_mux::container::Producer::new(
999			track,
1000			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1001		);
1002		let mut consumer = Consumer::new(
1003			&subscriber,
1004			&catalog,
1005			"audio",
1006			Options {
1007				output: Output {
1008					sample_rate: Some(44_100),
1009					..Output::default()
1010				},
1011				..Options::new()
1012			},
1013		)
1014		.await
1015		.unwrap();
1016
1017		// A 20 ms Opus packet decodes 960 frames, less the pre-skip on the first one,
1018		// so it doesn't fill the 960-frame chunk and is held whole.
1019		let write = |producer: &mut moq_mux::container::Producer<_>, frames: u64, payload: Bytes, keyframe: bool| {
1020			producer
1021				.write(moq_mux::container::Frame {
1022					timestamp: Timestamp::from_scale(frames, 48_000).unwrap(),
1023					duration: None,
1024					payload,
1025					keyframe,
1026				})
1027				.unwrap();
1028		};
1029		write(&mut producer, 0, active.payload, true);
1030		// The end marker, then the terminal packet a second past where it belongs.
1031		// Same group: a new group at 1s would sit below the marker's live edge.
1032		write(&mut producer, 3 * 48_000, Bytes::new(), false);
1033		write(&mut producer, 48_000, dtx.payload, false);
1034		producer.finish().unwrap();
1035
1036		let frame = consumer.read().await.unwrap().expect("decoded frame");
1037		// The output begins with the first packet's samples, so it is stamped and
1038		// labelled from that packet. Rewinding from the terminal one instead drops it
1039		// most of a second into the future, carrying the DTX label with it.
1040		assert_eq!(frame.timestamp.as_micros(), 0, "held samples moved with the jump");
1041		assert!(
1042			frame.activity.is_active(),
1043			"held samples took the terminal packet's label"
1044		);
1045	}
1046
1047	/// Every packet on the RTMP path lands beside where the last one ended: FLV
1048	/// stamps in whole milliseconds and a 1024-sample AAC frame at 44.1 kHz runs
1049	/// 23.22 ms, so the stamps drift up to a millisecond either way. Reading that as
1050	/// a hole would reset the codec and the resampler on nearly every packet.
1051	///
1052	/// PCM stands in for AAC, which needs an encoder this crate doesn't have: the
1053	/// arithmetic that matters is the packet length and the millisecond stamps.
1054	#[tokio::test]
1055	async fn millisecond_stamps_are_not_a_gap() {
1056		const FRAMES: u64 = 1024;
1057		const PACKETS: u64 = 32;
1058
1059		// What an FLV ingest sends: each packet stamped in whole milliseconds.
1060		let stamps: Vec<_> = (0..PACKETS)
1061			.map(|packet| Timestamp::from_millis(packet * FRAMES * 1_000 / 44_100).unwrap())
1062			.collect();
1063		let read = pcm_gaps(44_100, 48_000, FRAMES as usize, &stamps).await;
1064
1065		// One frame per packet, since 1024 frames always fill at least one 882-frame
1066		// chunk, plus the flush at the end of the track. Reading a gap would drain
1067		// the resampler as well, adding a frame at every packet it fired on.
1068		assert_eq!(read.len(), stamps.len() + 1, "unexpected frames: {read:?}");
1069
1070		// And the output stays continuous across all of them, within the millisecond
1071		// the stamps themselves are quantized to.
1072		for pair in read.windows(2) {
1073			let ends_at = pair[0].0 + (pair[0].1 as u128) * 1_000_000 / 48_000;
1074			assert!(
1075				pair[1].0.abs_diff(ends_at) <= 1_100,
1076				"frames at {} and {} do not meet",
1077				pair[0].0,
1078				pair[1].0
1079			);
1080		}
1081	}
1082
1083	/// The tolerance can't come from a frame duration. Opus packets run from 2.5 ms
1084	/// to 60 ms with nothing in the catalog to say which, so a rule read off the
1085	/// 20 ms packet before it would splice straight across a lost 2.5 ms one.
1086	#[tokio::test]
1087	async fn a_lost_opus_packet_shorter_than_its_neighbour_is_a_gap() {
1088		let input = Input {
1089			format: Format::F32,
1090			sample_rate: 48_000,
1091			layout: Layout::Mono,
1092		};
1093		let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1094		let catalog = encoder.catalog();
1095
1096		let broadcast = moq_net::broadcast::Info::new().produce();
1097		let track = broadcast
1098			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1099			.unwrap();
1100		let subscriber = broadcast.consume();
1101		let mut producer = moq_mux::container::Producer::new(
1102			track,
1103			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1104		);
1105		let mut consumer = Consumer::new(
1106			&subscriber,
1107			&catalog,
1108			"audio",
1109			Options {
1110				max_age: std::time::Duration::from_secs(1),
1111				..Options::new()
1112			},
1113		)
1114		.await
1115		.unwrap();
1116
1117		// A 20 ms packet at 0, then the next one at 22.5 ms: the 2.5 ms packet
1118		// between them was lost.
1119		let pcm = vec![0.25f32; encoder.frame_size()];
1120		for timestamp in [
1121			Timestamp::from_micros(0).unwrap(),
1122			Timestamp::from_micros(22_500).unwrap(),
1123			Timestamp::from_micros(42_500).unwrap(),
1124		] {
1125			producer
1126				.write(moq_mux::container::Frame {
1127					timestamp,
1128					duration: None,
1129					payload: encoder.encode(&pcm).unwrap().payload,
1130					keyframe: true,
1131				})
1132				.unwrap();
1133			producer.cut(None).unwrap();
1134		}
1135
1136		// The pre-skip is trimmed off the first packet, so it decodes short. That
1137		// shortfall is codec delay, not a hole: without counting it the packet after
1138		// every stream start would read as a gap.
1139		let first = consumer.read().await.unwrap().expect("decoded frame");
1140		let frames = first.data.len() / size_of::<f32>();
1141		assert!(frames < 960, "the pre-skip should be trimmed, got {frames} frames");
1142
1143		// The hole is real, so codec prediction starts over but stream-level pre-skip
1144		// does not. The audio is stamped where the packet says rather than 2.5 ms early.
1145		let second = consumer.read().await.unwrap().expect("decoded frame");
1146		assert_eq!(second.timestamp.as_micros(), 22_500);
1147		assert_eq!(second.data.len() / size_of::<f32>(), 960, "pre-skip was reapplied");
1148
1149		let third = consumer.read().await.unwrap().expect("decoded frame after gap");
1150		let second_frames = second.data.len() / size_of::<f32>();
1151		assert_eq!(
1152			third.timestamp,
1153			advance(second.timestamp, second_frames, 48_000).unwrap()
1154		);
1155	}
1156
1157	#[tokio::test]
1158	async fn max_age_is_clamped_to_publisher_retention() {
1159		let broadcast = moq_net::broadcast::Info::new().produce();
1160		let info = hang::container::track_info(hang::catalog::PRIORITY.audio)
1161			.with_max_age(std::time::Duration::from_millis(100));
1162		let _track = broadcast.create_track("audio", info).unwrap();
1163		let subscriber = broadcast.consume();
1164		let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
1165
1166		let consumer = Consumer::new(
1167			&subscriber,
1168			&catalog,
1169			"audio",
1170			Options {
1171				max_age: std::time::Duration::from_millis(500),
1172				..Options::new()
1173			},
1174		)
1175		.await
1176		.unwrap();
1177
1178		assert_eq!(consumer.max_age(), std::time::Duration::from_millis(100));
1179	}
1180
1181	/// Opus pre-skip is padding before the decoded epoch, not missing media after
1182	/// the first short frame. The second frame must meet the first or playback
1183	/// fills the codec delay with silence and creates a startup glitch.
1184	#[tokio::test]
1185	async fn opus_pre_skip_does_not_leave_a_timestamp_hole() {
1186		let input = Input {
1187			format: Format::F32,
1188			sample_rate: 48_000,
1189			layout: Layout::Mono,
1190		};
1191		let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1192		let catalog = encoder.catalog();
1193
1194		let broadcast = moq_net::broadcast::Info::new().produce();
1195		let track = broadcast
1196			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1197			.unwrap();
1198		let subscriber = broadcast.consume();
1199		let mut producer = moq_mux::container::Producer::new(
1200			track,
1201			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1202		);
1203		let mut consumer = Consumer::new(
1204			&subscriber,
1205			&catalog,
1206			"audio",
1207			Options {
1208				max_age: std::time::Duration::from_secs(1),
1209				..Options::new()
1210			},
1211		)
1212		.await
1213		.unwrap();
1214
1215		let pcm = vec![0.25f32; encoder.frame_size()];
1216		for packet in 0..2 {
1217			producer
1218				.write(moq_mux::container::Frame {
1219					timestamp: Timestamp::from_scale(packet * encoder.frame_size() as u64, 48_000).unwrap(),
1220					duration: None,
1221					payload: encoder.encode(&pcm).unwrap().payload,
1222					keyframe: true,
1223				})
1224				.unwrap();
1225			producer.cut(None).unwrap();
1226		}
1227
1228		let first = consumer.read().await.unwrap().expect("first decoded frame");
1229		let second = consumer.read().await.unwrap().expect("second decoded frame");
1230		let first_frames = first.data.len() / size_of::<f32>();
1231		let expected = advance(first.timestamp, first_frames, 48_000).unwrap();
1232		assert_eq!(second.timestamp, expected);
1233	}
1234
1235	#[tokio::test]
1236	async fn a_playhead_event_reapplies_opus_pre_skip() {
1237		let input = Input {
1238			format: Format::F32,
1239			sample_rate: 48_000,
1240			layout: Layout::Mono,
1241		};
1242		let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1243		let catalog = encoder.catalog();
1244		let frame_size = encoder.frame_size();
1245
1246		let broadcast = moq_net::broadcast::Info::new().produce();
1247		let track = broadcast
1248			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1249			.unwrap();
1250		let subscriber = broadcast.consume();
1251		let mut producer = moq_mux::container::Producer::new(
1252			track,
1253			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1254		);
1255		let mut consumer = Consumer::new(
1256			&subscriber,
1257			&catalog,
1258			"audio",
1259			Options {
1260				max_age: std::time::Duration::from_secs(1),
1261				..Options::new()
1262			},
1263		)
1264		.await
1265		.unwrap();
1266
1267		let pcm = vec![0.25f32; frame_size];
1268		let write = |producer: &mut moq_mux::container::Producer<_>, packet: u64, payload: bytes::Bytes| {
1269			producer
1270				.write(moq_mux::container::Frame {
1271					timestamp: Timestamp::from_scale(packet * frame_size as u64, 48_000).unwrap(),
1272					duration: None,
1273					payload,
1274					keyframe: true,
1275				})
1276				.unwrap();
1277			producer.cut(None).unwrap();
1278		};
1279		write(&mut producer, 0, encoder.encode(&pcm).unwrap().payload);
1280		write(&mut producer, 1, encoder.encode(&pcm).unwrap().payload);
1281		producer.discontinuity().unwrap();
1282		write(&mut producer, 2, encoder.encode(&pcm).unwrap().payload);
1283		producer.finish().unwrap();
1284
1285		let first = consumer.read().await.unwrap().expect("first decoded frame");
1286		let _second = consumer.read().await.unwrap().expect("second decoded frame");
1287		let resumed = consumer.read().await.unwrap().expect("resumed decoded frame");
1288		let first_frames = first.data.len() / size_of::<f32>();
1289		let resumed_frames = resumed.data.len() / size_of::<f32>();
1290		assert!(first_frames < frame_size, "the first epoch trims pre-skip");
1291		assert_eq!(
1292			resumed_frames, first_frames,
1293			"a playhead event reapplies pre-skip without flushing the decoder"
1294		);
1295	}
1296
1297	#[tokio::test]
1298	async fn reads_the_container_the_catalog_declares() {
1299		let broadcast = moq_net::broadcast::Info::new().produce();
1300		let track = broadcast
1301			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1302			.unwrap();
1303		let observed = track.clone();
1304		let subscriber = broadcast.consume();
1305
1306		let mut catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
1307		catalog.container = hang::catalog::Container::Loc;
1308
1309		let mut producer = moq_mux::container::Producer::new(
1310			track,
1311			moq_mux::catalog::hang::Container::Loc(moq_mux::container::Kind::Audio),
1312		);
1313		let max_age = std::time::Duration::from_millis(250);
1314		let mut consumer = Consumer::new(
1315			&subscriber,
1316			&catalog,
1317			"audio",
1318			Options {
1319				max_age,
1320				..Options::new()
1321			},
1322		)
1323		.await
1324		.unwrap();
1325		assert_eq!(observed.subscription().unwrap().max_age, max_age);
1326
1327		let samples = [0.25f32, -0.5, 0.75, -1.0];
1328		let payload: Vec<u8> = samples.iter().flat_map(|sample| sample.to_le_bytes()).collect();
1329		producer
1330			.write(moq_mux::container::Frame {
1331				timestamp: Timestamp::ZERO,
1332				duration: None,
1333				payload: payload.into(),
1334				keyframe: true,
1335			})
1336			.unwrap();
1337
1338		let frame = consumer.read().await.unwrap().expect("decoded frame");
1339		assert_eq!(
1340			Format::F32.as_interleaved_f32(&frame.data, 1).unwrap().as_ref(),
1341			samples
1342		);
1343	}
1344
1345	/// The catalog picks the framing, not this crate. Hardcoding the legacy wire
1346	/// read a CMAF fragment as a varint timestamp plus a payload, which handed the
1347	/// codec garbage instead of failing, so anything published by `moq import
1348	/// fmp4` was undecodable.
1349	#[tokio::test]
1350	async fn decodes_a_cmaf_framed_track() {
1351		let input = Input {
1352			format: Format::F32,
1353			sample_rate: 48_000,
1354			layout: Layout::Stereo,
1355		};
1356
1357		// One real Opus packet, so a mis-framed read can't accidentally decode.
1358		let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1359		let mut catalog = encoder.catalog();
1360		let pcm = vec![0.0f32; encoder.frame_size() * encoder.codec_channels() as usize];
1361		let packet = encoder.encode(&pcm).unwrap();
1362
1363		// Re-describe the same rendition as CMAF and publish it that way.
1364		let muxer = moq_mux::container::fmp4::Muxer::audio(&catalog).unwrap();
1365		let init = muxer.init().unwrap().expect("an out-of-band codec has an init segment");
1366		catalog.container = hang::catalog::Container::Cmaf { init };
1367
1368		let broadcast = moq_net::broadcast::Info::new().produce();
1369		let subscriber = broadcast.consume();
1370		let track = broadcast
1371			.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1372			.unwrap();
1373		let container = moq_mux::catalog::hang::Container::try_from(&catalog).unwrap();
1374		let mut producer = moq_mux::container::Producer::new(track, container);
1375
1376		let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Options::new())
1377			.await
1378			.unwrap();
1379
1380		producer
1381			.write(moq_mux::container::Frame {
1382				timestamp: Timestamp::ZERO,
1383				payload: packet.payload,
1384				keyframe: true,
1385				duration: None,
1386			})
1387			.unwrap();
1388		producer.cut(None).unwrap();
1389
1390		// The whole packet decodes: one 20 ms Opus frame at 48 kHz, less the pre-skip
1391		// trimmed off the first packet. Reading the fragment as legacy hands the codec
1392		// a slice of the moof instead, which still decodes, just to a shorter buffer.
1393		let frame = consumer.read().await.unwrap().expect("decoded frame");
1394		// `as_micros`, not `==`: the CMAF path carries the fmp4 timescale and
1395		// `Timestamp`'s equality is structural, so the scales would have to match too.
1396		assert_eq!(frame.timestamp.as_micros(), 0);
1397		let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
1398		assert_eq!(samples.len(), (960 - 312) * 2);
1399	}
1400}