Skip to main content

moq_audio/encode/
producer.rs

1//! Encode raw PCM and publish it as a moq audio track.
2
3use std::time::Instant;
4
5use bytes::Bytes;
6
7use moq_mux::catalog::hang::CatalogExt;
8use moq_mux::container::Frame as MuxFrame;
9use moq_net::Timestamp;
10
11use super::encoded::Encoded;
12use super::encoder::{Encoder, Input, Settings};
13use crate::resample::{Remix, Resampler};
14use crate::{Activity, Error, Frame};
15
16/// Encode and publication policy for [`Producer`].
17///
18/// `#[non_exhaustive]`: construct via [`Options::default`] and set fields, so
19/// new knobs can be added without breaking callers.
20#[derive(Clone, Debug)]
21#[non_exhaustive]
22pub struct Options {
23	/// Track name to publish under. `None` derives a unique one from the codec
24	/// (`0.opus`, then `1.opus`, ...), matching how the video side names its
25	/// track. Subscribers find it through the catalog either way.
26	pub track: Option<String>,
27	/// Codec input and encoded output settings.
28	pub settings: Settings,
29	/// The connection's bandwidth, as an allocator over
30	/// [`Session::send_bandwidth`](moq_net::Session::send_bandwidth).
31	///
32	/// The audio track reserves its bitrate against it, so the video encoder sharing
33	/// the connection sizes itself against what's actually left rather than against
34	/// the whole uplink. Pass the same allocator to both.
35	///
36	/// Defaults to [`Allocator::unlimited`](moq_net::bandwidth::Allocator::unlimited),
37	/// which reserves nothing and leaves every sender at its configured rate.
38	///
39	/// Audio reserves but does not follow its share: Opus can retune live and PCM
40	/// can't at all, and at `hang`'s priorities audio outranks video, so it is only
41	/// ever squeezed on a link that can't carry audio alone.
42	pub bandwidth: moq_net::bandwidth::Allocator,
43}
44
45impl Default for Options {
46	fn default() -> Self {
47		Self {
48			track: None,
49			settings: Settings::default(),
50			bandwidth: moq_net::bandwidth::Allocator::unlimited(),
51		}
52	}
53}
54
55/// Encode raw PCM and publish it as a moq-mux audio track.
56///
57/// The input PCM layout is fixed at construction via [`Input`]; the codec
58/// settings via [`Options`]. Subsequent [`write`](Self::write) calls just pass a
59/// [`Frame`]: payload bytes and a timestamp.
60///
61/// The catalog rendition is registered at construction (not on first write), so
62/// a subscriber that opens the catalog before any frames arrive still sees the
63/// track.
64pub struct Producer<E: CatalogExt = ()> {
65	encoder: Encoder,
66	input: Input,
67	/// Converts the input layout to the codec's, when they differ.
68	remix: Option<Remix>,
69	resampler: Option<Resampler>,
70	track: moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
71	_ext: std::marker::PhantomData<fn() -> E>,
72	pending: Vec<f32>,
73	/// Samples emitted since the current epoch (reset by [`reset_epoch`](Self::reset_epoch)).
74	frames_produced: u64,
75	/// Wall-clock anchor in microseconds, taken from the first frame after each
76	/// (re)start. Emitted PTS = `epoch + frames_produced / codec_rate`. `None`
77	/// until the first write so the next frame re-anchors to its timestamp.
78	epoch_us: Option<u64>,
79	/// An encoder reset that still needs a marker group before its next packet.
80	pending_discontinuity: bool,
81	/// Whether a marker group already separates the next packet from prior codec state.
82	decoder_boundary: bool,
83	/// How the encoder classified the packet it published most recently.
84	activity: Activity,
85	/// Set after a successful [`finish`](Self::finish). Writes then fail with
86	/// [`moq_net::Error::Closed`]; [`abort`](Self::abort) can still run.
87	finished: bool,
88}
89
90struct Terminal {
91	packets: Vec<Encoded>,
92	end: Timestamp,
93	start: Timestamp,
94	frame_size: usize,
95	codec_rate: u32,
96}
97
98/// A published track whose PCM layout is not known yet.
99///
100/// The track exists in the broadcast immediately, so its name and subscriber
101/// state are available, while the catalog rendition waits for
102/// [`encode`](Self::encode) to supply the layout it describes. Unlike a catalog
103/// [`Reserved`](moq_mux::catalog::Reserved) this does not withhold the catalog:
104/// subscribers see the broadcast without this rendition until it resolves.
105pub(crate) struct Reserved<E: CatalogExt = ()> {
106	track: moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
107	_ext: std::marker::PhantomData<fn() -> E>,
108}
109
110impl<E: CatalogExt> Reserved<E> {
111	pub(crate) fn new(
112		broadcast: &mut moq_net::broadcast::Producer,
113		catalog: moq_mux::catalog::Producer<E>,
114		options: &Options,
115	) -> Result<Self, Error> {
116		let track = match &options.track {
117			// The catalog's info carries the microsecond timescale audio hang frames stamp, so
118			// Lite05 subscribers know what scale to expect and the model layer accepts
119			// Frame::timestamp on append, plus whatever retention the broadcast declared.
120			Some(name) => broadcast.create_track(name.clone(), catalog.track_info(hang::catalog::PRIORITY.audio))?,
121			// Mirrors the video side, which derives a unique name from the codec
122			// rather than making every caller invent one.
123			None => broadcast.unique_track(
124				&format!(".{}", options.settings.codec),
125				catalog.track_info(hang::catalog::PRIORITY.audio),
126			)?,
127		};
128		let track = catalog.audio(
129			track,
130			moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio),
131			None,
132		)?;
133
134		Ok(Self {
135			track,
136			_ext: std::marker::PhantomData,
137		})
138	}
139
140	/// Build the encoder for `input` and register the rendition describing it.
141	///
142	/// Separate from [`encode`](Self::encode), which cannot fail, so a layout the
143	/// codec rejects leaves the reservation intact for another input.
144	pub(crate) fn register(&mut self, input: Input, options: &Options) -> Result<Registered, Error> {
145		let remix = (input.layout != options.settings.layout)
146			.then(|| Remix::new(input.layout, options.settings.layout))
147			.transpose()?;
148		let encoder = Encoder::new(&options.settings)?;
149
150		let resampler = if input.sample_rate == encoder.codec_rate() {
151			None
152		} else {
153			// Use microsecond precision so 2.5 ms frame_duration (supported by
154			// libopus) doesn't truncate to 2 ms.
155			let chunk_frames =
156				((input.sample_rate as u128 * encoder.settings().frame_duration.as_micros()) / 1_000_000) as usize;
157			Some(Resampler::new(
158				input.sample_rate,
159				encoder.codec_rate(),
160				encoder.codec_channels(),
161				chunk_frames,
162			)?)
163		};
164
165		self.track.set(encoder.catalog())?;
166
167		Ok(Registered {
168			encoder,
169			input,
170			remix,
171			resampler,
172		})
173	}
174
175	/// Spend the reservation on a registered encoder, publishing through the track.
176	pub(crate) fn encode(self, registered: Registered) -> Producer<E> {
177		Producer {
178			encoder: registered.encoder,
179			input: registered.input,
180			remix: registered.remix,
181			resampler: registered.resampler,
182			track: self.track,
183			_ext: self._ext,
184			pending: Vec::new(),
185			frames_produced: 0,
186			epoch_us: None,
187			pending_discontinuity: false,
188			decoder_boundary: true,
189			activity: Activity::Active,
190			finished: false,
191		}
192	}
193}
194
195/// A registered rendition, waiting for its [`Reserved`] to be spent on it.
196///
197/// Proof that [`Reserved::register`] succeeded, so [`Reserved::encode`] takes no
198/// fallible step and cannot strand the track it consumes.
199pub(crate) struct Registered {
200	encoder: Encoder,
201	input: Input,
202	remix: Option<Remix>,
203	resampler: Option<Resampler>,
204}
205
206/// What a capture publication needs while its layout is still undiscovered.
207#[cfg(feature = "capture")]
208impl<E: CatalogExt> Reserved<E> {
209	/// The resolved track name, available before the layout is.
210	pub(crate) fn name(&self) -> &str {
211		self.track.name()
212	}
213
214	/// The underlying track producer, e.g. to watch subscriber state.
215	pub(crate) fn track(&self) -> &moq_net::track::Producer {
216		self.track.track()
217	}
218
219	/// Finalize a track that never got a rendition.
220	pub(crate) fn finish(&mut self) -> Result<(), Error> {
221		self.track.finish()?;
222		Ok(())
223	}
224
225	/// Abort a track that never got a rendition, so subscribers see `err`.
226	pub(crate) fn abort(self, err: moq_net::Error) {
227		self.track.abort(err);
228	}
229}
230
231impl<E: CatalogExt> Producer<E> {
232	/// Publish a track encoding `input` into `broadcast`, registering its
233	/// rendition in `catalog` immediately.
234	pub fn new(
235		broadcast: &mut moq_net::broadcast::Producer,
236		catalog: moq_mux::catalog::Producer<E>,
237		input: Input,
238		options: &Options,
239	) -> Result<Self, Error> {
240		let mut reserved = Reserved::new(broadcast, catalog, options)?;
241		let registered = reserved.register(input, options)?;
242		Ok(reserved.encode(registered))
243	}
244
245	/// The name of the published track, which is [`Options::track`] resolved.
246	pub fn track_name(&self) -> &str {
247		self.track.name()
248	}
249
250	/// A watch-only handle to the track's subscriber demand, created eagerly so
251	/// subscription state is observable before any frames arrive. Watch it via
252	/// [`used`](moq_net::track::Demand::used) / [`unused`](moq_net::track::Demand::unused).
253	pub fn demand(&self) -> moq_net::track::Demand {
254		self.track.track().demand()
255	}
256
257	#[cfg(feature = "capture")]
258	pub(crate) fn track(&self) -> &moq_net::track::Producer {
259		self.track.track()
260	}
261
262	/// Whether the packet published most recently coded audio, or withheld it
263	/// because the input was silent.
264	///
265	/// A local "am I talking" indicator without running a second voice detector
266	/// over the microphone, though a silent run is punctuated by coded frames
267	/// that read [`Activity::Active`], so hold the indicator across those rather
268	/// than following it packet by packet. [`Activity::Active`] until the first
269	/// packet, and for codecs without a discontinuous mode.
270	pub fn activity(&self) -> Activity {
271		self.activity
272	}
273
274	/// Current encoder target bitrate.
275	pub fn bitrate(&self) -> moq_net::bandwidth::Rate {
276		self.encoder.bitrate()
277	}
278
279	/// Retune the live encoder to `bitrate`.
280	pub fn set_bitrate(&mut self, bitrate: moq_net::bandwidth::Rate) -> Result<(), Error> {
281		self.encoder.set_bitrate(bitrate)
282	}
283
284	/// Re-anchor the timeline to the next frame's timestamp, dropping any
285	/// buffered samples. Call this when resuming after an idle gap (e.g. a
286	/// released-then-reopened microphone) so the gap appears in the PTS and
287	/// audio stays aligned with a wall-clock video track, rather than the gap
288	/// being compressed out by the running sample count. Mirrors moq-boy's
289	/// `reset_epoch`. If the codec had started, a marker group is published before
290	/// the next packet so subscribers jump the playhead.
291	pub fn reset_epoch(&mut self) {
292		if self.encoder.started() && !self.decoder_boundary {
293			self.pending_discontinuity = true;
294		}
295		self.reset_state();
296	}
297
298	fn reset_state(&mut self) {
299		self.epoch_us = None;
300		self.activity = Activity::Active;
301		self.frames_produced = 0;
302		self.pending.clear();
303		self.encoder.reset();
304		// The resampler holds samples of its own, plus filter state primed by them.
305		// Left alone, `finish` would flush that pre-reset audio onto the track
306		// (stamped at an epoch that no longer exists), and the next write would run
307		// the new audio through a filter still ringing with the old.
308		if let Some(resampler) = self.resampler.as_mut() {
309			resampler.reset();
310		}
311	}
312
313	/// Push one [`Frame`] of PCM in the layout declared by [`Input`]. Encodes and
314	/// publishes as many packets as the input contains; any partial trailing
315	/// frame is carried to the next call.
316	///
317	/// The first frame after construction (or [`reset_epoch`](Self::reset_epoch))
318	/// anchors the timeline: its timestamp becomes the epoch, and emitted PTS
319	/// then advances purely by the running sample count, so subsequent frames'
320	/// timestamps are ignored. An idle gap is only reflected in the PTS if you
321	/// call [`reset_epoch`](Self::reset_epoch) on resume (which re-anchors from
322	/// the next frame's wall-clock stamp); writing straight across a gap without
323	/// resetting compresses it out.
324	///
325	/// [`Frame::activity`] is ignored: the encoder classifies what it actually
326	/// produced, which [`activity`](Self::activity) reports.
327	///
328	/// Writes after a successful [`finish`](Self::finish) fail with
329	/// [`moq_net::Error::Closed`].
330	pub fn write(&mut self, frame: &Frame) -> Result<(), Error> {
331		if self.finished {
332			return Err(moq_net::Error::Closed.into());
333		}
334		if self.pending_discontinuity {
335			self.track.discontinuity()?;
336			self.pending_discontinuity = false;
337			self.decoder_boundary = true;
338		}
339
340		let timestamp_us = u64::try_from(frame.timestamp.as_micros())
341			.map_err(|_| Error::Unsupported(format!("frame timestamp {:?} out of range", frame.timestamp)))?;
342		let epoch_us = *self.epoch_us.get_or_insert(timestamp_us);
343
344		let input = &self.input;
345		let (format, channels) = (input.format, input.layout.channels());
346		let pcm = format.as_interleaved_f32(frame.data.as_ref(), channels)?;
347		let pcm = match &self.remix {
348			Some(remix) => remix.process(&pcm),
349			None => pcm.into_owned(),
350		};
351		let pcm: Vec<f32> = match self.resampler.as_mut() {
352			Some(r) => r.process(&pcm, frame.timestamp)?,
353			None => pcm,
354		};
355
356		self.pending.extend(pcm);
357
358		self.publish_full_frames(epoch_us)
359	}
360
361	/// Encode and publish every full frame in `pending`, keeping any partial
362	/// trailing frame for the next call.
363	fn publish_full_frames(&mut self, epoch_us: u64) -> Result<(), Error> {
364		let frame_samples = self.encoder.frame_size() * self.encoder.codec_channels() as usize;
365		while self.pending.len() >= frame_samples {
366			let chunk: Vec<f32> = self.pending.drain(..frame_samples).collect();
367			let packet = self.encoder.encode(&chunk)?;
368
369			let timestamp = Self::timestamp(
370				epoch_us,
371				self.frames_produced,
372				self.encoder.folded_delay(),
373				self.encoder.codec_rate(),
374			)?;
375			self.frames_produced += self.encoder.frame_size() as u64;
376			self.activity = packet.activity;
377			Self::publish(&mut self.track, packet, timestamp)?;
378			self.decoder_boundary = false;
379		}
380
381		Ok(())
382	}
383
384	/// PTS of the frame `frames` samples past the epoch, stamped `delay` samples
385	/// earlier to fold in codec priming the catalog can't signal.
386	///
387	/// Priming that would land before a zero epoch is stamped at zero instead: it
388	/// decodes to the codec's warm-up rather than to input, so only its spacing is
389	/// lost.
390	fn timestamp(epoch_us: u64, frames: u64, delay: usize, codec_rate: u32) -> Result<Timestamp, Error> {
391		let frames = i128::from(frames) - delay as i128;
392		let offset_us = (frames * 1_000_000).div_euclid(i128::from(codec_rate));
393		let micros = (i128::from(epoch_us) + offset_us).max(0);
394		let micros = u64::try_from(micros).map_err(|_| moq_net::TimeOverflow)?;
395		Ok(Timestamp::from_micros(micros)?)
396	}
397
398	fn publish(
399		track: &mut moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
400		encoded: Encoded,
401		timestamp: Timestamp,
402	) -> Result<(), Error> {
403		// Publish each audio packet as its own moq-lite group: write it as a keyframe, then cut
404		// (below) so the relay forwards it without waiting for the next. Codecs can recover
405		// independently after a dropped group.
406		let mux_frame = MuxFrame {
407			timestamp,
408			payload: encoded.payload,
409			keyframe: true,
410			duration: None,
411		};
412		track.write(mux_frame)?;
413		// No boundary to give: the next packet bounds this one, and Opus frames have a
414		// deterministic duration anyway. Cut before observing the flush so a failed observation
415		// never leaves the group open.
416		track.cut(None)?;
417		track.flush(timestamp, Instant::now())?;
418		Ok(())
419	}
420
421	/// Publish terminal packets after an empty frame that carries their logical endpoint.
422	fn publish_terminal(
423		track: &mut moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
424		terminal: Terminal,
425	) -> Result<(), Error> {
426		track.write(MuxFrame {
427			timestamp: terminal.end,
428			payload: Bytes::new(),
429			keyframe: true,
430			duration: None,
431		})?;
432
433		for (index, packet) in terminal.packets.into_iter().enumerate() {
434			let offset = Timestamp::from_scale((index * terminal.frame_size) as u64, terminal.codec_rate as u64)?
435				.convert(terminal.start.scale())?;
436			let timestamp = terminal.start.checked_add(offset)?;
437			track.write(MuxFrame {
438				timestamp,
439				payload: packet.payload,
440				keyframe: false,
441				duration: None,
442			})?;
443			track.flush(timestamp, Instant::now())?;
444		}
445
446		track.cut(Some(terminal.end))?;
447		Ok(())
448	}
449
450	/// Mark a break in the published timeline and reset codec state.
451	///
452	/// Call this when capture stops rather than merely gapping between packets: going idle,
453	/// switching source, or anything else that resumes on a re-anchored epoch. Buffered samples
454	/// are dropped, and the next frame anchors a fresh codec epoch. See
455	/// [`Producer::discontinuity`](moq_mux::container::Producer::discontinuity).
456	pub fn discontinuity(&mut self) -> Result<(), Error> {
457		self.track.discontinuity()?;
458		self.pending_discontinuity = false;
459		self.decoder_boundary = true;
460		self.reset_state();
461		Ok(())
462	}
463
464	/// Flush pending samples, resampler output, and codec lookahead, then finalize
465	/// the track.
466	///
467	/// Borrows rather than consumes, so a later [`abort`](Self::abort) can still
468	/// run after a successful finish. Writes after this fail with
469	/// [`moq_net::Error::Closed`].
470	pub fn finish(&mut self) -> Result<(), Error> {
471		if self.finished {
472			return Ok(());
473		}
474		// Whatever the resampler still holds belongs to this track: its last partial
475		// chunk, plus the audio its filter is running behind on. Dropping it here
476		// would publish a track that ends before its source did.
477		if let Some(resampler) = self.resampler.take() {
478			self.pending.extend(resampler.flush()?);
479		}
480
481		// The drained resampler tail can span multiple frames. Publish those first
482		// so only the final partial frame reaches the encoder's terminal drain.
483		let epoch_us = self.epoch_us.unwrap_or(0);
484		self.publish_full_frames(epoch_us)?;
485
486		let frame_size = self.encoder.frame_size();
487		let codec_rate = self.encoder.codec_rate();
488		let channels = self.encoder.codec_channels() as usize;
489		let source_frames = self.pending.len() / channels;
490		let delay = self.encoder.folded_delay();
491		let start = Self::timestamp(epoch_us, self.frames_produced, delay, codec_rate)?;
492		// The source ends where it ends: priming only moves the packets carrying it.
493		let end = Self::timestamp(epoch_us, self.frames_produced + source_frames as u64, 0, codec_rate)?;
494		let finish = self.encoder.drain(&self.pending)?;
495		let discard_padding = finish.discard_padding();
496		let packets = finish.into_packets();
497
498		if discard_padding > 0 {
499			Self::publish_terminal(
500				&mut self.track,
501				Terminal {
502					packets,
503					end,
504					start,
505					frame_size,
506					codec_rate,
507				},
508			)?;
509		} else {
510			for packet in packets {
511				let timestamp = Self::timestamp(epoch_us, self.frames_produced, delay, codec_rate)?;
512				self.activity = packet.activity;
513				Self::publish(&mut self.track, packet, timestamp)?;
514				self.frames_produced += frame_size as u64;
515			}
516		}
517
518		self.track.finish()?;
519		self.finished = true;
520		Ok(())
521	}
522
523	/// Abort the track with `err` instead of finishing it, so subscribers see the
524	/// real cause rather than [`moq_net::Error::Dropped`]. Pending samples are dropped.
525	///
526	/// Consumes the producer. Still callable after [`finish`](Self::finish).
527	pub fn abort(self, err: moq_net::Error) {
528		self.track.abort(err);
529	}
530}
531
532#[cfg(test)]
533mod tests {
534	use std::time::Duration;
535
536	use super::*;
537	use crate::decode::{Consumer as AudioConsumer, Options as DecodeOptions};
538	use crate::{Activity, Format, Layout};
539
540	#[tokio::test]
541	async fn demand_follows_subscribers_and_closes_with_the_producer() {
542		let mut broadcast = moq_net::broadcast::Info::new().produce();
543		let consumer = broadcast.consume();
544		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
545		let options = Options {
546			track: Some("audio".into()),
547			..Options::default()
548		};
549		let mut producer = Producer::new(&mut broadcast, catalog, Input::default(), &options).unwrap();
550		let demand = producer.demand();
551
552		assert_eq!(demand.name(), "audio");
553		assert!(!demand.is_used());
554		let subscriber = consumer.track("audio").unwrap().subscribe(None).await.unwrap();
555		tokio::time::timeout(Duration::from_secs(1), demand.used())
556			.await
557			.expect("subscription demand")
558			.unwrap();
559		assert!(demand.is_used());
560
561		drop(subscriber);
562		drop(consumer);
563		tokio::time::timeout(Duration::from_secs(1), demand.unused())
564			.await
565			.expect("subscription released")
566			.unwrap();
567		assert!(!demand.is_used());
568
569		producer.finish().unwrap();
570		drop(producer);
571		let closed = tokio::time::timeout(Duration::from_secs(1), demand.closed())
572			.await
573			.expect("producer closed");
574		assert!(matches!(closed, moq_net::Error::Dropped));
575	}
576
577	/// Terminal Opus lookahead samples survive both exact-frame and partial-frame input.
578	#[tokio::test]
579	async fn finish_publishes_the_opus_lookahead_tail() {
580		for frames in [960, 860] {
581			let input = Input {
582				format: Format::F32,
583				sample_rate: 48_000,
584				layout: Layout::Mono,
585			};
586			let options = Options {
587				track: Some("audio".to_string()),
588				settings: Settings {
589					layout: Layout::Mono,
590					bitrate: Some(moq_net::bandwidth::Rate::from_bps(128_000)),
591					..Settings::default()
592				},
593				..Options::default()
594			};
595			let decoder_config = Encoder::new(&options.settings).unwrap().catalog();
596
597			let mut broadcast = moq_net::broadcast::Info::new().produce();
598			let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
599			let consumer = broadcast.consume();
600			let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
601			let mut audio = AudioConsumer::new(
602				&consumer,
603				&decoder_config,
604				"audio",
605				DecodeOptions {
606					max_age: Duration::from_secs(1),
607					..DecodeOptions::new()
608				},
609			)
610			.await
611			.unwrap();
612
613			let mut pcm = vec![0.0f32; frames];
614			let impulse = pcm.len() - 100;
615			pcm[impulse] = 1.0;
616			let data: Vec<u8> = pcm.iter().flat_map(|sample| sample.to_le_bytes()).collect();
617			producer.write(&Frame::new(Bytes::from(data), Timestamp::ZERO)).unwrap();
618			producer.finish().unwrap();
619
620			let mut decoded = Vec::new();
621			while let Some(frame) = audio.read().await.unwrap() {
622				let pcm = Format::F32.as_interleaved_f32(&frame.data, 1).unwrap();
623				decoded.extend_from_slice(&pcm);
624			}
625			assert_eq!(decoded.len(), frames, "terminal padding extended the source");
626			let peak = decoded.iter().fold(0.0f32, |peak, sample| peak.max(sample.abs()));
627			assert!(peak > 0.1, "the {frames}-frame Opus tail lost the impulse: peak {peak}");
628		}
629	}
630
631	/// A resampled publisher used to end its track early: `finish` flushed the
632	/// encoder's own buffer but left the resampler holding its last partial chunk,
633	/// plus the audio its filter runs behind on.
634	#[tokio::test]
635	async fn finish_publishes_the_resampled_tail() {
636		let input = Input {
637			format: Format::F32,
638			sample_rate: 44_100,
639			layout: Layout::Mono,
640		};
641
642		let mut broadcast = moq_net::broadcast::Info::new().produce();
643		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
644		let consumer = broadcast.consume();
645		let options = Options {
646			track: Some("audio".to_string()),
647			settings: Settings::new(48_000, Layout::Mono),
648			..Options::default()
649		};
650		let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
651
652		// Subscribe before the track ends, or there is nothing left to subscribe to.
653		let mut track = moq_mux::container::Consumer::new(
654			consumer
655				.track("audio")
656				.unwrap()
657				.subscribe(moq_net::track::Subscription::default().with_max_age(Duration::from_secs(1)))
658				.await
659				.unwrap(),
660			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
661		);
662
663		// Chosen so the tail decides a whole packet: 8838 frames at 44.1 kHz is ~9620
664		// at 48 kHz, just past ten 960-sample Opus frames. Losing the resampler's
665		// remainder and its filter delay drops back under ten, costing a packet.
666		let data: Vec<u8> = vec![0.25f32; 8_838].iter().flat_map(|s| s.to_le_bytes()).collect();
667		producer
668			.write(&Frame::new(data.into(), moq_net::Timestamp::ZERO))
669			.unwrap();
670		producer.finish().unwrap();
671
672		let mut packets = 0;
673		while track.read().await.unwrap().is_some() {
674			packets += 1;
675		}
676		assert_eq!(packets, 11);
677	}
678
679	/// `reset_epoch` promises to drop buffered samples, and the resampler buffers
680	/// samples of its own. Leaving those behind let `finish` flush pre-reset audio
681	/// onto the track, stamped at an epoch that no longer exists.
682	#[tokio::test]
683	async fn reset_epoch_drops_the_resampler_buffer_too() {
684		let input = Input {
685			format: Format::F32,
686			sample_rate: 44_100,
687			layout: Layout::Mono,
688		};
689
690		let mut broadcast = moq_net::broadcast::Info::new().produce();
691		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
692		let consumer = broadcast.consume();
693		let options = Options {
694			track: Some("audio".to_string()),
695			..Options::default()
696		};
697		let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
698
699		let mut track = moq_mux::container::Consumer::new(
700			consumer
701				.track("audio")
702				.unwrap()
703				.subscribe(moq_net::track::Subscription::default())
704				.await
705				.unwrap(),
706			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
707		);
708
709		// Too little to publish a packet, so it all sits in the resampler.
710		let data: Vec<u8> = vec![0.25f32; 441].iter().flat_map(|s| s.to_le_bytes()).collect();
711		producer
712			.write(&Frame::new(data.into(), moq_net::Timestamp::ZERO))
713			.unwrap();
714
715		producer.reset_epoch();
716		producer.finish().unwrap();
717
718		// The reset dropped everything, so the track ends without a packet.
719		assert!(track.read().await.unwrap().is_none());
720	}
721
722	/// Resetting after a full frame drops codec lookahead as well as producer buffers.
723	#[tokio::test]
724	async fn reset_epoch_drops_the_encoder_lookahead() {
725		let mut broadcast = moq_net::broadcast::Info::new().produce();
726		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
727		let consumer = broadcast.consume();
728		let options = Options {
729			track: Some("audio".to_string()),
730			..Options::default()
731		};
732		let mut producer = Producer::new(
733			&mut broadcast,
734			catalog,
735			Input {
736				layout: Layout::Mono,
737				..Input::default()
738			},
739			&options,
740		)
741		.unwrap();
742		let mut track = moq_mux::container::Consumer::new(
743			consumer
744				.track("audio")
745				.unwrap()
746				.subscribe(moq_net::track::Subscription::default())
747				.await
748				.unwrap(),
749			moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
750		);
751
752		producer.write(&full_frame(1_000_000)).unwrap();
753		producer.reset_epoch();
754		producer.finish().unwrap();
755
756		assert!(track.read().await.unwrap().is_some());
757		assert!(track.read().await.unwrap().is_none());
758	}
759
760	/// A codec reset starts a new pre-skip interval at the receiver too.
761	#[tokio::test]
762	async fn reset_epoch_restarts_the_decoder() {
763		let input = Input {
764			format: Format::F32,
765			sample_rate: 48_000,
766			layout: Layout::Mono,
767		};
768		let options = Options {
769			track: Some("audio".to_string()),
770			settings: Settings::new(48_000, Layout::Mono),
771			..Options::default()
772		};
773		let decoder_config = Encoder::new(&options.settings).unwrap().catalog();
774
775		let mut broadcast = moq_net::broadcast::Info::new().produce();
776		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
777		let subscriber = broadcast.consume();
778		let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
779		let mut audio = AudioConsumer::new(
780			&subscriber,
781			&decoder_config,
782			"audio",
783			DecodeOptions {
784				max_age: Duration::from_millis(500),
785				..DecodeOptions::new()
786			},
787		)
788		.await
789		.unwrap();
790
791		producer.write(&full_frame(0)).unwrap();
792		let first = audio.read().await.unwrap().expect("first epoch packet");
793		assert_eq!(first.data.len() / size_of::<f32>(), 960 - 312);
794
795		producer.reset_epoch();
796		producer.write(&full_frame(1_000_000)).unwrap();
797		producer.finish().unwrap();
798
799		let mut resumed_frames = 0;
800		while let Some(frame) = audio.read().await.unwrap() {
801			if frame.timestamp.as_micros() >= 1_000_000 {
802				resumed_frames += frame.data.len() / size_of::<f32>();
803			}
804		}
805		assert!(resumed_frames > 0, "the resumed epoch still decodes");
806	}
807
808	#[tokio::test]
809	async fn producer_and_consumer_keep_activity_on_the_audio_stream() {
810		let input = Input {
811			format: Format::F32,
812			sample_rate: 48_000,
813			layout: Layout::Mono,
814		};
815		let options = Options {
816			track: Some("audio".to_string()),
817			settings: Settings {
818				layout: Layout::Mono,
819				bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
820				dtx: true,
821				..Settings::default()
822			},
823			..Options::default()
824		};
825		let decoder_config = Encoder::new(&options.settings).unwrap().catalog();
826
827		let mut broadcast = moq_net::broadcast::Info::new().produce();
828		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
829		let subscriber = broadcast.consume();
830		let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
831		let mut consumer = AudioConsumer::new(&subscriber, &decoder_config, "audio", DecodeOptions::new())
832			.await
833			.unwrap();
834
835		let silence = vec![0.0; 960];
836		let mut entered_dtx = false;
837		for index in 0..100 {
838			producer.write(&pcm_frame(&silence, index * 20_000)).unwrap();
839			let consumed = consumer.read().await.unwrap().expect("one decoded frame");
840			assert_eq!(producer.activity(), consumed.activity);
841			if consumed.activity.is_dtx() {
842				entered_dtx = true;
843				break;
844			}
845		}
846		assert!(entered_dtx, "silence should enter Opus DTX");
847
848		let active: Vec<f32> = (0..960)
849			.map(|sample| {
850				let phase = sample as f32 * 440.0 * 2.0 * std::f32::consts::PI / 48_000.0;
851				phase.sin() * 0.5
852			})
853			.collect();
854		producer.write(&pcm_frame(&active, 2_000_000)).unwrap();
855		let consumed = consumer.read().await.unwrap().expect("one decoded frame");
856		assert_eq!(producer.activity(), Activity::Active);
857		assert_eq!(consumed.activity, Activity::Active);
858	}
859
860	// One 20 ms Opus frame at 48 kHz mono is exactly 960 f32 samples, so each
861	// `write` of this drains precisely one packet (no resampler, no leftover).
862	fn full_frame(timestamp_us: u64) -> Frame {
863		pcm_frame(&vec![0.1; 960], timestamp_us)
864	}
865
866	fn pcm_frame(samples: &[f32], timestamp_us: u64) -> Frame {
867		let data: Vec<u8> = samples.iter().flat_map(|sample| sample.to_le_bytes()).collect();
868		Frame::new(Bytes::from(data), Timestamp::from_micros(timestamp_us).unwrap())
869	}
870
871	/// Publish each frame and read back the resulting packet PTS (microseconds).
872	/// If `reset_before` contains an index, `reset_epoch()` is called before that
873	/// frame's `write`.
874	async fn published_pts(frames: &[Frame], reset_before: Option<usize>) -> Vec<u128> {
875		let mut broadcast = moq_net::broadcast::Info::new().produce();
876		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
877		let consumer = broadcast.consume();
878
879		// Input rate == Opus codec rate, so there's no resampler and sample
880		// counts stay exact, making the PTS assertions deterministic.
881		let input = Input {
882			format: Format::F32,
883			sample_rate: 48_000,
884			layout: Layout::Mono,
885		};
886		let options = Options {
887			track: Some("audio".to_string()),
888			..Options::default()
889		};
890		let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
891
892		let track = consumer.track("audio").unwrap().subscribe(None).await.unwrap();
893		let mut reader =
894			moq_mux::container::Consumer::new(track, moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio));
895
896		let mut pts = Vec::new();
897		for (i, frame) in frames.iter().enumerate() {
898			if reset_before == Some(i) {
899				producer.reset_epoch();
900			}
901			producer.write(frame).unwrap();
902			let read = reader.read().await.unwrap().expect("a packet per full frame");
903			pts.push(read.timestamp.as_micros());
904		}
905		pts
906	}
907
908	#[tokio::test]
909	async fn epoch_anchors_to_first_frame_timestamp() {
910		// The first frame's timestamp becomes the epoch (regression guard: the
911		// old code derived PTS purely from the sample count, always near 0).
912		let pts = published_pts(&[full_frame(1_000_000)], None).await;
913		assert_eq!(pts, vec![1_000_000]);
914	}
915
916	/// AAC can't signal its encoder delay, so each packet is stamped that much
917	/// earlier and the first input sample still decodes at the epoch.
918	#[tokio::test]
919	async fn aac_folds_the_encoder_delay_into_timestamps() {
920		async fn pts(epoch_us: u64) -> Vec<u128> {
921			let mut broadcast = moq_net::broadcast::Info::new().produce();
922			let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
923			let consumer = broadcast.consume();
924
925			let input = Input::new(48_000, Layout::Mono);
926			let options = Options {
927				track: Some("audio".to_string()),
928				settings: Settings::from_input(crate::encode::Codec::Aac, &input),
929				..Options::default()
930			};
931			let stub = crate::encode::backend::stub::install();
932			let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
933			drop(stub);
934
935			let track = consumer
936				.track("audio")
937				.unwrap()
938				.subscribe(moq_net::track::Subscription::default().with_max_age(Duration::from_secs(1)))
939				.await
940				.unwrap();
941			let mut reader = moq_mux::container::Consumer::new(
942				track,
943				moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio),
944			);
945
946			producer.write(&pcm_frame(&[0.1; 4 * 1024], epoch_us)).unwrap();
947			let mut pts = Vec::new();
948			for _ in 0..4 {
949				pts.push(reader.read().await.unwrap().expect("a packet").timestamp.as_micros());
950			}
951			pts
952		}
953
954		// 2112 frames of delay at 48 kHz is 44 ms, rounded down per packet.
955		assert_eq!(pts(1_000_000).await, vec![956_000, 977_333, 998_666, 1_020_000]);
956		// Priming before a zero epoch stamps at zero; the input still starts on time.
957		assert_eq!(pts(0).await, vec![0, 0, 0, 20_000]);
958	}
959
960	/// The encoder needs no correction for the resampler's own delay: it anchors
961	/// the epoch to the first input timestamp and advances by emitted samples,
962	/// while `Resampler::process` drops its startup silence rather than passing it
963	/// on, so the first sample published is the first sample written. Anything
964	/// that let that delay through would shift every PTS on a resampled track,
965	/// which no `reset_epoch` would ever correct.
966	#[tokio::test]
967	async fn resampling_does_not_shift_the_first_pts() {
968		let mut broadcast = moq_net::broadcast::Info::new().produce();
969		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
970		let consumer = broadcast.consume();
971
972		// 44.1 kHz in, and Opus only runs at 48 kHz, so this one resamples.
973		let input = Input {
974			format: Format::F32,
975			sample_rate: 44_100,
976			layout: Layout::Mono,
977		};
978		let options = Options {
979			track: Some("audio".to_string()),
980			..Options::default()
981		};
982		let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
983
984		// A whole second goes out before the first read, so the default budget of zero
985		// would shed all but the last packet and the first PTS read would be the tail's.
986		let track = consumer
987			.track("audio")
988			.unwrap()
989			.subscribe(moq_net::track::Subscription::default().with_max_age(Duration::from_secs(1)))
990			.await
991			.unwrap();
992		let mut reader =
993			moq_mux::container::Consumer::new(track, moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio));
994
995		// A second of audio, so the filter's delay is nowhere near the whole write.
996		producer.write(&pcm_frame(&vec![0.1; 44_100], 1_000_000)).unwrap();
997
998		let first = reader.read().await.unwrap().expect("a packet");
999		assert_eq!(first.timestamp.as_micros(), 1_000_000);
1000	}
1001
1002	#[tokio::test]
1003	async fn pts_advances_by_frame_duration_ignoring_later_timestamps() {
1004		// Second frame's own timestamp (way ahead) is ignored; PTS advances by
1005		// exactly one 20 ms frame from the epoch.
1006		let pts = published_pts(&[full_frame(1_000), full_frame(999_999)], None).await;
1007		assert_eq!(pts, vec![1_000, 1_000 + 20_000]);
1008	}
1009
1010	#[tokio::test]
1011	async fn reset_epoch_reanchors_so_the_gap_lands_in_pts() {
1012		// Frame at t=0, then reset_epoch + a frame at t=5s: the 5 s idle gap must
1013		// appear in the PTS (otherwise audio drifts behind a wall-clock video track).
1014		let pts = published_pts(&[full_frame(0), full_frame(5_000_000)], Some(1)).await;
1015		assert_eq!(pts, vec![0, 5_000_000]);
1016	}
1017
1018	/// Finish leaves the handle, so abort can still run.
1019	#[tokio::test]
1020	async fn abort_after_finish() {
1021		let mut broadcast = moq_net::broadcast::Info::new().produce();
1022		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
1023		let consumer = broadcast.consume();
1024		let options = Options {
1025			track: Some("audio".to_string()),
1026			..Options::default()
1027		};
1028		let mut producer = Producer::new(
1029			&mut broadcast,
1030			catalog,
1031			Input {
1032				layout: Layout::Mono,
1033				..Input::default()
1034			},
1035			&options,
1036		)
1037		.unwrap();
1038		let mut track = moq_mux::container::Consumer::new(
1039			consumer
1040				.track("audio")
1041				.unwrap()
1042				.subscribe(moq_net::track::Subscription::default())
1043				.await
1044				.unwrap(),
1045			moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio),
1046		);
1047
1048		producer.write(&full_frame(0)).unwrap();
1049		producer.finish().unwrap();
1050		assert!(track.read().await.unwrap().is_some());
1051		producer.abort(moq_net::Error::Cancel);
1052	}
1053
1054	/// A sub-frame write never reaches the closed track, so the producer must
1055	/// refuse it itself rather than buffering samples that cannot be published.
1056	#[tokio::test]
1057	async fn write_after_finish_is_closed() {
1058		let mut broadcast = moq_net::broadcast::Info::new().produce();
1059		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
1060		let options = Options {
1061			track: Some("audio".to_string()),
1062			..Options::default()
1063		};
1064		let mut producer = Producer::new(
1065			&mut broadcast,
1066			catalog,
1067			Input {
1068				layout: Layout::Mono,
1069				..Input::default()
1070			},
1071			&options,
1072		)
1073		.unwrap();
1074
1075		producer.finish().unwrap();
1076		let err = producer.write(&pcm_frame(&[0.1; 100], 0)).unwrap_err();
1077		assert!(matches!(err, Error::Net(moq_net::Error::Closed)));
1078		producer.abort(moq_net::Error::Cancel);
1079	}
1080
1081	/// `Options::track = None` derives a codec-suffixed name rather than making
1082	/// the caller invent one, mirroring the video side. Pins the exact name the
1083	/// docs promise, and that a second producer doesn't collide with the first.
1084	#[tokio::test]
1085	async fn default_options_derive_the_track_name() {
1086		let mut broadcast = moq_net::broadcast::Info::new().produce();
1087		let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
1088
1089		let first = Producer::new(&mut broadcast, catalog.clone(), Input::default(), &Options::default()).unwrap();
1090		assert_eq!(first.track_name(), "0.opus");
1091
1092		let second = Producer::new(&mut broadcast, catalog, Input::default(), &Options::default()).unwrap();
1093		assert_eq!(second.track_name(), "1.opus");
1094	}
1095}