1use std::collections::VecDeque;
4
5use bytes::Bytes;
6
7use super::decoder::{Config, Decoder};
8use crate::resample::{Resampler, remix, validate_remix};
9use crate::{Activity, Error, Format, Frame, Layout};
10
11#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
13#[non_exhaustive]
14pub enum Start {
15 #[default]
17 Oldest,
18 Latest,
20}
21
22#[derive(Clone, Debug, Default)]
24#[non_exhaustive]
25pub struct Output {
26 pub format: Format,
28 pub sample_rate: Option<u32>,
30 pub layout: Option<Layout>,
32}
33
34#[derive(Clone, Debug, Default)]
36#[non_exhaustive]
37pub struct Options {
38 pub decoder: Config,
40 pub output: Output,
42 pub max_age: std::time::Duration,
44 pub start: Start,
46}
47
48impl Options {
49 pub fn new() -> Self {
51 Self::default()
52 }
53}
54
55pub struct Consumer {
63 decoder: Decoder,
64 track: moq_mux::container::Consumer<moq_mux::catalog::hang::Container>,
65 resampler: Option<Resampler>,
66 options: Options,
67 max_age: std::time::Duration,
68 resolved_sample_rate: u32,
69 resolved_layout: Layout,
70 next_start: Option<moq_net::Timestamp>,
74 ready: VecDeque<Frame>,
77 spans: VecDeque<ActivitySpan>,
79 trailing: Activity,
82 epoch: Option<moq_net::Timestamp>,
85 delay_trimmed: usize,
87 frames_decoded: usize,
89 end: Option<moq_net::Timestamp>,
91 terminal_start: Option<moq_net::Timestamp>,
93 discontinuity: u64,
95}
96
97struct ActivitySpan {
98 end: moq_net::Timestamp,
99 activity: Activity,
100}
101
102impl Consumer {
103 pub async fn new(
106 broadcast: &moq_net::broadcast::Consumer,
107 catalog: &hang::catalog::AudioConfig,
108 name: impl Into<String>,
109 options: Options,
110 ) -> Result<Self, Error> {
111 let decoder = Decoder::new(catalog, &options.decoder)?;
112 let sample_rate = options.output.sample_rate.unwrap_or_else(|| decoder.sample_rate());
113 let layout = options.output.layout.unwrap_or_else(|| decoder.layout());
114 validate_remix(decoder.layout(), layout)?;
115
116 let resampler = if sample_rate == decoder.sample_rate() {
117 None
118 } else {
119 let chunk_frames = (decoder.sample_rate() as usize * 20) / 1000;
120 Some(Resampler::new(
121 decoder.sample_rate(),
122 sample_rate,
123 decoder.layout().channels(),
124 chunk_frames,
125 )?)
126 };
127
128 let name = name.into();
129 let track = broadcast.track(&name)?;
130 let mut subscriber = track
131 .subscribe(
132 moq_net::track::Subscription::default()
133 .with_priority(hang::catalog::PRIORITY.audio)
134 .with_max_age(options.max_age),
135 )
136 .await?;
137 if options.start == Start::Latest
152 && let Some(live_edge) = track.latest()
153 {
154 subscriber.set_groups(live_edge..);
155 }
156 let track = subscriber;
157 let max_age = options.max_age.min(track.info().max_age);
158 let container = moq_mux::catalog::hang::Container::try_from(catalog)?;
162 let track = moq_mux::container::Consumer::new(track, container);
163
164 Ok(Self {
165 decoder,
166 track,
167 resampler,
168 options,
169 max_age,
170 resolved_sample_rate: sample_rate,
171 resolved_layout: layout,
172 next_start: None,
173 ready: VecDeque::new(),
174 spans: VecDeque::new(),
175 trailing: Activity::Active,
176 epoch: None,
177 delay_trimmed: 0,
178 frames_decoded: 0,
179 end: None,
180 terminal_start: None,
181 discontinuity: 0,
182 })
183 }
184
185 pub fn options(&self) -> &Options {
187 &self.options
188 }
189
190 pub fn max_age(&self) -> std::time::Duration {
192 self.max_age
193 }
194
195 pub fn sample_rate(&self) -> u32 {
198 self.resolved_sample_rate
199 }
200
201 pub fn layout(&self) -> Layout {
203 self.resolved_layout
204 }
205
206 pub async fn read(&mut self) -> Result<Option<Frame>, Error> {
219 loop {
220 if let Some(frame) = self.ready.pop_front() {
221 return Ok(Some(frame));
222 }
223
224 let mux_frame = self.track.read().await?;
225 self.apply_discontinuity()?;
226 let Some(mux_frame) = mux_frame else {
227 return self.flush();
228 };
229
230 if let Some(end) = self.track.end()
231 && self.end != Some(end)
232 {
233 self.end = Some(end);
234 self.frames_decoded = 0;
235 self.terminal_start = None;
236 }
237
238 if self.end.is_none()
246 && self
247 .next_start
248 .is_some_and(|next| discontinuous(next, mux_frame.timestamp))
249 && let Some(frame) = self.gap()?
250 {
251 self.ready.push_back(frame);
252 }
253
254 let rate = self.decoder.sample_rate();
255 let epoch = *self.epoch.get_or_insert(mux_frame.timestamp);
256 let delay = self.decoder.delay_remaining();
257 let decoded = self.decoder.decode(&mux_frame.payload)?;
258 let trimmed = delay - self.decoder.delay_remaining();
261 self.delay_trimmed += trimmed;
262 let activity = decoded.activity;
263 let mut decoded = decoded.samples;
264 if let Some(end) = self.end {
265 let terminal_start = *self
266 .terminal_start
267 .get_or_insert(rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch));
268 let total = frames_between(terminal_start, end, rate)?;
269 let remaining = total.saturating_sub(self.frames_decoded);
270 decoded.truncate(remaining.saturating_mul(self.decoder.layout().channels() as usize));
271 }
272
273 let frames = decoded.len() / self.decoder.layout().channels() as usize;
274 let decoded_at = if let Some(terminal_start) = self.terminal_start {
275 advance(terminal_start, self.frames_decoded, rate)?
276 } else {
277 rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch)
281 };
282 if self.end.is_some() {
283 self.frames_decoded += frames;
284 }
285 self.next_start = Some(advance(mux_frame.timestamp, frames + trimmed, rate)?);
288 if decoded.is_empty() {
289 continue;
290 }
291
292 let (pcm, timestamp) = match self.resampler.as_mut() {
293 Some(r) => {
300 let held = if r.pending_frames() == 0 {
301 decoded_at
302 } else {
303 r.held_at().unwrap_or(decoded_at)
304 };
305 let skipped = r.skipped();
306 let pcm = r.process(&decoded, decoded_at)?;
307 (pcm, rewind(held, skipped, self.resolved_sample_rate)?)
308 }
309 None => (decoded, decoded_at),
310 };
311
312 let decoded_end = advance(decoded_at, frames, rate)?;
313
314 let resampled = self.resampler.is_some();
319 if resampled {
320 self.spans.push_back(ActivitySpan {
321 end: decoded_end,
322 activity,
323 });
324 }
325
326 if pcm.is_empty() {
330 continue;
331 }
332
333 let activity = if resampled {
334 self.activity_at(timestamp)
335 } else {
336 activity
337 };
338 let frame = self.frame(pcm, timestamp, activity)?;
341 self.ready.push_back(frame);
342 }
343 }
344
345 fn apply_discontinuity(&mut self) -> Result<(), Error> {
348 let discontinuity = self.track.discontinuity();
349 if discontinuity == self.discontinuity {
350 return Ok(());
351 }
352
353 self.discontinuity = discontinuity;
354 self.next_start = None;
355 self.spans.clear();
356 self.trailing = Activity::Active;
357 self.frames_decoded = 0;
358 self.end = None;
359 self.terminal_start = None;
360 self.epoch = None;
361 self.delay_trimmed = 0;
362 self.decoder.reapply_delay();
363 Ok(())
364 }
365
366 fn gap(&mut self) -> Result<Option<Frame>, Error> {
375 self.decoder.reset_prediction()?;
376
377 let mut frame = None;
378 if let Some(resampler) = self.resampler.as_mut() {
379 let held = resampler.held_at();
380 let skipped = resampler.skipped();
381 let pcm = resampler.drain()?;
382 frame = self.tail(pcm, held, skipped)?;
383 }
384
385 self.next_start = None;
386 self.spans.clear();
387 self.trailing = Activity::Active;
388 self.epoch = None;
389 self.delay_trimmed = 0;
390 Ok(frame)
391 }
392
393 fn flush(&mut self) -> Result<Option<Frame>, Error> {
400 let Some(resampler) = self.resampler.take() else {
401 return Ok(None);
402 };
403
404 let held = resampler.held_at();
405 let skipped = resampler.skipped();
406 self.tail(resampler.flush()?, held, skipped)
407 }
408
409 fn tail(
415 &mut self,
416 pcm: Vec<f32>,
417 held: Option<moq_net::Timestamp>,
418 skipped: usize,
419 ) -> Result<Option<Frame>, Error> {
420 let Some(held) = held.filter(|_| !pcm.is_empty()) else {
421 return Ok(None);
422 };
423
424 let timestamp = rewind(held, skipped, self.resolved_sample_rate)?;
425 let activity = self.activity_at(timestamp);
426 Ok(Some(self.frame(pcm, timestamp, activity)?))
427 }
428
429 fn activity_at(&mut self, timestamp: moq_net::Timestamp) -> Activity {
431 while let Some(span) = self.spans.front().filter(|span| span.end <= timestamp) {
432 self.trailing = span.activity;
433 self.spans.pop_front();
434 }
435
436 self.spans.front().map_or(self.trailing, |span| span.activity)
437 }
438
439 fn frame(&self, pcm: Vec<f32>, timestamp: moq_net::Timestamp, activity: Activity) -> Result<Frame, Error> {
441 let pcm = if self.decoder.layout() == self.resolved_layout {
442 pcm
443 } else {
444 remix(&pcm, self.decoder.layout(), self.resolved_layout)?
445 };
446
447 let bytes = self
448 .options
449 .output
450 .format
451 .from_interleaved_f32(&pcm, self.resolved_layout.channels())?;
452 Ok(Frame {
453 timestamp,
454 data: Bytes::from(bytes),
455 activity,
456 })
457 }
458}
459
460fn discontinuous(expected: moq_net::Timestamp, timestamp: moq_net::Timestamp) -> bool {
484 let scale = expected.scale().max(timestamp.scale());
485 let quantum = scale.min(moq_net::Timescale::default());
486 let tolerance = (scale.as_u64() as u128).div_ceil(quantum.as_u64() as u128) + 1;
487 expected.as_scale(scale).abs_diff(timestamp.as_scale(scale)) > tolerance
488}
489
490fn advance(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
492 if frames == 0 {
493 return Ok(timestamp);
494 }
495
496 let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
497 Ok(timestamp.checked_add(offset)?)
498}
499
500fn frames_between(start: moq_net::Timestamp, end: moq_net::Timestamp, sample_rate: u32) -> Result<usize, Error> {
502 let duration = end.checked_sub(start)?;
503 let frames = (std::time::Duration::from(duration).as_nanos() * sample_rate as u128 + 500_000_000) / 1_000_000_000;
504 usize::try_from(frames).map_err(|_| Error::Unsupported("audio duration does not fit in memory".into()))
505}
506
507fn rewind(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
512 if frames == 0 {
513 return Ok(timestamp);
514 }
515
516 let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
517 Ok(timestamp
518 .checked_sub(offset)
519 .unwrap_or(moq_net::Timestamp::new(0, timestamp.scale())?))
520}
521
522#[cfg(test)]
523mod tests {
524 use moq_net::Timestamp;
525
526 use super::*;
527 use crate::encode::{Encoder, Input, Options as EncodeOptions, Producer, Settings};
528 use crate::{Format, Layout};
529
530 #[tokio::test]
531 async fn remixes_mono_stream_to_stereo_output() {
532 let mut broadcast = moq_net::broadcast::Info::new().produce();
533 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
534 let subscriber = broadcast.consume();
535 let input = Input {
536 format: Format::F32,
537 sample_rate: 48_000,
538 layout: Layout::Mono,
539 };
540 let options = EncodeOptions {
541 track: Some("audio".to_string()),
542 settings: Settings::new(48_000, Layout::Mono),
543 ..EncodeOptions::default()
544 };
545 let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
546 let catalog = Encoder::new(&Settings::new(input.sample_rate, input.layout))
547 .unwrap()
548 .catalog();
549 let mut consumer = Consumer::new(
550 &subscriber,
551 &catalog,
552 "audio",
553 Options {
554 output: Output {
555 layout: Some(Layout::Stereo),
556 ..Output::default()
557 },
558 ..Options::new()
559 },
560 )
561 .await
562 .unwrap();
563
564 let samples = vec![0.1f32; 960];
565 let mut data = Vec::with_capacity(samples.len() * size_of::<f32>());
566 for sample in samples {
567 data.extend_from_slice(&sample.to_le_bytes());
568 }
569 producer.write(&Frame::new(data.into(), Timestamp::ZERO)).unwrap();
570
571 let frame = consumer.read().await.unwrap().expect("decoded frame");
572 let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
573 assert_eq!(samples.len(), (960 - 312) * 2);
574 for pair in samples.as_chunks::<2>().0.iter() {
575 assert_eq!(pair[0], pair[1]);
576 }
577 }
578
579 #[tokio::test]
583 async fn opus_timestamps_follow_the_48k_clock() {
584 use crate::decode::decoder::tests::{opus_catalog, opus_packets};
585
586 let broadcast = moq_net::broadcast::Info::new().produce();
587 let track = broadcast
588 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
589 .unwrap();
590 let subscriber = broadcast.consume();
591
592 let catalog = opus_catalog(moq_mux::codec::opus::Config::new(44_100, 1).with_pre_skip(312));
593 let mut producer = moq_mux::container::Producer::new(
594 track,
595 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
596 );
597 let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Options::new())
598 .await
599 .unwrap();
600 assert_eq!(consumer.sample_rate(), 48_000);
601
602 for (packet, payload) in opus_packets(3).into_iter().enumerate() {
603 producer
604 .write(moq_mux::container::Frame {
605 timestamp: Timestamp::from_micros(packet as u64 * 20_000).unwrap(),
606 duration: None,
607 payload,
608 keyframe: packet == 0,
609 })
610 .unwrap();
611 }
612
613 for (micros, frames) in [(0, 960 - 312), (13_500, 960), (33_500, 960)] {
615 let frame = consumer.read().await.unwrap().expect("decoded frame");
616 assert_eq!(frame.timestamp.as_micros(), micros);
617 assert_eq!(frame.data.len() / size_of::<f32>(), frames);
618 }
619 }
620
621 #[tokio::test]
628 async fn resampled_timestamps_follow_the_samples() {
629 let broadcast = moq_net::broadcast::Info::new().produce();
630 let track = broadcast
631 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
632 .unwrap();
633 let subscriber = broadcast.consume();
634
635 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
636 let mut producer = moq_mux::container::Producer::new(
637 track,
638 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
639 );
640
641 let mut consumer = Consumer::new(
642 &subscriber,
643 &catalog,
644 "audio",
645 Options {
646 output: Output {
647 sample_rate: Some(48_000),
648 ..Output::default()
649 },
650 max_age: std::time::Duration::from_secs(1),
651 ..Options::new()
652 },
653 )
654 .await
655 .unwrap();
656
657 const FRAMES: u64 = 1024;
659 let payload: Bytes = vec![0u8; FRAMES as usize * size_of::<f32>()].into();
660 for packet in 0..2 {
661 producer
662 .write(moq_mux::container::Frame {
663 timestamp: moq_net::Timestamp::from_scale(packet * FRAMES, 44_100).unwrap(),
664 duration: None,
665 payload: payload.clone(),
666 keyframe: true,
667 })
668 .unwrap();
669 }
670
671 let first = consumer.read().await.unwrap().expect("decoded frame");
672 assert_eq!(first.timestamp.as_micros(), 0);
673
674 let second = consumer.read().await.unwrap().expect("decoded frame");
680 let first_frames = (first.data.len() / size_of::<f32>()) as u128;
681 let ends_at = first_frames * 1_000_000 / 48_000;
682 let gap = second.timestamp.as_micros().abs_diff(ends_at);
683 assert!(gap < 100, "expected the frames to meet, got a {gap} us gap");
684 }
685
686 #[tokio::test]
690 async fn resampled_tail_survives_the_end_of_the_track() {
691 let broadcast = moq_net::broadcast::Info::new().produce();
692 let track = broadcast
693 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
694 .unwrap();
695 let subscriber = broadcast.consume();
696
697 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
698 let mut producer = moq_mux::container::Producer::new(
699 track,
700 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
701 );
702
703 let mut consumer = Consumer::new(
704 &subscriber,
705 &catalog,
706 "audio",
707 Options {
708 output: Output {
709 sample_rate: Some(48_000),
710 ..Output::default()
711 },
712 ..Options::new()
713 },
714 )
715 .await
716 .unwrap();
717
718 const FRAMES: usize = 1024;
720 let payload: Bytes = vec![0u8; FRAMES * size_of::<f32>()].into();
721 producer
722 .write(moq_mux::container::Frame {
723 timestamp: moq_net::Timestamp::ZERO,
724 duration: None,
725 payload,
726 keyframe: true,
727 })
728 .unwrap();
729 producer.finish().unwrap();
730
731 let first = consumer.read().await.unwrap().expect("decoded frame");
732 let first_frames = first.data.len() / size_of::<f32>();
733
734 let tail = consumer.read().await.unwrap().expect("flushed tail");
735 let tail_frames = tail.data.len() / size_of::<f32>();
736
737 assert!((215..=230).contains(&tail_frames), "unexpected tail: {tail_frames}");
741 let ends_at = (first_frames as u128) * 1_000_000 / 48_000;
744 let gap = tail.timestamp.as_micros().abs_diff(ends_at);
745 assert!(gap < 100, "expected the tail to meet the body, got a {gap} us gap");
746
747 let total = first_frames + tail_frames;
751 assert!((1105..=1120).contains(&total), "unexpected total: {total}");
752 assert!(consumer.read().await.unwrap().is_none());
753 }
754
755 #[tokio::test]
756 async fn resampling_keeps_the_activity_boundary_on_its_source() {
757 let mut encoder = Encoder::new(&Settings {
758 dtx: true,
759 bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
760 frame_duration: std::time::Duration::from_millis(10),
761 ..Settings::new(48_000, Layout::Mono)
762 })
763 .unwrap();
764 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Opus, 48_000, 1);
765
766 let broadcast = moq_net::broadcast::Info::new().produce();
767 let track = broadcast
768 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
769 .unwrap();
770 let subscriber = broadcast.consume();
771 let mut producer = moq_mux::container::Producer::new(
772 track,
773 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
774 );
775 let mut consumer = Consumer::new(
776 &subscriber,
777 &catalog,
778 "audio",
779 Options {
780 output: Output {
781 sample_rate: Some(44_100),
782 ..Output::default()
783 },
784 max_age: std::time::Duration::from_secs(1),
785 ..Options::new()
786 },
787 )
788 .await
789 .unwrap();
790
791 let active = vec![0.5; encoder.frame_size()];
792 let silence = vec![0.0; encoder.frame_size()];
793 let mut first_dtx = None;
794 for index in 0..40u64 {
795 let packet = encoder.encode(if index == 0 { &active } else { &silence }).unwrap();
796 let timestamp = Timestamp::from_scale(index * encoder.frame_size() as u64, 48_000).unwrap();
797 if first_dtx.is_none() && packet.activity.is_dtx() {
798 first_dtx = Some(timestamp);
799 }
800 producer
801 .write(moq_mux::container::Frame {
802 timestamp,
803 payload: packet.payload,
804 keyframe: true,
805 duration: None,
806 })
807 .unwrap();
808 producer.cut(None).unwrap();
809 }
810 producer.finish().unwrap();
811
812 let expected = first_dtx.expect("silence should enter Opus DTX");
813 let mut actual = None;
814 while let Some(frame) = consumer.read().await.unwrap() {
815 assert!(!frame.data.is_empty(), "read returned a frame with no samples");
820 if frame.activity.is_dtx() {
821 actual = Some(frame.timestamp);
822 break;
823 }
824 }
825 let actual = actual.expect("consumer should report Opus DTX");
826
827 let delay = actual.as_micros() as i128 - expected.as_micros() as i128;
834 let chunk_us = 20_000i128;
835 assert!(
836 (0..chunk_us).contains(&delay),
837 "DTX label landed {delay} us from its source, outside [0, {chunk_us})"
838 );
839 }
840
841 async fn pcm_gaps(rate: u32, out_rate: u32, frames: usize, stamps: &[Timestamp]) -> Vec<(u128, usize)> {
844 let broadcast = moq_net::broadcast::Info::new().produce();
845 let track = broadcast
846 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
847 .unwrap();
848 let subscriber = broadcast.consume();
849
850 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, rate, 1);
851 let mut producer = moq_mux::container::Producer::new(
852 track,
853 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
854 );
855 let mut consumer = Consumer::new(
856 &subscriber,
857 &catalog,
858 "audio",
859 Options {
863 output: Output {
864 sample_rate: Some(out_rate),
865 ..Output::default()
866 },
867 max_age: std::time::Duration::from_secs(1),
868 ..Options::new()
869 },
870 )
871 .await
872 .unwrap();
873
874 let payload: Bytes = vec![0u8; frames * size_of::<f32>()].into();
875 for stamp in stamps {
876 producer
877 .write(moq_mux::container::Frame {
878 timestamp: *stamp,
879 duration: None,
880 payload: payload.clone(),
881 keyframe: true,
882 })
883 .unwrap();
884 }
885 producer.finish().unwrap();
886
887 let mut read = Vec::new();
888 while let Some(frame) = consumer.read().await.unwrap() {
889 read.push((frame.timestamp.as_micros(), frame.data.len() / size_of::<f32>()));
890 }
891 read
892 }
893
894 #[tokio::test]
899 async fn a_missing_packet_leaves_a_hole() {
900 const FRAMES: usize = 1024;
901 let stamps = [
903 Timestamp::from_scale(0, 44_100).unwrap(),
904 Timestamp::from_scale(2 * FRAMES as u64, 44_100).unwrap(),
905 ];
906 let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
907
908 assert_eq!(read.len(), 4, "unexpected frames: {read:?}");
911
912 let before: usize = read[..2].iter().map(|(_, frames)| frames).sum();
915 assert!((1105..=1120).contains(&before), "unexpected pre-gap audio: {before}");
916
917 assert_eq!(read[2].0, stamps[1].as_micros());
921
922 let ends_at = read[1].0 + (read[1].1 as u128) * 1_000_000 / 48_000;
924 let hole = read[2].0 - ends_at;
925 assert!((23_100..=23_350).contains(&hole), "unexpected hole: {hole} us");
926 }
927
928 #[tokio::test]
934 async fn a_jump_inside_the_slack_leaves_the_held_samples_alone() {
935 const FRAMES: usize = 441;
938 let stamps = [
941 Timestamp::from_micros(0).unwrap(),
942 Timestamp::from_micros(11_000).unwrap(),
943 ];
944 let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
945
946 assert_eq!(read.len(), 2, "unexpected frames: {read:?}");
949 assert_eq!(read[0].0, 0, "held samples moved with the jump: {read:?}");
952 }
953
954 #[tokio::test]
955 async fn a_jump_after_a_full_chunk_uses_the_new_packet_timestamp() {
956 let stamps = [
957 Timestamp::from_micros(0).unwrap(),
958 Timestamp::from_micros(21_000).unwrap(),
959 ];
960 let read = pcm_gaps(44_100, 48_000, 882, &stamps).await;
961 let mut r = crate::resample::Resampler::new(44_100, 48_000, 1, 882).unwrap();
962 r.process(&[0.25; 882], stamps[0]).unwrap();
963 let expected = rewind(stamps[1], r.skipped(), 48_000).unwrap().as_micros();
964 assert_eq!(read[1].0, expected);
965 }
966
967 #[tokio::test]
973 async fn a_terminal_jump_leaves_the_held_samples_alone() {
974 let mut encoder = Encoder::new(&Settings {
975 dtx: true,
976 bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
977 ..Settings::new(48_000, Layout::Mono)
978 })
979 .unwrap();
980 let catalog = encoder.catalog();
981
982 let active = encoder.encode(&vec![0.5f32; encoder.frame_size()]).unwrap();
986 assert!(active.activity.is_active());
987 let silence = vec![0.0f32; encoder.frame_size()];
988 let dtx = (0..200)
989 .map(|_| encoder.encode(&silence).unwrap())
990 .find(|packet| packet.activity.is_dtx())
991 .expect("silence should enter Opus DTX");
992
993 let broadcast = moq_net::broadcast::Info::new().produce();
994 let track = broadcast
995 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
996 .unwrap();
997 let subscriber = broadcast.consume();
998 let mut producer = moq_mux::container::Producer::new(
999 track,
1000 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1001 );
1002 let mut consumer = Consumer::new(
1003 &subscriber,
1004 &catalog,
1005 "audio",
1006 Options {
1007 output: Output {
1008 sample_rate: Some(44_100),
1009 ..Output::default()
1010 },
1011 ..Options::new()
1012 },
1013 )
1014 .await
1015 .unwrap();
1016
1017 let write = |producer: &mut moq_mux::container::Producer<_>, frames: u64, payload: Bytes, keyframe: bool| {
1020 producer
1021 .write(moq_mux::container::Frame {
1022 timestamp: Timestamp::from_scale(frames, 48_000).unwrap(),
1023 duration: None,
1024 payload,
1025 keyframe,
1026 })
1027 .unwrap();
1028 };
1029 write(&mut producer, 0, active.payload, true);
1030 write(&mut producer, 3 * 48_000, Bytes::new(), false);
1033 write(&mut producer, 48_000, dtx.payload, false);
1034 producer.finish().unwrap();
1035
1036 let frame = consumer.read().await.unwrap().expect("decoded frame");
1037 assert_eq!(frame.timestamp.as_micros(), 0, "held samples moved with the jump");
1041 assert!(
1042 frame.activity.is_active(),
1043 "held samples took the terminal packet's label"
1044 );
1045 }
1046
1047 #[tokio::test]
1055 async fn millisecond_stamps_are_not_a_gap() {
1056 const FRAMES: u64 = 1024;
1057 const PACKETS: u64 = 32;
1058
1059 let stamps: Vec<_> = (0..PACKETS)
1061 .map(|packet| Timestamp::from_millis(packet * FRAMES * 1_000 / 44_100).unwrap())
1062 .collect();
1063 let read = pcm_gaps(44_100, 48_000, FRAMES as usize, &stamps).await;
1064
1065 assert_eq!(read.len(), stamps.len() + 1, "unexpected frames: {read:?}");
1069
1070 for pair in read.windows(2) {
1073 let ends_at = pair[0].0 + (pair[0].1 as u128) * 1_000_000 / 48_000;
1074 assert!(
1075 pair[1].0.abs_diff(ends_at) <= 1_100,
1076 "frames at {} and {} do not meet",
1077 pair[0].0,
1078 pair[1].0
1079 );
1080 }
1081 }
1082
1083 #[tokio::test]
1087 async fn a_lost_opus_packet_shorter_than_its_neighbour_is_a_gap() {
1088 let input = Input {
1089 format: Format::F32,
1090 sample_rate: 48_000,
1091 layout: Layout::Mono,
1092 };
1093 let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1094 let catalog = encoder.catalog();
1095
1096 let broadcast = moq_net::broadcast::Info::new().produce();
1097 let track = broadcast
1098 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1099 .unwrap();
1100 let subscriber = broadcast.consume();
1101 let mut producer = moq_mux::container::Producer::new(
1102 track,
1103 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1104 );
1105 let mut consumer = Consumer::new(
1106 &subscriber,
1107 &catalog,
1108 "audio",
1109 Options {
1110 max_age: std::time::Duration::from_secs(1),
1111 ..Options::new()
1112 },
1113 )
1114 .await
1115 .unwrap();
1116
1117 let pcm = vec![0.25f32; encoder.frame_size()];
1120 for timestamp in [
1121 Timestamp::from_micros(0).unwrap(),
1122 Timestamp::from_micros(22_500).unwrap(),
1123 Timestamp::from_micros(42_500).unwrap(),
1124 ] {
1125 producer
1126 .write(moq_mux::container::Frame {
1127 timestamp,
1128 duration: None,
1129 payload: encoder.encode(&pcm).unwrap().payload,
1130 keyframe: true,
1131 })
1132 .unwrap();
1133 producer.cut(None).unwrap();
1134 }
1135
1136 let first = consumer.read().await.unwrap().expect("decoded frame");
1140 let frames = first.data.len() / size_of::<f32>();
1141 assert!(frames < 960, "the pre-skip should be trimmed, got {frames} frames");
1142
1143 let second = consumer.read().await.unwrap().expect("decoded frame");
1146 assert_eq!(second.timestamp.as_micros(), 22_500);
1147 assert_eq!(second.data.len() / size_of::<f32>(), 960, "pre-skip was reapplied");
1148
1149 let third = consumer.read().await.unwrap().expect("decoded frame after gap");
1150 let second_frames = second.data.len() / size_of::<f32>();
1151 assert_eq!(
1152 third.timestamp,
1153 advance(second.timestamp, second_frames, 48_000).unwrap()
1154 );
1155 }
1156
1157 #[tokio::test]
1158 async fn max_age_is_clamped_to_publisher_retention() {
1159 let broadcast = moq_net::broadcast::Info::new().produce();
1160 let info = hang::container::track_info(hang::catalog::PRIORITY.audio)
1161 .with_max_age(std::time::Duration::from_millis(100));
1162 let _track = broadcast.create_track("audio", info).unwrap();
1163 let subscriber = broadcast.consume();
1164 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
1165
1166 let consumer = Consumer::new(
1167 &subscriber,
1168 &catalog,
1169 "audio",
1170 Options {
1171 max_age: std::time::Duration::from_millis(500),
1172 ..Options::new()
1173 },
1174 )
1175 .await
1176 .unwrap();
1177
1178 assert_eq!(consumer.max_age(), std::time::Duration::from_millis(100));
1179 }
1180
1181 #[tokio::test]
1185 async fn opus_pre_skip_does_not_leave_a_timestamp_hole() {
1186 let input = Input {
1187 format: Format::F32,
1188 sample_rate: 48_000,
1189 layout: Layout::Mono,
1190 };
1191 let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1192 let catalog = encoder.catalog();
1193
1194 let broadcast = moq_net::broadcast::Info::new().produce();
1195 let track = broadcast
1196 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1197 .unwrap();
1198 let subscriber = broadcast.consume();
1199 let mut producer = moq_mux::container::Producer::new(
1200 track,
1201 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1202 );
1203 let mut consumer = Consumer::new(
1204 &subscriber,
1205 &catalog,
1206 "audio",
1207 Options {
1208 max_age: std::time::Duration::from_secs(1),
1209 ..Options::new()
1210 },
1211 )
1212 .await
1213 .unwrap();
1214
1215 let pcm = vec![0.25f32; encoder.frame_size()];
1216 for packet in 0..2 {
1217 producer
1218 .write(moq_mux::container::Frame {
1219 timestamp: Timestamp::from_scale(packet * encoder.frame_size() as u64, 48_000).unwrap(),
1220 duration: None,
1221 payload: encoder.encode(&pcm).unwrap().payload,
1222 keyframe: true,
1223 })
1224 .unwrap();
1225 producer.cut(None).unwrap();
1226 }
1227
1228 let first = consumer.read().await.unwrap().expect("first decoded frame");
1229 let second = consumer.read().await.unwrap().expect("second decoded frame");
1230 let first_frames = first.data.len() / size_of::<f32>();
1231 let expected = advance(first.timestamp, first_frames, 48_000).unwrap();
1232 assert_eq!(second.timestamp, expected);
1233 }
1234
1235 #[tokio::test]
1236 async fn a_playhead_event_reapplies_opus_pre_skip() {
1237 let input = Input {
1238 format: Format::F32,
1239 sample_rate: 48_000,
1240 layout: Layout::Mono,
1241 };
1242 let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1243 let catalog = encoder.catalog();
1244 let frame_size = encoder.frame_size();
1245
1246 let broadcast = moq_net::broadcast::Info::new().produce();
1247 let track = broadcast
1248 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1249 .unwrap();
1250 let subscriber = broadcast.consume();
1251 let mut producer = moq_mux::container::Producer::new(
1252 track,
1253 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
1254 );
1255 let mut consumer = Consumer::new(
1256 &subscriber,
1257 &catalog,
1258 "audio",
1259 Options {
1260 max_age: std::time::Duration::from_secs(1),
1261 ..Options::new()
1262 },
1263 )
1264 .await
1265 .unwrap();
1266
1267 let pcm = vec![0.25f32; frame_size];
1268 let write = |producer: &mut moq_mux::container::Producer<_>, packet: u64, payload: bytes::Bytes| {
1269 producer
1270 .write(moq_mux::container::Frame {
1271 timestamp: Timestamp::from_scale(packet * frame_size as u64, 48_000).unwrap(),
1272 duration: None,
1273 payload,
1274 keyframe: true,
1275 })
1276 .unwrap();
1277 producer.cut(None).unwrap();
1278 };
1279 write(&mut producer, 0, encoder.encode(&pcm).unwrap().payload);
1280 write(&mut producer, 1, encoder.encode(&pcm).unwrap().payload);
1281 producer.discontinuity().unwrap();
1282 write(&mut producer, 2, encoder.encode(&pcm).unwrap().payload);
1283 producer.finish().unwrap();
1284
1285 let first = consumer.read().await.unwrap().expect("first decoded frame");
1286 let _second = consumer.read().await.unwrap().expect("second decoded frame");
1287 let resumed = consumer.read().await.unwrap().expect("resumed decoded frame");
1288 let first_frames = first.data.len() / size_of::<f32>();
1289 let resumed_frames = resumed.data.len() / size_of::<f32>();
1290 assert!(first_frames < frame_size, "the first epoch trims pre-skip");
1291 assert_eq!(
1292 resumed_frames, first_frames,
1293 "a playhead event reapplies pre-skip without flushing the decoder"
1294 );
1295 }
1296
1297 #[tokio::test]
1298 async fn reads_the_container_the_catalog_declares() {
1299 let broadcast = moq_net::broadcast::Info::new().produce();
1300 let track = broadcast
1301 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1302 .unwrap();
1303 let observed = track.clone();
1304 let subscriber = broadcast.consume();
1305
1306 let mut catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
1307 catalog.container = hang::catalog::Container::Loc;
1308
1309 let mut producer = moq_mux::container::Producer::new(
1310 track,
1311 moq_mux::catalog::hang::Container::Loc(moq_mux::container::Kind::Audio),
1312 );
1313 let max_age = std::time::Duration::from_millis(250);
1314 let mut consumer = Consumer::new(
1315 &subscriber,
1316 &catalog,
1317 "audio",
1318 Options {
1319 max_age,
1320 ..Options::new()
1321 },
1322 )
1323 .await
1324 .unwrap();
1325 assert_eq!(observed.subscription().unwrap().max_age, max_age);
1326
1327 let samples = [0.25f32, -0.5, 0.75, -1.0];
1328 let payload: Vec<u8> = samples.iter().flat_map(|sample| sample.to_le_bytes()).collect();
1329 producer
1330 .write(moq_mux::container::Frame {
1331 timestamp: Timestamp::ZERO,
1332 duration: None,
1333 payload: payload.into(),
1334 keyframe: true,
1335 })
1336 .unwrap();
1337
1338 let frame = consumer.read().await.unwrap().expect("decoded frame");
1339 assert_eq!(
1340 Format::F32.as_interleaved_f32(&frame.data, 1).unwrap().as_ref(),
1341 samples
1342 );
1343 }
1344
1345 #[tokio::test]
1350 async fn decodes_a_cmaf_framed_track() {
1351 let input = Input {
1352 format: Format::F32,
1353 sample_rate: 48_000,
1354 layout: Layout::Stereo,
1355 };
1356
1357 let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
1359 let mut catalog = encoder.catalog();
1360 let pcm = vec![0.0f32; encoder.frame_size() * encoder.codec_channels() as usize];
1361 let packet = encoder.encode(&pcm).unwrap();
1362
1363 let muxer = moq_mux::container::fmp4::Muxer::audio(&catalog).unwrap();
1365 let init = muxer.init().unwrap().expect("an out-of-band codec has an init segment");
1366 catalog.container = hang::catalog::Container::Cmaf { init };
1367
1368 let broadcast = moq_net::broadcast::Info::new().produce();
1369 let subscriber = broadcast.consume();
1370 let track = broadcast
1371 .create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
1372 .unwrap();
1373 let container = moq_mux::catalog::hang::Container::try_from(&catalog).unwrap();
1374 let mut producer = moq_mux::container::Producer::new(track, container);
1375
1376 let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Options::new())
1377 .await
1378 .unwrap();
1379
1380 producer
1381 .write(moq_mux::container::Frame {
1382 timestamp: Timestamp::ZERO,
1383 payload: packet.payload,
1384 keyframe: true,
1385 duration: None,
1386 })
1387 .unwrap();
1388 producer.cut(None).unwrap();
1389
1390 let frame = consumer.read().await.unwrap().expect("decoded frame");
1394 assert_eq!(frame.timestamp.as_micros(), 0);
1397 let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
1398 assert_eq!(samples.len(), (960 - 312) * 2);
1399 }
1400}