1use 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#[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 remix: Option<Remix>,
68 options: Options,
69 max_age: std::time::Duration,
70 resolved_sample_rate: u32,
71 resolved_layout: Layout,
72 next_start: Option<moq_net::Timestamp>,
76 ready: VecDeque<Frame>,
79 spans: VecDeque<ActivitySpan>,
81 trailing: Activity,
84 epoch: Option<moq_net::Timestamp>,
87 delay_trimmed: usize,
89 frames_decoded: usize,
91 end: Option<moq_net::Timestamp>,
93 terminal_start: Option<moq_net::Timestamp>,
95 discontinuity: u64,
97}
98
99struct ActivitySpan {
100 end: moq_net::Timestamp,
101 activity: Activity,
102}
103
104impl Consumer {
105 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 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 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 pub fn name(&self) -> &str {
192 self.decoder.name()
193 }
194
195 pub fn options(&self) -> &Options {
197 &self.options
198 }
199
200 pub fn max_age(&self) -> std::time::Duration {
202 self.max_age
203 }
204
205 pub fn sample_rate(&self) -> u32 {
208 self.resolved_sample_rate
209 }
210
211 pub fn layout(&self) -> Layout {
213 self.resolved_layout
214 }
215
216 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 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 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 rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch)
291 };
292 if self.end.is_some() {
293 self.frames_decoded += frames;
294 }
295 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 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 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 if pcm.is_empty() {
340 continue;
341 }
342
343 let activity = if resampled {
344 self.activity_at(timestamp)
345 } else {
346 activity
347 };
348 let frame = self.frame(pcm, timestamp, activity)?;
351 self.ready.push_back(frame);
352 }
353 }
354
355 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 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 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 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 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 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
469fn 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
499fn 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
509fn 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
516fn 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 #[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 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 #[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 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 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 #[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 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 assert!((215..=230).contains(&tail_frames), "unexpected tail: {tail_frames}");
750 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 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 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 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 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 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 #[tokio::test]
908 async fn a_missing_packet_leaves_a_hole() {
909 const FRAMES: usize = 1024;
910 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 assert_eq!(read.len(), 4, "unexpected frames: {read:?}");
920
921 let before: usize = read[..2].iter().map(|(_, frames)| frames).sum();
924 assert!((1105..=1120).contains(&before), "unexpected pre-gap audio: {before}");
925
926 assert_eq!(read[2].0, stamps[1].as_micros());
930
931 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 #[tokio::test]
943 async fn a_jump_inside_the_slack_leaves_the_held_samples_alone() {
944 const FRAMES: usize = 441;
947 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 assert_eq!(read.len(), 2, "unexpected frames: {read:?}");
958 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 #[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 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 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 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 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 #[tokio::test]
1064 async fn millisecond_stamps_are_not_a_gap() {
1065 const FRAMES: u64 = 1024;
1066 const PACKETS: u64 = 32;
1067
1068 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 assert_eq!(read.len(), stamps.len() + 1, "unexpected frames: {read:?}");
1078
1079 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 #[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 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 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 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 #[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 #[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 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 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 let frame = consumer.read().await.unwrap().expect("decoded frame");
1403 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}