1use std::collections::VecDeque;
4
5use bytes::Bytes;
6
7use super::decoder::{Config, Decoder};
8use crate::resample::{Resampler, remix, validate_channels};
9use crate::{Activity, Error, Frame};
10
11pub struct Consumer {
19 decoder: Decoder,
20 track: moq_mux::container::Consumer<moq_mux::catalog::hang::Container>,
21 resampler: Option<Resampler>,
22 config: Config,
23 latency_max: std::time::Duration,
24 resolved_sample_rate: u32,
25 resolved_channels: u32,
26 next_start: Option<moq_net::Timestamp>,
30 ready: VecDeque<Frame>,
33 spans: VecDeque<ActivitySpan>,
35 trailing: Activity,
38 epoch: Option<moq_net::Timestamp>,
41 delay_trimmed: usize,
43 frames_decoded: usize,
45 end: Option<moq_net::Timestamp>,
47 terminal_start: Option<moq_net::Timestamp>,
49 discontinuity: u64,
51}
52
53struct ActivitySpan {
54 end: moq_net::Timestamp,
55 activity: Activity,
56}
57
58impl Consumer {
59 pub async fn new(
62 broadcast: &moq_net::broadcast::Consumer,
63 catalog: &hang::catalog::AudioConfig,
64 name: impl Into<String>,
65 config: Config,
66 ) -> Result<Self, Error> {
67 let decoder = Decoder::new(catalog)?;
68 let sample_rate = config.sample_rate.unwrap_or_else(|| decoder.sample_rate());
69 let channels = config.channels.unwrap_or_else(|| decoder.channel_count());
70 validate_channels(channels)?;
71
72 let resampler = if sample_rate == decoder.sample_rate() {
73 None
74 } else {
75 let chunk_frames = (decoder.sample_rate() as usize * 20) / 1000;
76 Some(Resampler::new(
77 decoder.sample_rate(),
78 sample_rate,
79 decoder.channel_count(),
80 chunk_frames,
81 )?)
82 };
83
84 let name = name.into();
85 let track = broadcast
86 .track(&name)?
87 .subscribe(moq_net::track::Subscription::default().with_priority(hang::catalog::PRIORITY.audio))
88 .await?;
89 let latency_max = config.latency_max.unwrap_or_default().min(track.info().latency_max);
90 let container = moq_mux::catalog::hang::Container::try_from(&catalog.container)?;
94 let mut track = moq_mux::container::Consumer::new(track, container);
95 if let Some(latency) = config.latency_max {
96 track = track.with_latency(latency);
97 }
98
99 Ok(Self {
100 decoder,
101 track,
102 resampler,
103 config,
104 latency_max,
105 resolved_sample_rate: sample_rate,
106 resolved_channels: channels,
107 next_start: None,
108 ready: VecDeque::new(),
109 spans: VecDeque::new(),
110 trailing: Activity::Active,
111 epoch: None,
112 delay_trimmed: 0,
113 frames_decoded: 0,
114 end: None,
115 terminal_start: None,
116 discontinuity: 0,
117 })
118 }
119
120 pub fn config(&self) -> &Config {
122 &self.config
123 }
124
125 pub fn latency_max(&self) -> std::time::Duration {
127 self.latency_max
128 }
129
130 pub fn sample_rate(&self) -> u32 {
133 self.resolved_sample_rate
134 }
135
136 pub fn channels(&self) -> u32 {
139 self.resolved_channels
140 }
141
142 pub async fn read(&mut self) -> Result<Option<Frame>, Error> {
155 loop {
156 if let Some(frame) = self.ready.pop_front() {
157 return Ok(Some(frame));
158 }
159
160 let mux_frame = self.track.read().await?;
161 self.apply_discontinuity()?;
162 let Some(mux_frame) = mux_frame else {
163 return self.flush();
164 };
165
166 if let Some(end) = self.track.end()
167 && self.end != Some(end)
168 {
169 self.end = Some(end);
170 self.frames_decoded = 0;
171 self.terminal_start = None;
172 }
173
174 if self.end.is_none()
182 && self
183 .next_start
184 .is_some_and(|next| discontinuous(next, mux_frame.timestamp))
185 && let Some(frame) = self.gap()?
186 {
187 self.ready.push_back(frame);
188 }
189
190 let rate = self.decoder.sample_rate();
191 let epoch = *self.epoch.get_or_insert(mux_frame.timestamp);
192 let delay = self.decoder.delay_remaining();
193 let decoded = self.decoder.decode(&mux_frame.payload)?;
194 let trimmed = delay - self.decoder.delay_remaining();
197 self.delay_trimmed += trimmed;
198 let activity = decoded.activity;
199 let mut decoded = decoded.samples;
200 if let Some(end) = self.end {
201 let terminal_start = *self
202 .terminal_start
203 .get_or_insert(rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch));
204 let total = frames_between(terminal_start, end, rate)?;
205 let remaining = total.saturating_sub(self.frames_decoded);
206 decoded.truncate(remaining.saturating_mul(self.decoder.channel_count() as usize));
207 }
208
209 let frames = decoded.len() / self.decoder.channel_count().max(1) as usize;
210 let decoded_at = if let Some(terminal_start) = self.terminal_start {
211 advance(terminal_start, self.frames_decoded, rate)?
212 } else {
213 rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch)
217 };
218 if self.end.is_some() {
219 self.frames_decoded += frames;
220 }
221 self.next_start = Some(advance(mux_frame.timestamp, frames + trimmed, rate)?);
224 if decoded.is_empty() {
225 continue;
226 }
227
228 let (pcm, timestamp) = match self.resampler.as_mut() {
229 Some(r) => {
236 let held = if r.pending_frames() == 0 {
237 decoded_at
238 } else {
239 r.held_at().unwrap_or(decoded_at)
240 };
241 let skipped = r.skipped();
242 let pcm = r.process(&decoded, decoded_at)?;
243 (pcm, rewind(held, skipped, self.resolved_sample_rate)?)
244 }
245 None => (decoded, decoded_at),
246 };
247
248 let decoded_end = advance(decoded_at, frames, rate)?;
249
250 let resampled = self.resampler.is_some();
255 if resampled {
256 self.spans.push_back(ActivitySpan {
257 end: decoded_end,
258 activity,
259 });
260 }
261
262 if pcm.is_empty() {
266 continue;
267 }
268
269 let activity = if resampled {
270 self.activity_at(timestamp)
271 } else {
272 activity
273 };
274 let frame = self.frame(pcm, timestamp, activity)?;
277 self.ready.push_back(frame);
278 }
279 }
280
281 fn apply_discontinuity(&mut self) -> Result<(), Error> {
283 let discontinuity = self.track.discontinuity();
284 if discontinuity == self.discontinuity {
285 return Ok(());
286 }
287
288 self.discontinuity = discontinuity;
289 self.decoder.reset()?;
290 if let Some(resampler) = self.resampler.as_mut() {
291 resampler.reset();
292 }
293 self.next_start = None;
294 self.spans.clear();
295 self.trailing = Activity::Active;
296 self.epoch = None;
297 self.delay_trimmed = 0;
298 self.frames_decoded = 0;
299 self.end = None;
300 self.terminal_start = None;
301 Ok(())
302 }
303
304 fn gap(&mut self) -> Result<Option<Frame>, Error> {
313 self.decoder.reset_prediction()?;
314
315 let mut frame = None;
316 if let Some(resampler) = self.resampler.as_mut() {
317 let held = resampler.held_at();
318 let skipped = resampler.skipped();
319 let pcm = resampler.drain()?;
320 frame = self.tail(pcm, held, skipped)?;
321 }
322
323 self.next_start = None;
324 self.spans.clear();
325 self.trailing = Activity::Active;
326 self.epoch = None;
327 self.delay_trimmed = 0;
328 Ok(frame)
329 }
330
331 fn flush(&mut self) -> Result<Option<Frame>, Error> {
338 let Some(resampler) = self.resampler.take() else {
339 return Ok(None);
340 };
341
342 let held = resampler.held_at();
343 let skipped = resampler.skipped();
344 self.tail(resampler.flush()?, held, skipped)
345 }
346
347 fn tail(
353 &mut self,
354 pcm: Vec<f32>,
355 held: Option<moq_net::Timestamp>,
356 skipped: usize,
357 ) -> Result<Option<Frame>, Error> {
358 let Some(held) = held.filter(|_| !pcm.is_empty()) else {
359 return Ok(None);
360 };
361
362 let timestamp = rewind(held, skipped, self.resolved_sample_rate)?;
363 let activity = self.activity_at(timestamp);
364 Ok(Some(self.frame(pcm, timestamp, activity)?))
365 }
366
367 fn activity_at(&mut self, timestamp: moq_net::Timestamp) -> Activity {
369 while let Some(span) = self.spans.front().filter(|span| span.end <= timestamp) {
370 self.trailing = span.activity;
371 self.spans.pop_front();
372 }
373
374 self.spans.front().map_or(self.trailing, |span| span.activity)
375 }
376
377 fn frame(&self, pcm: Vec<f32>, timestamp: moq_net::Timestamp, activity: Activity) -> Result<Frame, Error> {
379 let pcm = if self.decoder.channel_count() == self.resolved_channels {
380 pcm
381 } else {
382 remix(&pcm, self.decoder.channel_count(), self.resolved_channels)?
383 };
384
385 let bytes = self.config.format.from_interleaved_f32(&pcm, self.resolved_channels)?;
386 Ok(Frame {
387 timestamp,
388 data: Bytes::from(bytes),
389 activity,
390 })
391 }
392}
393
394fn discontinuous(expected: moq_net::Timestamp, timestamp: moq_net::Timestamp) -> bool {
418 let scale = expected.scale().max(timestamp.scale());
419 let quantum = scale.min(moq_net::Timescale::default());
420 let tolerance = (scale.as_u64() as u128).div_ceil(quantum.as_u64() as u128) + 1;
421 expected.as_scale(scale).abs_diff(timestamp.as_scale(scale)) > tolerance
422}
423
424fn advance(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
426 if frames == 0 {
427 return Ok(timestamp);
428 }
429
430 let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
431 Ok(timestamp.checked_add(offset)?)
432}
433
434fn frames_between(start: moq_net::Timestamp, end: moq_net::Timestamp, sample_rate: u32) -> Result<usize, Error> {
436 let duration = end.checked_sub(start)?;
437 let frames = (std::time::Duration::from(duration).as_nanos() * sample_rate as u128 + 500_000_000) / 1_000_000_000;
438 usize::try_from(frames).map_err(|_| Error::Unsupported("audio duration does not fit in memory".into()))
439}
440
441fn rewind(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
446 if frames == 0 {
447 return Ok(timestamp);
448 }
449
450 let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
451 Ok(timestamp
452 .checked_sub(offset)
453 .unwrap_or(moq_net::Timestamp::new(0, timestamp.scale())?))
454}
455
456#[cfg(test)]
457mod tests {
458 use moq_net::Timestamp;
459
460 use super::*;
461 use crate::Format;
462 use crate::encode::{Encoder, Input, Options, Producer};
463
464 #[tokio::test]
465 async fn remixes_mono_stream_to_stereo_output() {
466 let mut broadcast = moq_net::broadcast::Info::new().produce();
467 let catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
468 let subscriber = broadcast.consume();
469 let input = Input {
470 format: Format::F32,
471 sample_rate: 48_000,
472 channels: 1,
473 };
474 let options = Options {
475 track: Some("audio".to_string()),
476 ..Options::default()
477 };
478 let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
479 let catalog = Encoder::new(&crate::encode::Config::new(input)).unwrap().catalog();
480 let mut consumer = Consumer::new(
481 &subscriber,
482 &catalog,
483 "audio",
484 Config {
485 channels: Some(2),
486 ..Config::new()
487 },
488 )
489 .await
490 .unwrap();
491
492 let samples = vec![0.1f32; 960];
493 let mut data = Vec::with_capacity(samples.len() * size_of::<f32>());
494 for sample in samples {
495 data.extend_from_slice(&sample.to_le_bytes());
496 }
497 producer.write(&Frame::new(data.into(), Timestamp::ZERO)).unwrap();
498
499 let frame = consumer.read().await.unwrap().expect("decoded frame");
500 let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
501 assert_eq!(samples.len(), (960 - 312) * 2);
502 for pair in samples.chunks_exact(2) {
503 assert_eq!(pair[0], pair[1]);
504 }
505 }
506
507 #[tokio::test]
514 async fn resampled_timestamps_follow_the_samples() {
515 let mut broadcast = moq_net::broadcast::Info::new().produce();
516 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
517 let subscriber = broadcast.consume();
518
519 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
520 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Legacy);
521
522 let mut consumer = Consumer::new(
523 &subscriber,
524 &catalog,
525 "audio",
526 Config {
527 sample_rate: Some(48_000),
528 ..Config::new()
529 },
530 )
531 .await
532 .unwrap();
533
534 const FRAMES: u64 = 1024;
536 let payload: Bytes = vec![0u8; FRAMES as usize * size_of::<f32>()].into();
537 for packet in 0..2 {
538 producer
539 .write(moq_mux::container::Frame {
540 timestamp: moq_net::Timestamp::from_scale(packet * FRAMES, 44_100).unwrap(),
541 duration: None,
542 payload: payload.clone(),
543 keyframe: true,
544 })
545 .unwrap();
546 }
547
548 let first = consumer.read().await.unwrap().expect("decoded frame");
549 assert_eq!(first.timestamp.as_micros(), 0);
550
551 let second = consumer.read().await.unwrap().expect("decoded frame");
557 let first_frames = (first.data.len() / size_of::<f32>()) as u128;
558 let ends_at = first_frames * 1_000_000 / 48_000;
559 let gap = second.timestamp.as_micros().abs_diff(ends_at);
560 assert!(gap < 100, "expected the frames to meet, got a {gap} us gap");
561 }
562
563 #[tokio::test]
567 async fn resampled_tail_survives_the_end_of_the_track() {
568 let mut broadcast = moq_net::broadcast::Info::new().produce();
569 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
570 let subscriber = broadcast.consume();
571
572 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
573 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Legacy);
574
575 let mut consumer = Consumer::new(
576 &subscriber,
577 &catalog,
578 "audio",
579 Config {
580 sample_rate: Some(48_000),
581 ..Config::new()
582 },
583 )
584 .await
585 .unwrap();
586
587 const FRAMES: usize = 1024;
589 let payload: Bytes = vec![0u8; FRAMES * size_of::<f32>()].into();
590 producer
591 .write(moq_mux::container::Frame {
592 timestamp: moq_net::Timestamp::ZERO,
593 duration: None,
594 payload,
595 keyframe: true,
596 })
597 .unwrap();
598 producer.finish().unwrap();
599
600 let first = consumer.read().await.unwrap().expect("decoded frame");
601 let first_frames = first.data.len() / size_of::<f32>();
602
603 let tail = consumer.read().await.unwrap().expect("flushed tail");
604 let tail_frames = tail.data.len() / size_of::<f32>();
605
606 assert!((215..=230).contains(&tail_frames), "unexpected tail: {tail_frames}");
610 let ends_at = (first_frames as u128) * 1_000_000 / 48_000;
613 let gap = tail.timestamp.as_micros().abs_diff(ends_at);
614 assert!(gap < 100, "expected the tail to meet the body, got a {gap} us gap");
615
616 let total = first_frames + tail_frames;
620 assert!((1105..=1120).contains(&total), "unexpected total: {total}");
621 assert!(consumer.read().await.unwrap().is_none());
622 }
623
624 #[tokio::test]
625 async fn resampling_keeps_the_activity_boundary_on_its_source() {
626 let mut encoder = Encoder::new(&crate::encode::Config {
627 dtx: true,
628 bitrate: Some(24_000),
629 frame_duration: std::time::Duration::from_millis(10),
630 ..crate::encode::Config::new(Input {
631 channels: 1,
632 ..Input::default()
633 })
634 })
635 .unwrap();
636 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Opus, 48_000, 1);
637
638 let mut broadcast = moq_net::broadcast::Info::new().produce();
639 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
640 let subscriber = broadcast.consume();
641 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Legacy);
642 let mut consumer = Consumer::new(
643 &subscriber,
644 &catalog,
645 "audio",
646 Config {
647 sample_rate: Some(44_100),
648 ..Config::new()
649 },
650 )
651 .await
652 .unwrap();
653
654 let active = vec![0.5; encoder.frame_size()];
655 let silence = vec![0.0; encoder.frame_size()];
656 let mut first_dtx = None;
657 for index in 0..40u64 {
658 let packet = encoder.encode(if index == 0 { &active } else { &silence }).unwrap();
659 let timestamp = Timestamp::from_scale(index * encoder.frame_size() as u64, 48_000).unwrap();
660 if first_dtx.is_none() && packet.activity.is_dtx() {
661 first_dtx = Some(timestamp);
662 }
663 producer
664 .write(moq_mux::container::Frame {
665 timestamp,
666 payload: packet.payload,
667 keyframe: true,
668 duration: None,
669 })
670 .unwrap();
671 producer.cut(None).unwrap();
672 }
673 producer.finish().unwrap();
674
675 let expected = first_dtx.expect("silence should enter Opus DTX");
676 let mut actual = None;
677 while let Some(frame) = consumer.read().await.unwrap() {
678 assert!(!frame.data.is_empty(), "read returned a frame with no samples");
683 if frame.activity.is_dtx() {
684 actual = Some(frame.timestamp);
685 break;
686 }
687 }
688 let actual = actual.expect("consumer should report Opus DTX");
689
690 let delay = actual.as_micros() as i128 - expected.as_micros() as i128;
697 let chunk_us = 20_000i128;
698 assert!(
699 (0..chunk_us).contains(&delay),
700 "DTX label landed {delay} us from its source, outside [0, {chunk_us})"
701 );
702 }
703
704 async fn pcm_gaps(rate: u32, out_rate: u32, frames: usize, stamps: &[Timestamp]) -> Vec<(u128, usize)> {
707 let mut broadcast = moq_net::broadcast::Info::new().produce();
708 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
709 let subscriber = broadcast.consume();
710
711 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, rate, 1);
712 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Legacy);
713 let mut consumer = Consumer::new(
714 &subscriber,
715 &catalog,
716 "audio",
717 Config {
718 sample_rate: Some(out_rate),
719 ..Config::new()
720 },
721 )
722 .await
723 .unwrap();
724
725 let payload: Bytes = vec![0u8; frames * size_of::<f32>()].into();
726 for stamp in stamps {
727 producer
728 .write(moq_mux::container::Frame {
729 timestamp: *stamp,
730 duration: None,
731 payload: payload.clone(),
732 keyframe: true,
733 })
734 .unwrap();
735 }
736 producer.finish().unwrap();
737
738 let mut read = Vec::new();
739 while let Some(frame) = consumer.read().await.unwrap() {
740 read.push((frame.timestamp.as_micros(), frame.data.len() / size_of::<f32>()));
741 }
742 read
743 }
744
745 #[tokio::test]
750 async fn a_missing_packet_leaves_a_hole() {
751 const FRAMES: usize = 1024;
752 let stamps = [
754 Timestamp::from_scale(0, 44_100).unwrap(),
755 Timestamp::from_scale(2 * FRAMES as u64, 44_100).unwrap(),
756 ];
757 let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
758
759 assert_eq!(read.len(), 4, "unexpected frames: {read:?}");
762
763 let before: usize = read[..2].iter().map(|(_, frames)| frames).sum();
766 assert!((1105..=1120).contains(&before), "unexpected pre-gap audio: {before}");
767
768 assert_eq!(read[2].0, stamps[1].as_micros());
772
773 let ends_at = read[1].0 + (read[1].1 as u128) * 1_000_000 / 48_000;
775 let hole = read[2].0 - ends_at;
776 assert!((23_100..=23_350).contains(&hole), "unexpected hole: {hole} us");
777 }
778
779 #[tokio::test]
785 async fn a_jump_inside_the_slack_leaves_the_held_samples_alone() {
786 const FRAMES: usize = 441;
789 let stamps = [
792 Timestamp::from_micros(0).unwrap(),
793 Timestamp::from_micros(11_000).unwrap(),
794 ];
795 let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
796
797 assert_eq!(read.len(), 2, "unexpected frames: {read:?}");
800 assert_eq!(read[0].0, 0, "held samples moved with the jump: {read:?}");
803 }
804
805 #[tokio::test]
806 async fn a_jump_after_a_full_chunk_uses_the_new_packet_timestamp() {
807 let stamps = [
808 Timestamp::from_micros(0).unwrap(),
809 Timestamp::from_micros(21_000).unwrap(),
810 ];
811 let read = pcm_gaps(44_100, 48_000, 882, &stamps).await;
812 let mut r = crate::Resampler::new(44_100, 48_000, 1, 882).unwrap();
813 r.process(&[0.25; 882], stamps[0]).unwrap();
814 let expected = rewind(stamps[1], r.skipped(), 48_000).unwrap().as_micros();
815 assert_eq!(read[1].0, expected);
816 }
817
818 #[tokio::test]
824 async fn a_terminal_jump_leaves_the_held_samples_alone() {
825 let mut encoder = Encoder::new(&crate::encode::Config {
826 dtx: true,
827 bitrate: Some(24_000),
828 ..crate::encode::Config::new(Input {
829 channels: 1,
830 ..Input::default()
831 })
832 })
833 .unwrap();
834 let catalog = encoder.catalog();
835
836 let active = encoder.encode(&vec![0.5f32; encoder.frame_size()]).unwrap();
840 assert!(active.activity.is_active());
841 let silence = vec![0.0f32; encoder.frame_size()];
842 let dtx = (0..200)
843 .map(|_| encoder.encode(&silence).unwrap())
844 .find(|packet| packet.activity.is_dtx())
845 .expect("silence should enter Opus DTX");
846
847 let mut broadcast = moq_net::broadcast::Info::new().produce();
848 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
849 let subscriber = broadcast.consume();
850 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Legacy);
851 let mut consumer = Consumer::new(
852 &subscriber,
853 &catalog,
854 "audio",
855 Config {
856 sample_rate: Some(44_100),
857 ..Config::new()
858 },
859 )
860 .await
861 .unwrap();
862
863 let write = |producer: &mut moq_mux::container::Producer<_>, frames: u64, payload: Bytes| {
866 producer
867 .write(moq_mux::container::Frame {
868 timestamp: Timestamp::from_scale(frames, 48_000).unwrap(),
869 duration: None,
870 payload,
871 keyframe: true,
872 })
873 .unwrap();
874 };
875 write(&mut producer, 0, active.payload);
876 write(&mut producer, 3 * 48_000, Bytes::new());
878 write(&mut producer, 48_000, dtx.payload);
879 producer.finish().unwrap();
880
881 let frame = consumer.read().await.unwrap().expect("decoded frame");
882 assert_eq!(frame.timestamp.as_micros(), 0, "held samples moved with the jump");
886 assert!(
887 frame.activity.is_active(),
888 "held samples took the terminal packet's label"
889 );
890 }
891
892 #[tokio::test]
900 async fn millisecond_stamps_are_not_a_gap() {
901 const FRAMES: u64 = 1024;
902 const PACKETS: u64 = 32;
903
904 let stamps: Vec<_> = (0..PACKETS)
906 .map(|packet| Timestamp::from_millis(packet * FRAMES * 1_000 / 44_100).unwrap())
907 .collect();
908 let read = pcm_gaps(44_100, 48_000, FRAMES as usize, &stamps).await;
909
910 assert_eq!(read.len(), stamps.len() + 1, "unexpected frames: {read:?}");
914
915 for pair in read.windows(2) {
918 let ends_at = pair[0].0 + (pair[0].1 as u128) * 1_000_000 / 48_000;
919 assert!(
920 pair[1].0.abs_diff(ends_at) <= 1_100,
921 "frames at {} and {} do not meet",
922 pair[0].0,
923 pair[1].0
924 );
925 }
926 }
927
928 #[tokio::test]
932 async fn a_lost_opus_packet_shorter_than_its_neighbour_is_a_gap() {
933 let input = Input {
934 format: Format::F32,
935 sample_rate: 48_000,
936 channels: 1,
937 };
938 let mut encoder = Encoder::new(&crate::encode::Config::new(input)).unwrap();
939 let catalog = encoder.catalog();
940
941 let mut broadcast = moq_net::broadcast::Info::new().produce();
942 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
943 let subscriber = broadcast.consume();
944 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Legacy);
945 let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Config::new())
946 .await
947 .unwrap();
948
949 let pcm = vec![0.25f32; encoder.frame_size()];
952 for timestamp in [
953 Timestamp::from_micros(0).unwrap(),
954 Timestamp::from_micros(22_500).unwrap(),
955 Timestamp::from_micros(42_500).unwrap(),
956 ] {
957 producer
958 .write(moq_mux::container::Frame {
959 timestamp,
960 duration: None,
961 payload: encoder.encode(&pcm).unwrap().payload,
962 keyframe: true,
963 })
964 .unwrap();
965 producer.cut(None).unwrap();
966 }
967
968 let first = consumer.read().await.unwrap().expect("decoded frame");
972 let frames = first.data.len() / size_of::<f32>();
973 assert!(frames < 960, "the pre-skip should be trimmed, got {frames} frames");
974
975 let second = consumer.read().await.unwrap().expect("decoded frame");
978 assert_eq!(second.timestamp.as_micros(), 22_500);
979 assert_eq!(second.data.len() / size_of::<f32>(), 960, "pre-skip was reapplied");
980
981 let third = consumer.read().await.unwrap().expect("decoded frame after gap");
982 let second_frames = second.data.len() / size_of::<f32>();
983 assert_eq!(
984 third.timestamp,
985 advance(second.timestamp, second_frames, 48_000).unwrap()
986 );
987 }
988
989 #[tokio::test]
990 async fn latency_max_is_clamped_to_publisher_retention() {
991 let mut broadcast = moq_net::broadcast::Info::new().produce();
992 let info = hang::container::track_info().with_latency_max(std::time::Duration::from_millis(100));
993 let _track = broadcast.create_track("audio", info).unwrap();
994 let subscriber = broadcast.consume();
995 let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
996
997 let consumer = Consumer::new(
998 &subscriber,
999 &catalog,
1000 "audio",
1001 Config {
1002 latency_max: Some(std::time::Duration::from_millis(500)),
1003 ..Config::new()
1004 },
1005 )
1006 .await
1007 .unwrap();
1008
1009 assert_eq!(consumer.latency_max(), std::time::Duration::from_millis(100));
1010 }
1011
1012 #[tokio::test]
1016 async fn opus_pre_skip_does_not_leave_a_timestamp_hole() {
1017 let input = Input {
1018 format: Format::F32,
1019 sample_rate: 48_000,
1020 channels: 1,
1021 };
1022 let mut encoder = Encoder::new(&crate::encode::Config::new(input)).unwrap();
1023 let catalog = encoder.catalog();
1024
1025 let mut broadcast = moq_net::broadcast::Info::new().produce();
1026 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
1027 let subscriber = broadcast.consume();
1028 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Legacy);
1029 let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Config::new())
1030 .await
1031 .unwrap();
1032
1033 let pcm = vec![0.25f32; encoder.frame_size()];
1034 for packet in 0..2 {
1035 producer
1036 .write(moq_mux::container::Frame {
1037 timestamp: Timestamp::from_scale(packet * encoder.frame_size() as u64, 48_000).unwrap(),
1038 duration: None,
1039 payload: encoder.encode(&pcm).unwrap().payload,
1040 keyframe: true,
1041 })
1042 .unwrap();
1043 producer.cut(None).unwrap();
1044 }
1045
1046 let first = consumer.read().await.unwrap().expect("first decoded frame");
1047 let second = consumer.read().await.unwrap().expect("second decoded frame");
1048 let first_frames = first.data.len() / size_of::<f32>();
1049 let expected = advance(first.timestamp, first_frames, 48_000).unwrap();
1050 assert_eq!(second.timestamp, expected);
1051 }
1052
1053 #[tokio::test]
1054 async fn reads_the_container_the_catalog_declares() {
1055 let mut broadcast = moq_net::broadcast::Info::new().produce();
1056 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
1057 let subscriber = broadcast.consume();
1058
1059 let mut catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
1060 catalog.container = hang::catalog::Container::Loc;
1061
1062 let mut producer = moq_mux::container::Producer::new(track, moq_mux::catalog::hang::Container::Loc);
1063 let mut consumer = Consumer::new(
1064 &subscriber,
1065 &catalog,
1066 "audio",
1067 Config {
1068 format: Format::F32,
1069 ..Config::new()
1070 },
1071 )
1072 .await
1073 .unwrap();
1074
1075 let samples = [0.25f32, -0.5, 0.75, -1.0];
1076 let payload: Vec<u8> = samples.iter().flat_map(|sample| sample.to_le_bytes()).collect();
1077 producer
1078 .write(moq_mux::container::Frame {
1079 timestamp: Timestamp::ZERO,
1080 duration: None,
1081 payload: payload.into(),
1082 keyframe: true,
1083 })
1084 .unwrap();
1085
1086 let frame = consumer.read().await.unwrap().expect("decoded frame");
1087 assert_eq!(
1088 Format::F32.as_interleaved_f32(&frame.data, 1).unwrap().as_ref(),
1089 samples
1090 );
1091 }
1092
1093 #[tokio::test]
1098 async fn decodes_a_cmaf_framed_track() {
1099 let input = Input {
1100 format: Format::F32,
1101 sample_rate: 48_000,
1102 channels: 2,
1103 };
1104
1105 let mut encoder = Encoder::new(&crate::encode::Config::new(input.clone())).unwrap();
1107 let mut catalog = encoder.catalog();
1108 let pcm = vec![0.0f32; encoder.frame_size() * encoder.codec_channels() as usize];
1109 let packet = encoder.encode(&pcm).unwrap();
1110
1111 let muxer = moq_mux::container::fmp4::Muxer::audio(&catalog).unwrap();
1113 let init = muxer.init().unwrap().expect("an out-of-band codec has an init segment");
1114 catalog.container = hang::catalog::Container::Cmaf { init };
1115
1116 let mut broadcast = moq_net::broadcast::Info::new().produce();
1117 let subscriber = broadcast.consume();
1118 let track = broadcast.create_track("audio", hang::container::track_info()).unwrap();
1119 let container = moq_mux::catalog::hang::Container::try_from(&catalog.container).unwrap();
1120 let mut producer = moq_mux::container::Producer::new(track, container);
1121
1122 let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Config::new())
1123 .await
1124 .unwrap();
1125
1126 producer
1127 .write(moq_mux::container::Frame {
1128 timestamp: Timestamp::ZERO,
1129 payload: packet.payload,
1130 keyframe: true,
1131 duration: None,
1132 })
1133 .unwrap();
1134 producer.cut(None).unwrap();
1135
1136 let frame = consumer.read().await.unwrap().expect("decoded frame");
1140 assert_eq!(frame.timestamp.as_micros(), 0);
1143 let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
1144 assert_eq!(samples.len(), (960 - 312) * 2);
1145 }
1146}