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