1use std::time::Instant;
4
5use bytes::Bytes;
6
7use moq_mux::catalog::hang::CatalogExt;
8use moq_mux::container::Frame as MuxFrame;
9use moq_net::Timestamp;
10
11use super::encoded::Encoded;
12use super::encoder::{Encoder, Input, Settings};
13use crate::resample::{Remix, Resampler};
14use crate::{Activity, Error, Frame};
15
16#[derive(Clone, Debug)]
21#[non_exhaustive]
22pub struct Options {
23 pub track: Option<String>,
27 pub settings: Settings,
29 pub bandwidth: moq_net::bandwidth::Allocator,
43}
44
45impl Default for Options {
46 fn default() -> Self {
47 Self {
48 track: None,
49 settings: Settings::default(),
50 bandwidth: moq_net::bandwidth::Allocator::unlimited(),
51 }
52 }
53}
54
55pub struct Producer<E: CatalogExt = ()> {
65 encoder: Encoder,
66 input: Input,
67 remix: Option<Remix>,
69 resampler: Option<Resampler>,
70 track: moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
71 _ext: std::marker::PhantomData<fn() -> E>,
72 pending: Vec<f32>,
73 frames_produced: u64,
75 epoch_us: Option<u64>,
79 pending_discontinuity: bool,
81 decoder_boundary: bool,
83 activity: Activity,
85 finished: bool,
88}
89
90struct Terminal {
91 packets: Vec<Encoded>,
92 end: Timestamp,
93 start: Timestamp,
94 frame_size: usize,
95 codec_rate: u32,
96}
97
98pub(crate) struct Reserved<E: CatalogExt = ()> {
106 track: moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
107 _ext: std::marker::PhantomData<fn() -> E>,
108}
109
110impl<E: CatalogExt> Reserved<E> {
111 pub(crate) fn new(
112 broadcast: &mut moq_net::broadcast::Producer,
113 catalog: moq_mux::catalog::Producer<E>,
114 options: &Options,
115 ) -> Result<Self, Error> {
116 let track = match &options.track {
117 Some(name) => broadcast.create_track(name.clone(), catalog.track_info(hang::catalog::PRIORITY.audio))?,
121 None => broadcast.unique_track(
124 &format!(".{}", options.settings.codec),
125 catalog.track_info(hang::catalog::PRIORITY.audio),
126 )?,
127 };
128 let track = catalog.audio(
129 track,
130 moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio),
131 None,
132 )?;
133
134 Ok(Self {
135 track,
136 _ext: std::marker::PhantomData,
137 })
138 }
139
140 pub(crate) fn register(&mut self, input: Input, options: &Options) -> Result<Registered, Error> {
145 let remix = (input.layout != options.settings.layout)
146 .then(|| Remix::new(input.layout, options.settings.layout))
147 .transpose()?;
148 let encoder = Encoder::new(&options.settings)?;
149
150 let resampler = if input.sample_rate == encoder.codec_rate() {
151 None
152 } else {
153 let chunk_frames =
156 ((input.sample_rate as u128 * encoder.settings().frame_duration.as_micros()) / 1_000_000) as usize;
157 Some(Resampler::new(
158 input.sample_rate,
159 encoder.codec_rate(),
160 encoder.codec_channels(),
161 chunk_frames,
162 )?)
163 };
164
165 self.track.set(encoder.catalog())?;
166
167 Ok(Registered {
168 encoder,
169 input,
170 remix,
171 resampler,
172 })
173 }
174
175 pub(crate) fn encode(self, registered: Registered) -> Producer<E> {
177 Producer {
178 encoder: registered.encoder,
179 input: registered.input,
180 remix: registered.remix,
181 resampler: registered.resampler,
182 track: self.track,
183 _ext: self._ext,
184 pending: Vec::new(),
185 frames_produced: 0,
186 epoch_us: None,
187 pending_discontinuity: false,
188 decoder_boundary: true,
189 activity: Activity::Active,
190 finished: false,
191 }
192 }
193}
194
195pub(crate) struct Registered {
200 encoder: Encoder,
201 input: Input,
202 remix: Option<Remix>,
203 resampler: Option<Resampler>,
204}
205
206#[cfg(feature = "capture")]
208impl<E: CatalogExt> Reserved<E> {
209 pub(crate) fn name(&self) -> &str {
211 self.track.name()
212 }
213
214 pub(crate) fn track(&self) -> &moq_net::track::Producer {
216 self.track.track()
217 }
218
219 pub(crate) fn finish(&mut self) -> Result<(), Error> {
221 self.track.finish()?;
222 Ok(())
223 }
224
225 pub(crate) fn abort(self, err: moq_net::Error) {
227 self.track.abort(err);
228 }
229}
230
231impl<E: CatalogExt> Producer<E> {
232 pub fn new(
235 broadcast: &mut moq_net::broadcast::Producer,
236 catalog: moq_mux::catalog::Producer<E>,
237 input: Input,
238 options: &Options,
239 ) -> Result<Self, Error> {
240 let mut reserved = Reserved::new(broadcast, catalog, options)?;
241 let registered = reserved.register(input, options)?;
242 Ok(reserved.encode(registered))
243 }
244
245 pub fn track_name(&self) -> &str {
247 self.track.name()
248 }
249
250 pub fn demand(&self) -> moq_net::track::Demand {
254 self.track.track().demand()
255 }
256
257 #[cfg(feature = "capture")]
258 pub(crate) fn track(&self) -> &moq_net::track::Producer {
259 self.track.track()
260 }
261
262 pub fn activity(&self) -> Activity {
271 self.activity
272 }
273
274 pub fn bitrate(&self) -> moq_net::bandwidth::Rate {
276 self.encoder.bitrate()
277 }
278
279 pub fn set_bitrate(&mut self, bitrate: moq_net::bandwidth::Rate) -> Result<(), Error> {
281 self.encoder.set_bitrate(bitrate)
282 }
283
284 pub fn reset_epoch(&mut self) {
292 if self.encoder.started() && !self.decoder_boundary {
293 self.pending_discontinuity = true;
294 }
295 self.reset_state();
296 }
297
298 fn reset_state(&mut self) {
299 self.epoch_us = None;
300 self.activity = Activity::Active;
301 self.frames_produced = 0;
302 self.pending.clear();
303 self.encoder.reset();
304 if let Some(resampler) = self.resampler.as_mut() {
309 resampler.reset();
310 }
311 }
312
313 pub fn write(&mut self, frame: &Frame) -> Result<(), Error> {
331 if self.finished {
332 return Err(moq_net::Error::Closed.into());
333 }
334 if self.pending_discontinuity {
335 self.track.discontinuity()?;
336 self.pending_discontinuity = false;
337 self.decoder_boundary = true;
338 }
339
340 let timestamp_us = u64::try_from(frame.timestamp.as_micros())
341 .map_err(|_| Error::Unsupported(format!("frame timestamp {:?} out of range", frame.timestamp)))?;
342 let epoch_us = *self.epoch_us.get_or_insert(timestamp_us);
343
344 let input = &self.input;
345 let (format, channels) = (input.format, input.layout.channels());
346 let pcm = format.as_interleaved_f32(frame.data.as_ref(), channels)?;
347 let pcm = match &self.remix {
348 Some(remix) => remix.process(&pcm),
349 None => pcm.into_owned(),
350 };
351 let pcm: Vec<f32> = match self.resampler.as_mut() {
352 Some(r) => r.process(&pcm, frame.timestamp)?,
353 None => pcm,
354 };
355
356 self.pending.extend(pcm);
357
358 self.publish_full_frames(epoch_us)
359 }
360
361 fn publish_full_frames(&mut self, epoch_us: u64) -> Result<(), Error> {
364 let frame_samples = self.encoder.frame_size() * self.encoder.codec_channels() as usize;
365 while self.pending.len() >= frame_samples {
366 let chunk: Vec<f32> = self.pending.drain(..frame_samples).collect();
367 let packet = self.encoder.encode(&chunk)?;
368
369 let timestamp = Self::timestamp(
370 epoch_us,
371 self.frames_produced,
372 self.encoder.folded_delay(),
373 self.encoder.codec_rate(),
374 )?;
375 self.frames_produced += self.encoder.frame_size() as u64;
376 self.activity = packet.activity;
377 Self::publish(&mut self.track, packet, timestamp)?;
378 self.decoder_boundary = false;
379 }
380
381 Ok(())
382 }
383
384 fn timestamp(epoch_us: u64, frames: u64, delay: usize, codec_rate: u32) -> Result<Timestamp, Error> {
391 let frames = i128::from(frames) - delay as i128;
392 let offset_us = (frames * 1_000_000).div_euclid(i128::from(codec_rate));
393 let micros = (i128::from(epoch_us) + offset_us).max(0);
394 let micros = u64::try_from(micros).map_err(|_| moq_net::TimeOverflow)?;
395 Ok(Timestamp::from_micros(micros)?)
396 }
397
398 fn publish(
399 track: &mut moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
400 encoded: Encoded,
401 timestamp: Timestamp,
402 ) -> Result<(), Error> {
403 let mux_frame = MuxFrame {
407 timestamp,
408 payload: encoded.payload,
409 keyframe: true,
410 duration: None,
411 };
412 track.write(mux_frame)?;
413 track.cut(None)?;
417 track.flush(timestamp, Instant::now())?;
418 Ok(())
419 }
420
421 fn publish_terminal(
423 track: &mut moq_mux::container::Producer<moq_mux::container::legacy::Wire, hang::catalog::AudioConfig>,
424 terminal: Terminal,
425 ) -> Result<(), Error> {
426 track.write(MuxFrame {
427 timestamp: terminal.end,
428 payload: Bytes::new(),
429 keyframe: true,
430 duration: None,
431 })?;
432
433 for (index, packet) in terminal.packets.into_iter().enumerate() {
434 let offset = Timestamp::from_scale((index * terminal.frame_size) as u64, terminal.codec_rate as u64)?
435 .convert(terminal.start.scale())?;
436 let timestamp = terminal.start.checked_add(offset)?;
437 track.write(MuxFrame {
438 timestamp,
439 payload: packet.payload,
440 keyframe: false,
441 duration: None,
442 })?;
443 track.flush(timestamp, Instant::now())?;
444 }
445
446 track.cut(Some(terminal.end))?;
447 Ok(())
448 }
449
450 pub fn discontinuity(&mut self) -> Result<(), Error> {
457 self.track.discontinuity()?;
458 self.pending_discontinuity = false;
459 self.decoder_boundary = true;
460 self.reset_state();
461 Ok(())
462 }
463
464 pub fn finish(&mut self) -> Result<(), Error> {
471 if self.finished {
472 return Ok(());
473 }
474 if let Some(resampler) = self.resampler.take() {
478 self.pending.extend(resampler.flush()?);
479 }
480
481 let epoch_us = self.epoch_us.unwrap_or(0);
484 self.publish_full_frames(epoch_us)?;
485
486 let frame_size = self.encoder.frame_size();
487 let codec_rate = self.encoder.codec_rate();
488 let channels = self.encoder.codec_channels() as usize;
489 let source_frames = self.pending.len() / channels;
490 let delay = self.encoder.folded_delay();
491 let start = Self::timestamp(epoch_us, self.frames_produced, delay, codec_rate)?;
492 let end = Self::timestamp(epoch_us, self.frames_produced + source_frames as u64, 0, codec_rate)?;
494 let finish = self.encoder.drain(&self.pending)?;
495 let discard_padding = finish.discard_padding();
496 let packets = finish.into_packets();
497
498 if discard_padding > 0 {
499 Self::publish_terminal(
500 &mut self.track,
501 Terminal {
502 packets,
503 end,
504 start,
505 frame_size,
506 codec_rate,
507 },
508 )?;
509 } else {
510 for packet in packets {
511 let timestamp = Self::timestamp(epoch_us, self.frames_produced, delay, codec_rate)?;
512 self.activity = packet.activity;
513 Self::publish(&mut self.track, packet, timestamp)?;
514 self.frames_produced += frame_size as u64;
515 }
516 }
517
518 self.track.finish()?;
519 self.finished = true;
520 Ok(())
521 }
522
523 pub fn abort(self, err: moq_net::Error) {
528 self.track.abort(err);
529 }
530}
531
532#[cfg(test)]
533mod tests {
534 use std::time::Duration;
535
536 use super::*;
537 use crate::decode::{Consumer as AudioConsumer, Options as DecodeOptions};
538 use crate::{Activity, Format, Layout};
539
540 #[tokio::test]
541 async fn demand_follows_subscribers_and_closes_with_the_producer() {
542 let mut broadcast = moq_net::broadcast::Info::new().produce();
543 let consumer = broadcast.consume();
544 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
545 let options = Options {
546 track: Some("audio".into()),
547 ..Options::default()
548 };
549 let mut producer = Producer::new(&mut broadcast, catalog, Input::default(), &options).unwrap();
550 let demand = producer.demand();
551
552 assert_eq!(demand.name(), "audio");
553 assert!(!demand.is_used());
554 let subscriber = consumer.track("audio").unwrap().subscribe(None).await.unwrap();
555 tokio::time::timeout(Duration::from_secs(1), demand.used())
556 .await
557 .expect("subscription demand")
558 .unwrap();
559 assert!(demand.is_used());
560
561 drop(subscriber);
562 drop(consumer);
563 tokio::time::timeout(Duration::from_secs(1), demand.unused())
564 .await
565 .expect("subscription released")
566 .unwrap();
567 assert!(!demand.is_used());
568
569 producer.finish().unwrap();
570 drop(producer);
571 let closed = tokio::time::timeout(Duration::from_secs(1), demand.closed())
572 .await
573 .expect("producer closed");
574 assert!(matches!(closed, moq_net::Error::Dropped));
575 }
576
577 #[tokio::test]
579 async fn finish_publishes_the_opus_lookahead_tail() {
580 for frames in [960, 860] {
581 let input = Input {
582 format: Format::F32,
583 sample_rate: 48_000,
584 layout: Layout::Mono,
585 };
586 let options = Options {
587 track: Some("audio".to_string()),
588 settings: Settings {
589 layout: Layout::Mono,
590 bitrate: Some(moq_net::bandwidth::Rate::from_bps(128_000)),
591 ..Settings::default()
592 },
593 ..Options::default()
594 };
595 let decoder_config = Encoder::new(&options.settings).unwrap().catalog();
596
597 let mut broadcast = moq_net::broadcast::Info::new().produce();
598 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
599 let consumer = broadcast.consume();
600 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
601 let mut audio = AudioConsumer::new(
602 &consumer,
603 &decoder_config,
604 "audio",
605 DecodeOptions {
606 max_age: Duration::from_secs(1),
607 ..DecodeOptions::new()
608 },
609 )
610 .await
611 .unwrap();
612
613 let mut pcm = vec![0.0f32; frames];
614 let impulse = pcm.len() - 100;
615 pcm[impulse] = 1.0;
616 let data: Vec<u8> = pcm.iter().flat_map(|sample| sample.to_le_bytes()).collect();
617 producer.write(&Frame::new(Bytes::from(data), Timestamp::ZERO)).unwrap();
618 producer.finish().unwrap();
619
620 let mut decoded = Vec::new();
621 while let Some(frame) = audio.read().await.unwrap() {
622 let pcm = Format::F32.as_interleaved_f32(&frame.data, 1).unwrap();
623 decoded.extend_from_slice(&pcm);
624 }
625 assert_eq!(decoded.len(), frames, "terminal padding extended the source");
626 let peak = decoded.iter().fold(0.0f32, |peak, sample| peak.max(sample.abs()));
627 assert!(peak > 0.1, "the {frames}-frame Opus tail lost the impulse: peak {peak}");
628 }
629 }
630
631 #[tokio::test]
635 async fn finish_publishes_the_resampled_tail() {
636 let input = Input {
637 format: Format::F32,
638 sample_rate: 44_100,
639 layout: Layout::Mono,
640 };
641
642 let mut broadcast = moq_net::broadcast::Info::new().produce();
643 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
644 let consumer = broadcast.consume();
645 let options = Options {
646 track: Some("audio".to_string()),
647 settings: Settings::new(48_000, Layout::Mono),
648 ..Options::default()
649 };
650 let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
651
652 let mut track = moq_mux::container::Consumer::new(
654 consumer
655 .track("audio")
656 .unwrap()
657 .subscribe(moq_net::track::Subscription::default().with_max_age(Duration::from_secs(1)))
658 .await
659 .unwrap(),
660 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
661 );
662
663 let data: Vec<u8> = vec![0.25f32; 8_838].iter().flat_map(|s| s.to_le_bytes()).collect();
667 producer
668 .write(&Frame::new(data.into(), moq_net::Timestamp::ZERO))
669 .unwrap();
670 producer.finish().unwrap();
671
672 let mut packets = 0;
673 while track.read().await.unwrap().is_some() {
674 packets += 1;
675 }
676 assert_eq!(packets, 11);
677 }
678
679 #[tokio::test]
683 async fn reset_epoch_drops_the_resampler_buffer_too() {
684 let input = Input {
685 format: Format::F32,
686 sample_rate: 44_100,
687 layout: Layout::Mono,
688 };
689
690 let mut broadcast = moq_net::broadcast::Info::new().produce();
691 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
692 let consumer = broadcast.consume();
693 let options = Options {
694 track: Some("audio".to_string()),
695 ..Options::default()
696 };
697 let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
698
699 let mut track = moq_mux::container::Consumer::new(
700 consumer
701 .track("audio")
702 .unwrap()
703 .subscribe(moq_net::track::Subscription::default())
704 .await
705 .unwrap(),
706 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
707 );
708
709 let data: Vec<u8> = vec![0.25f32; 441].iter().flat_map(|s| s.to_le_bytes()).collect();
711 producer
712 .write(&Frame::new(data.into(), moq_net::Timestamp::ZERO))
713 .unwrap();
714
715 producer.reset_epoch();
716 producer.finish().unwrap();
717
718 assert!(track.read().await.unwrap().is_none());
720 }
721
722 #[tokio::test]
724 async fn reset_epoch_drops_the_encoder_lookahead() {
725 let mut broadcast = moq_net::broadcast::Info::new().produce();
726 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
727 let consumer = broadcast.consume();
728 let options = Options {
729 track: Some("audio".to_string()),
730 ..Options::default()
731 };
732 let mut producer = Producer::new(
733 &mut broadcast,
734 catalog,
735 Input {
736 layout: Layout::Mono,
737 ..Input::default()
738 },
739 &options,
740 )
741 .unwrap();
742 let mut track = moq_mux::container::Consumer::new(
743 consumer
744 .track("audio")
745 .unwrap()
746 .subscribe(moq_net::track::Subscription::default())
747 .await
748 .unwrap(),
749 moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
750 );
751
752 producer.write(&full_frame(1_000_000)).unwrap();
753 producer.reset_epoch();
754 producer.finish().unwrap();
755
756 assert!(track.read().await.unwrap().is_some());
757 assert!(track.read().await.unwrap().is_none());
758 }
759
760 #[tokio::test]
762 async fn reset_epoch_restarts_the_decoder() {
763 let input = Input {
764 format: Format::F32,
765 sample_rate: 48_000,
766 layout: Layout::Mono,
767 };
768 let options = Options {
769 track: Some("audio".to_string()),
770 settings: Settings::new(48_000, Layout::Mono),
771 ..Options::default()
772 };
773 let decoder_config = Encoder::new(&options.settings).unwrap().catalog();
774
775 let mut broadcast = moq_net::broadcast::Info::new().produce();
776 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
777 let subscriber = broadcast.consume();
778 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
779 let mut audio = AudioConsumer::new(
780 &subscriber,
781 &decoder_config,
782 "audio",
783 DecodeOptions {
784 max_age: Duration::from_millis(500),
785 ..DecodeOptions::new()
786 },
787 )
788 .await
789 .unwrap();
790
791 producer.write(&full_frame(0)).unwrap();
792 let first = audio.read().await.unwrap().expect("first epoch packet");
793 assert_eq!(first.data.len() / size_of::<f32>(), 960 - 312);
794
795 producer.reset_epoch();
796 producer.write(&full_frame(1_000_000)).unwrap();
797 producer.finish().unwrap();
798
799 let mut resumed_frames = 0;
800 while let Some(frame) = audio.read().await.unwrap() {
801 if frame.timestamp.as_micros() >= 1_000_000 {
802 resumed_frames += frame.data.len() / size_of::<f32>();
803 }
804 }
805 assert!(resumed_frames > 0, "the resumed epoch still decodes");
806 }
807
808 #[tokio::test]
809 async fn producer_and_consumer_keep_activity_on_the_audio_stream() {
810 let input = Input {
811 format: Format::F32,
812 sample_rate: 48_000,
813 layout: Layout::Mono,
814 };
815 let options = Options {
816 track: Some("audio".to_string()),
817 settings: Settings {
818 layout: Layout::Mono,
819 bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
820 dtx: true,
821 ..Settings::default()
822 },
823 ..Options::default()
824 };
825 let decoder_config = Encoder::new(&options.settings).unwrap().catalog();
826
827 let mut broadcast = moq_net::broadcast::Info::new().produce();
828 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
829 let subscriber = broadcast.consume();
830 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
831 let mut consumer = AudioConsumer::new(&subscriber, &decoder_config, "audio", DecodeOptions::new())
832 .await
833 .unwrap();
834
835 let silence = vec![0.0; 960];
836 let mut entered_dtx = false;
837 for index in 0..100 {
838 producer.write(&pcm_frame(&silence, index * 20_000)).unwrap();
839 let consumed = consumer.read().await.unwrap().expect("one decoded frame");
840 assert_eq!(producer.activity(), consumed.activity);
841 if consumed.activity.is_dtx() {
842 entered_dtx = true;
843 break;
844 }
845 }
846 assert!(entered_dtx, "silence should enter Opus DTX");
847
848 let active: Vec<f32> = (0..960)
849 .map(|sample| {
850 let phase = sample as f32 * 440.0 * 2.0 * std::f32::consts::PI / 48_000.0;
851 phase.sin() * 0.5
852 })
853 .collect();
854 producer.write(&pcm_frame(&active, 2_000_000)).unwrap();
855 let consumed = consumer.read().await.unwrap().expect("one decoded frame");
856 assert_eq!(producer.activity(), Activity::Active);
857 assert_eq!(consumed.activity, Activity::Active);
858 }
859
860 fn full_frame(timestamp_us: u64) -> Frame {
863 pcm_frame(&vec![0.1; 960], timestamp_us)
864 }
865
866 fn pcm_frame(samples: &[f32], timestamp_us: u64) -> Frame {
867 let data: Vec<u8> = samples.iter().flat_map(|sample| sample.to_le_bytes()).collect();
868 Frame::new(Bytes::from(data), Timestamp::from_micros(timestamp_us).unwrap())
869 }
870
871 async fn published_pts(frames: &[Frame], reset_before: Option<usize>) -> Vec<u128> {
875 let mut broadcast = moq_net::broadcast::Info::new().produce();
876 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
877 let consumer = broadcast.consume();
878
879 let input = Input {
882 format: Format::F32,
883 sample_rate: 48_000,
884 layout: Layout::Mono,
885 };
886 let options = Options {
887 track: Some("audio".to_string()),
888 ..Options::default()
889 };
890 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
891
892 let track = consumer.track("audio").unwrap().subscribe(None).await.unwrap();
893 let mut reader =
894 moq_mux::container::Consumer::new(track, moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio));
895
896 let mut pts = Vec::new();
897 for (i, frame) in frames.iter().enumerate() {
898 if reset_before == Some(i) {
899 producer.reset_epoch();
900 }
901 producer.write(frame).unwrap();
902 let read = reader.read().await.unwrap().expect("a packet per full frame");
903 pts.push(read.timestamp.as_micros());
904 }
905 pts
906 }
907
908 #[tokio::test]
909 async fn epoch_anchors_to_first_frame_timestamp() {
910 let pts = published_pts(&[full_frame(1_000_000)], None).await;
913 assert_eq!(pts, vec![1_000_000]);
914 }
915
916 #[tokio::test]
919 async fn aac_folds_the_encoder_delay_into_timestamps() {
920 async fn pts(epoch_us: u64) -> Vec<u128> {
921 let mut broadcast = moq_net::broadcast::Info::new().produce();
922 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
923 let consumer = broadcast.consume();
924
925 let input = Input::new(48_000, Layout::Mono);
926 let options = Options {
927 track: Some("audio".to_string()),
928 settings: Settings::from_input(crate::encode::Codec::Aac, &input),
929 ..Options::default()
930 };
931 let stub = crate::encode::backend::stub::install();
932 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
933 drop(stub);
934
935 let track = consumer
936 .track("audio")
937 .unwrap()
938 .subscribe(moq_net::track::Subscription::default().with_max_age(Duration::from_secs(1)))
939 .await
940 .unwrap();
941 let mut reader = moq_mux::container::Consumer::new(
942 track,
943 moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio),
944 );
945
946 producer.write(&pcm_frame(&[0.1; 4 * 1024], epoch_us)).unwrap();
947 let mut pts = Vec::new();
948 for _ in 0..4 {
949 pts.push(reader.read().await.unwrap().expect("a packet").timestamp.as_micros());
950 }
951 pts
952 }
953
954 assert_eq!(pts(1_000_000).await, vec![956_000, 977_333, 998_666, 1_020_000]);
956 assert_eq!(pts(0).await, vec![0, 0, 0, 20_000]);
958 }
959
960 #[tokio::test]
967 async fn resampling_does_not_shift_the_first_pts() {
968 let mut broadcast = moq_net::broadcast::Info::new().produce();
969 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
970 let consumer = broadcast.consume();
971
972 let input = Input {
974 format: Format::F32,
975 sample_rate: 44_100,
976 layout: Layout::Mono,
977 };
978 let options = Options {
979 track: Some("audio".to_string()),
980 ..Options::default()
981 };
982 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
983
984 let track = consumer
987 .track("audio")
988 .unwrap()
989 .subscribe(moq_net::track::Subscription::default().with_max_age(Duration::from_secs(1)))
990 .await
991 .unwrap();
992 let mut reader =
993 moq_mux::container::Consumer::new(track, moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio));
994
995 producer.write(&pcm_frame(&vec![0.1; 44_100], 1_000_000)).unwrap();
997
998 let first = reader.read().await.unwrap().expect("a packet");
999 assert_eq!(first.timestamp.as_micros(), 1_000_000);
1000 }
1001
1002 #[tokio::test]
1003 async fn pts_advances_by_frame_duration_ignoring_later_timestamps() {
1004 let pts = published_pts(&[full_frame(1_000), full_frame(999_999)], None).await;
1007 assert_eq!(pts, vec![1_000, 1_000 + 20_000]);
1008 }
1009
1010 #[tokio::test]
1011 async fn reset_epoch_reanchors_so_the_gap_lands_in_pts() {
1012 let pts = published_pts(&[full_frame(0), full_frame(5_000_000)], Some(1)).await;
1015 assert_eq!(pts, vec![0, 5_000_000]);
1016 }
1017
1018 #[tokio::test]
1020 async fn abort_after_finish() {
1021 let mut broadcast = moq_net::broadcast::Info::new().produce();
1022 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
1023 let consumer = broadcast.consume();
1024 let options = Options {
1025 track: Some("audio".to_string()),
1026 ..Options::default()
1027 };
1028 let mut producer = Producer::new(
1029 &mut broadcast,
1030 catalog,
1031 Input {
1032 layout: Layout::Mono,
1033 ..Input::default()
1034 },
1035 &options,
1036 )
1037 .unwrap();
1038 let mut track = moq_mux::container::Consumer::new(
1039 consumer
1040 .track("audio")
1041 .unwrap()
1042 .subscribe(moq_net::track::Subscription::default())
1043 .await
1044 .unwrap(),
1045 moq_mux::container::legacy::Wire(moq_mux::container::Kind::Audio),
1046 );
1047
1048 producer.write(&full_frame(0)).unwrap();
1049 producer.finish().unwrap();
1050 assert!(track.read().await.unwrap().is_some());
1051 producer.abort(moq_net::Error::Cancel);
1052 }
1053
1054 #[tokio::test]
1057 async fn write_after_finish_is_closed() {
1058 let mut broadcast = moq_net::broadcast::Info::new().produce();
1059 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
1060 let options = Options {
1061 track: Some("audio".to_string()),
1062 ..Options::default()
1063 };
1064 let mut producer = Producer::new(
1065 &mut broadcast,
1066 catalog,
1067 Input {
1068 layout: Layout::Mono,
1069 ..Input::default()
1070 },
1071 &options,
1072 )
1073 .unwrap();
1074
1075 producer.finish().unwrap();
1076 let err = producer.write(&pcm_frame(&[0.1; 100], 0)).unwrap_err();
1077 assert!(matches!(err, Error::Net(moq_net::Error::Closed)));
1078 producer.abort(moq_net::Error::Cancel);
1079 }
1080
1081 #[tokio::test]
1085 async fn default_options_derive_the_track_name() {
1086 let mut broadcast = moq_net::broadcast::Info::new().produce();
1087 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
1088
1089 let first = Producer::new(&mut broadcast, catalog.clone(), Input::default(), &Options::default()).unwrap();
1090 assert_eq!(first.track_name(), "0.opus");
1091
1092 let second = Producer::new(&mut broadcast, catalog, Input::default(), &Options::default()).unwrap();
1093 assert_eq!(second.track_name(), "1.opus");
1094 }
1095}