1#[cfg(feature = "capture")]
12use std::time::Instant;
13
14use moq_mux::catalog::hang::CatalogExt;
15#[cfg(feature = "capture")]
16use moq_mux::rate::{Control, Policy};
17#[cfg(any(feature = "capture", test))]
18use moq_net::Timestamp;
19
20use crate::Error;
21#[cfg(feature = "capture")]
22use crate::Rate;
23#[cfg(feature = "capture")]
24use crate::capture;
25
26use super::Encoded;
27#[cfg(feature = "capture")]
28use super::Sink;
29#[cfg(any(feature = "capture", test))]
30use super::encoder;
31#[cfg(feature = "capture")]
32use super::encoder::Codec;
33
34#[cfg(feature = "capture")]
36const DEFAULT_FRAMERATE: Rate = Rate::integer(30);
37
38fn rendition_hint(rendition: hang::catalog::VideoConfig) -> moq_mux::catalog::VideoHint {
40 let mut hint = moq_mux::catalog::VideoHint::default();
41 hint.codec = Some(rendition.codec);
42 hint.coded_width = rendition.coded_width;
43 hint.coded_height = rendition.coded_height;
44 hint.display_aspect_width = rendition.display_aspect_width;
45 hint.display_aspect_height = rendition.display_aspect_height;
46 hint.framerate = rendition.framerate;
47 hint.bitrate = rendition.bitrate;
48 hint.optimize_for_latency = rendition.optimize_for_latency;
49 hint.container = rendition.container;
52 hint
53}
54
55enum Codecs {
58 H264 {
59 split: moq_mux::codec::h264::Split,
60 import: moq_mux::codec::h264::Import,
61 },
62 H265 {
63 split: moq_mux::codec::h265::Split,
64 import: moq_mux::codec::h265::Import,
65 },
66}
67
68pub struct Producer<E: CatalogExt = ()> {
79 codecs: Codecs,
80 _ext: std::marker::PhantomData<fn() -> E>,
81}
82
83impl<E: CatalogExt> Producer<E> {
84 pub fn new(
96 broadcast: moq_net::broadcast::Producer,
97 catalog: moq_mux::catalog::Producer<E>,
98 rendition: hang::catalog::VideoConfig,
99 ) -> Result<Self, Error> {
100 let suffix = match &rendition.codec {
101 hang::catalog::VideoCodec::H264(_) => ".avc3",
102 hang::catalog::VideoCodec::H265(_) => ".hev1",
103 other => {
104 return Err(Error::Codec(anyhow::anyhow!(
105 "{other} is not a codec this producer can publish"
106 )));
107 }
108 };
109 let track = broadcast.unique_track(suffix, catalog.track_info(hang::catalog::PRIORITY.video))?;
110 Self::with_track(track, catalog, rendition)
111 }
112
113 pub fn with_track(
118 track: moq_net::track::Producer,
119 catalog: moq_mux::catalog::Producer<E>,
120 rendition: hang::catalog::VideoConfig,
121 ) -> Result<Self, Error> {
122 let codecs = match &rendition.codec {
123 hang::catalog::VideoCodec::H264(_) => Codecs::H264 {
124 split: moq_mux::codec::h264::Split::new(),
125 import: moq_mux::codec::h264::Import::new(track, catalog.reserve(), rendition_hint(rendition))?,
126 },
127 hang::catalog::VideoCodec::H265(_) => Codecs::H265 {
128 split: moq_mux::codec::h265::Split::new(),
129 import: moq_mux::codec::h265::Import::new(track, catalog.reserve(), rendition_hint(rendition))?,
130 },
131 other => {
133 return Err(Error::Codec(anyhow::anyhow!(
134 "{other} is not a codec this producer can publish"
135 )));
136 }
137 };
138 Ok(Self {
139 codecs,
140 _ext: std::marker::PhantomData,
141 })
142 }
143
144 pub fn demand(&self) -> moq_net::track::Demand {
148 match &self.codecs {
149 Codecs::H264 { import, .. } => import.demand(),
150 Codecs::H265 { import, .. } => import.demand(),
151 }
152 }
153
154 pub fn publish(&mut self, encoded: &[Encoded]) -> Result<(), Error> {
157 for frame in encoded {
158 let timestamp = Some(frame.timestamp);
159 match &mut self.codecs {
161 Codecs::H264 { split, import } => {
162 let mut frames = split.decode(&frame.payload, timestamp)?;
163 frames.extend(split.flush(timestamp)?);
164 import.decode(frames)?;
165 import.flush(frame.timestamp, std::time::Instant::now())?;
166 }
167 Codecs::H265 { split, import } => {
168 let mut frames = split.decode(&frame.payload, timestamp)?;
169 frames.extend(split.flush(timestamp)?);
170 import.decode(frames)?;
171 import.flush(frame.timestamp, std::time::Instant::now())?;
172 }
173 }
174 }
175 Ok(())
176 }
177
178 pub fn observe_lag(&mut self, lag: std::time::Duration) -> Result<(), Error> {
180 match &mut self.codecs {
181 Codecs::H264 { import, .. } => import.observe_lag(lag)?,
182 Codecs::H265 { import, .. } => import.observe_lag(lag)?,
183 }
184 Ok(())
185 }
186
187 pub fn tick(&mut self) -> Result<(), Error> {
189 match &mut self.codecs {
190 Codecs::H264 { import, .. } => import.tick()?,
191 Codecs::H265 { import, .. } => import.tick()?,
192 }
193 Ok(())
194 }
195
196 pub fn idle(&mut self) -> Result<(), Error> {
198 match &mut self.codecs {
199 Codecs::H264 { import, .. } => import.idle()?,
200 Codecs::H265 { import, .. } => import.idle()?,
201 }
202 Ok(())
203 }
204
205 pub fn discontinuity(&mut self) -> Result<(), Error> {
213 match &mut self.codecs {
214 Codecs::H264 { import, .. } => import.discontinuity()?,
215 Codecs::H265 { import, .. } => import.discontinuity()?,
216 }
217 Ok(())
218 }
219
220 pub fn finish(&mut self) -> Result<(), Error> {
225 match &mut self.codecs {
226 Codecs::H264 { import, .. } => import.finish()?,
227 Codecs::H265 { import, .. } => import.finish()?,
228 }
229 Ok(())
230 }
231
232 pub fn abort(self, err: moq_net::Error) {
237 match self.codecs {
238 Codecs::H264 { import, .. } => import.abort(err),
239 Codecs::H265 { import, .. } => import.abort(err),
240 }
241 }
242}
243
244#[derive(Clone, Debug, Default)]
252#[non_exhaustive]
253#[cfg(feature = "capture")]
254pub struct Options {
255 pub bitrate: Option<moq_net::bandwidth::Rate>,
261 pub codec: Codec,
263 pub kind: encoder::Kind,
265 pub bandwidth: moq_net::bandwidth::Allocator,
280}
281
282#[cfg(feature = "capture")]
294pub async fn publish_capture<E: CatalogExt>(
295 broadcast: moq_net::broadcast::Producer,
296 catalog: moq_mux::catalog::Producer<E>,
297 capture: capture::Config,
298 encode: Options,
299 clock: moq_mux::Clock,
300) -> Result<(), Error> {
301 let rendition = {
306 let camera = capture::open(&capture).await?;
307 let mut probe_config = encoder::Config::new(
308 camera.width(),
309 camera.height(),
310 capture
311 .framerate
312 .or_else(|| camera.framerate())
313 .unwrap_or(DEFAULT_FRAMERATE),
314 );
315 probe_config.bitrate = encode.bitrate;
316 probe_config.codec = encode.codec;
317 probe_config.kind = encode.kind.clone();
318 probe_config.color = camera.color();
319 probe_config.probe().await?
320 };
321
322 let mut producer = Producer::new(broadcast, catalog, rendition)?;
323 let demand = producer.demand();
324
325 let result = capture_loop(&mut producer, &demand, &mut DeviceSource, &capture, &encode, &clock).await;
326
327 match &result {
331 Ok(()) => {
333 if let Err(err) = producer.finish() {
334 tracing::debug!(error = %err, "video track finish after capture ended");
335 }
336 }
337 Err(err) => producer.abort(moq_net::Error::Transport(err.to_string())),
339 }
340 result
341}
342
343#[cfg(all(feature = "capture", not(target_os = "macos")))]
349#[allow(dead_code)]
350fn assert_publish_capture_send(
351 broadcast: moq_net::broadcast::Producer,
352 catalog: moq_mux::catalog::Producer,
353 capture: capture::Config,
354 encode: Options,
355 clock: moq_mux::Clock,
356) {
357 fn is_send<T: Send>(_: &T) {}
358 is_send(&publish_capture(broadcast, catalog, capture, encode, clock));
359}
360
361#[cfg(feature = "capture")]
364trait CaptureSource {
365 async fn open(&mut self, config: &capture::Config) -> Result<capture::Stream, Error>;
366}
367
368#[cfg(feature = "capture")]
369struct DeviceSource;
370
371#[cfg(feature = "capture")]
372impl CaptureSource for DeviceSource {
373 async fn open(&mut self, config: &capture::Config) -> Result<capture::Stream, Error> {
374 capture::open(config).await
375 }
376}
377
378#[cfg(feature = "capture")]
384type RateControl = Option<(moq_net::bandwidth::Consumer, Control)>;
385
386#[cfg(feature = "capture")]
391async fn next_estimate(rate: &mut RateControl) -> Option<Option<moq_net::bandwidth::Rate>> {
392 match rate {
393 Some((bandwidth, _)) => bandwidth.changed().await.ok(),
394 None => std::future::pending().await,
396 }
397}
398
399#[cfg(feature = "capture")]
405async fn apply_estimate(
406 encoder: &mut Sink,
407 rate: &mut RateControl,
408 estimate: Option<Option<moq_net::bandwidth::Rate>>,
409) {
410 let Some((_, control)) = rate.as_mut() else { return };
411
412 let Some(estimate) = estimate else {
413 tracing::debug!("bandwidth estimate ended; holding the current encoder bitrate");
414 *rate = None;
415 return;
416 };
417
418 let Some(bitrate) = control.update(estimate, Instant::now()) else {
419 return;
420 };
421
422 match encoder.set_bitrate(bitrate).await {
423 Ok(()) => tracing::debug!(bitrate = bitrate.as_bps(), "adjusted encoder bitrate"),
424 Err(Error::BitrateUnsupported(name)) => {
428 tracing::warn!(encoder = name, "encoder cannot follow the bandwidth estimate");
429 *rate = None;
430 }
431 Err(err) => tracing::warn!(error = %err, bitrate = bitrate.as_bps(), "failed to adjust encoder bitrate"),
435 }
436}
437
438#[cfg(feature = "capture")]
442fn log_track_ended(err: moq_net::Error) {
443 if matches!(err, moq_net::Error::Dropped | moq_net::Error::Closed) {
444 tracing::debug!("video track no longer announced; stopping capture");
445 } else {
446 tracing::warn!(error = %err, "video track aborted; stopping capture");
447 }
448}
449
450#[cfg(any(feature = "capture", all(test, feature = "openh264")))]
451fn capture_stopped<E: CatalogExt>(producer: &mut Producer<E>) -> Result<(), Error> {
452 producer.discontinuity()
455}
456
457#[cfg(feature = "capture")]
460async fn wait_capture<E: CatalogExt, T>(
461 producer: &mut Producer<E>,
462 demand: &moq_net::track::Demand,
463 work: impl std::future::Future<Output = Result<T, Error>>,
464) -> Result<Option<T>, Error> {
465 let mut work = std::pin::pin!(work);
466 let mut timer = tokio::time::interval(hang::catalog::stalled::DEFAULT_INTERVAL);
467 loop {
468 tokio::select! {
469 biased;
470 res = demand.unused() => {
471 if let Err(err) = res {
472 log_track_ended(err);
473 }
474 producer.idle()?;
475 return Ok(None);
476 }
477 _ = timer.tick() => producer.tick()?,
478 res = &mut work => return res.map(Some),
479 }
480 }
481}
482
483#[cfg(feature = "capture")]
493async fn capture_loop<E: CatalogExt, S: CaptureSource>(
494 producer: &mut Producer<E>,
495 demand: &moq_net::track::Demand,
496 source: &mut S,
497 capture: &capture::Config,
498 encode: &Options,
499 clock: &moq_mux::Clock,
500) -> Result<(), Error> {
501 let mut reservation: Option<moq_net::bandwidth::Reservation> = None;
505
506 loop {
507 if let Err(err) = demand.used().await {
511 log_track_ended(err);
512 return Ok(());
513 }
514
515 let Some(mut camera) = wait_capture(producer, demand, source.open(capture)).await? else {
517 continue;
518 };
519 let capture_epoch =
523 u64::try_from(clock.now().as_micros().saturating_sub(camera.now().as_micros())).unwrap_or(u64::MAX);
524 let framerate = capture
527 .framerate
528 .or_else(|| camera.framerate())
529 .unwrap_or(DEFAULT_FRAMERATE);
530 let mut encoder_config = encoder::Config::new(camera.width(), camera.height(), framerate);
531 encoder_config.bitrate = encode.bitrate;
532 encoder_config.codec = encode.codec;
533 encoder_config.kind = encode.kind.clone();
534 encoder_config.color = camera.color();
535 let Some(mut encoder) = wait_capture(producer, demand, Sink::open(&encoder_config)).await? else {
540 continue;
541 };
542 tracing::info!(encoder = encoder.name(), device = camera.label(), "capturing");
543
544 let ceiling = encoder_config.resolved_bitrate();
549 let reservation = reservation.get_or_insert_with(|| encode.bandwidth.reserve(demand, ceiling));
550 reservation.update(ceiling);
551
552 let mut rate = Some((reservation.consumer(), Control::new(Policy::new(ceiling))));
557
558 loop {
559 let interval = hang::catalog::stalled::interval_from_fps(Some(framerate.as_f64()));
563 let frame = tokio::select! {
564 biased;
565 res = demand.unused() => {
566 if let Err(err) = res {
567 log_track_ended(err);
568 return Ok(());
569 }
570 break; }
572 estimate = next_estimate(&mut rate) => {
575 apply_estimate(&mut encoder, &mut rate, estimate).await;
576 continue;
577 }
578 frame = tokio::time::timeout(interval, camera.read()) => match frame {
582 Ok(frame) => frame?,
583 Err(_) => {
584 producer.tick()?;
585 continue;
586 }
587 },
588 };
589
590 let Some(mut frame) = frame else { break };
591 frame.timestamp = map_capture_timestamp(capture_epoch, frame.timestamp)?;
592 let started = Instant::now();
593 let Some(encoded) = wait_capture(producer, demand, encoder.encode(frame)).await? else {
594 break;
595 };
596 let lag = started.elapsed();
597 producer.observe_lag(lag)?;
598 producer.publish(&encoded)?;
599 }
600
601 drop(camera);
603 drop(encoder);
604 producer.idle()?;
605 capture_stopped(producer)?;
606 tracing::info!("capture stopped; released source");
607 }
608}
609
610#[cfg(feature = "capture")]
611fn map_capture_timestamp(epoch_micros: u64, timestamp: Timestamp) -> Result<Timestamp, Error> {
612 let capture_micros = u64::try_from(timestamp.as_micros()).unwrap_or(u64::MAX);
613 Ok(Timestamp::from_micros(epoch_micros.saturating_add(capture_micros))?)
614}
615
616#[cfg(test)]
617mod tests {
618 #![cfg_attr(not(feature = "openh264"), allow(dead_code, unused_imports))]
619
620 use moq_mux::catalog::Stream as _;
621
622 use super::*;
623 use crate::Frame;
624 use crate::encode::{Codec, Config, Encoder};
625
626 #[cfg(feature = "capture")]
627 #[test]
628 fn capture_clock_mapping_is_monotonic() {
629 let first = map_capture_timestamp(10_000, Timestamp::from_micros(2_000).unwrap()).unwrap();
630 let second = map_capture_timestamp(10_000, Timestamp::from_micros(2_001).unwrap()).unwrap();
631 assert!(second > first);
632 }
633
634 async fn roundtrip_rendition(codec: Codec, kind: encoder::Kind) -> (String, hang::catalog::VideoConfig) {
644 let mut broadcast = moq_net::broadcast::Info::new().produce();
645 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
646
647 let mut config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
648 config.codec = codec;
649 config.kind = kind;
650
651 let mut producer = Producer::new(broadcast, catalog.clone(), config.probe().await.unwrap()).unwrap();
652 let advertised = rendition(&catalog).expect("the rendition publishes before any frame").1;
653
654 let mut encoder = Encoder::new(&config).unwrap();
655 assert_eq!(encoder.codec(), codec);
656
657 let rgba = vec![0x80u8; 320 * 240 * 4];
658 for i in 0..10u64 {
659 let surface = crate::Surface::rgba(&rgba, crate::Size::new(320, 240)).unwrap();
660 let frame = Frame::new(surface, Timestamp::from_micros(i * 33_333).unwrap());
661 producer.publish(&encoder.encode(&frame).unwrap()).unwrap();
662 }
663 producer.publish(&encoder.finish().unwrap()).unwrap();
664
665 let (name, resolved) = rendition(&catalog).expect("the importer should have registered a video rendition");
666 let (mut before, mut after) = (advertised, resolved.clone());
668 (before.jitter, before.delay) = (None, None);
669 (after.jitter, after.delay) = (None, None);
670 assert_eq!(
671 before, after,
672 "the first keyframe should confirm the advertised rendition, not correct it"
673 );
674 (name, resolved)
675 }
676
677 fn rendition(catalog: &moq_mux::catalog::Producer) -> Option<(String, hang::catalog::VideoConfig)> {
679 let snapshot = catalog.snapshot();
680 let (name, config) = snapshot.video.renditions.iter().next()?;
681 Some((name.clone(), config.clone()))
682 }
683
684 async fn collect_groups(mut consumer: moq_net::track::Subscriber) -> Vec<usize> {
685 let mut groups = Vec::new();
686 while let Some(mut group) = consumer.recv_group().await.unwrap() {
687 let mut frames = 0;
688 while group.next_frame().await.unwrap().is_some() {
689 frames += 1;
690 }
691 groups.push(frames);
692 }
693 groups
694 }
695
696 #[tokio::test]
700 #[cfg(feature = "openh264")]
701 async fn idle_capture_publishes_a_discontinuity_before_resume() {
702 let mut broadcast = moq_net::broadcast::Info::new().produce();
703 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
704 let replay = std::time::Duration::from_secs(11);
707 let track = broadcast
708 .create_track(
709 "video",
710 catalog.track_info(hang::catalog::PRIORITY.video).with_max_age(replay),
711 )
712 .unwrap();
713 let consumer = track.subscribe(moq_net::track::Subscription::default().with_max_age(replay));
714
715 let mut config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
716 config.kind = encoder::Kind::Software;
717 let mut producer = Producer::with_track(track, catalog, config.probe().await.unwrap()).unwrap();
718 let mut encoder = Encoder::new(&config).unwrap();
719 let rgba = vec![0x80u8; 320 * 240 * 4];
720
721 for timestamp in [0, 10_000_000] {
722 if timestamp > 0 {
723 capture_stopped(&mut producer).unwrap();
724 }
725 encoder.cut().unwrap();
726 let surface = crate::Surface::rgba(&rgba, crate::Size::new(320, 240)).unwrap();
727 let frame = Frame::new(surface, Timestamp::from_micros(timestamp).unwrap());
728 producer.publish(&encoder.encode(&frame).unwrap()).unwrap();
729 }
730 producer.finish().unwrap();
731
732 assert_eq!(collect_groups(consumer).await, vec![1, 1, 1]);
733 }
734
735 #[tokio::test]
736 #[cfg(feature = "openh264")]
737 async fn source_resize_updates_the_published_rendition() {
738 let mut broadcast = moq_net::broadcast::Info::new().produce();
739 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
740 let mut initial = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
741 initial.kind = encoder::Kind::Software;
742 let mut producer = Producer::new(broadcast, catalog.clone(), initial.probe().await.unwrap()).unwrap();
743
744 for (timestamp, config) in [
745 (0, initial),
746 (33_333, Config::new(640, 360, crate::Rate::new(30, 1).unwrap())),
747 ] {
748 let mut config = config;
749 config.kind = encoder::Kind::Software;
750 let mut encoder = Encoder::new(&config).unwrap();
751 encoder.cut().unwrap();
752 let rgba = vec![0x80u8; usize::try_from(config.width * config.height * 4).unwrap()];
753 let surface = crate::Surface::rgba(&rgba, crate::Size::new(config.width, config.height)).unwrap();
754 let frame = Frame::new(surface, Timestamp::from_micros(timestamp).unwrap());
755 producer.publish(&encoder.encode(&frame).unwrap()).unwrap();
756 capture_stopped(&mut producer).unwrap();
757 }
758
759 let (_, rendition) = rendition(&catalog).expect("the resized rendition should be published");
760 assert_eq!(rendition.coded_width, Some(640));
761 assert_eq!(rendition.coded_height, Some(360));
762 }
763
764 #[tokio::test]
770 #[cfg(feature = "openh264")]
771 async fn a_selected_container_survives_the_rendition_hint() {
772 let mut broadcast = moq_net::broadcast::Info::new().produce();
773 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
774
775 let mut config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
776 config.kind = encoder::Kind::Software;
778 let mut selected = config.probe().await.unwrap();
779 selected.container = hang::catalog::Container::Loc;
780
781 let _producer = Producer::new(broadcast, catalog.clone(), selected).unwrap();
782
783 let (_, published) = rendition(&catalog).expect("the rendition publishes before any frame");
784 assert_eq!(published.container, hang::catalog::Container::Loc);
785 }
786
787 #[tokio::test]
795 #[cfg(feature = "openh264")]
796 async fn the_rendition_reaches_the_wire_before_the_first_frame() {
797 let mut broadcast = moq_net::broadcast::Info::new().produce();
798 let consumer = broadcast.consume();
799 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
800
801 let mut config = Config::new(1920, 1080, crate::Rate::new(30, 1).unwrap());
802 config.bitrate = Some(moq_net::bandwidth::Rate::from_mbps(6));
803 config.kind = encoder::Kind::Software;
805 let _producer = Producer::new(broadcast, catalog, config.probe().await.unwrap()).unwrap();
806
807 let mut stream = moq_mux::catalog::Consumer::<()>::new(&consumer, moq_mux::catalog::CatalogFormat::Hang)
809 .await
810 .unwrap();
811 let snapshot = stream.next().await.unwrap().expect("a catalog before any frame");
812
813 let (name, rendition) = snapshot
814 .video
815 .renditions
816 .iter()
817 .next()
818 .expect("the track must be discoverable before it has encoded anything");
819 assert!(name.ends_with(".avc3"));
820
821 let hang::catalog::VideoCodec::H264(h264) = &rendition.codec else {
824 panic!("expected H.264, got {}", rendition.codec)
825 };
826 assert!(h264.inline, "an avc3 track carries its parameter sets in band");
827 assert_eq!(rendition.coded_width, Some(1920));
828 assert_eq!(rendition.coded_height, Some(1080));
829 assert_eq!(rendition.framerate, Some(30.0));
831 assert_eq!(rendition.bitrate, Some(6_000_000));
832 }
833
834 #[tokio::test]
836 #[cfg(feature = "openh264")]
837 async fn abort_after_finish() {
838 let mut broadcast = moq_net::broadcast::Info::new().produce();
839 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
840 let mut config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
841 config.kind = encoder::Kind::Software;
842 let track = broadcast
843 .create_track("video", catalog.track_info(hang::catalog::PRIORITY.video))
844 .unwrap();
845 let mut subscriber = track.subscribe(None);
846 let mut producer = Producer::with_track(track, catalog, config.probe().await.unwrap()).unwrap();
847 let mut encoder = Encoder::new(&config).unwrap();
848 let rgba = vec![0x80u8; 320 * 240 * 4];
849 let surface = crate::Surface::rgba(&rgba, crate::Size::new(320, 240)).unwrap();
850 let frame = Frame::new(surface, Timestamp::from_micros(0).unwrap());
851 producer.publish(&encoder.encode(&frame).unwrap()).unwrap();
852
853 producer.finish().unwrap();
854 assert!(subscriber.recv_group().await.unwrap().is_some());
855 producer.abort(moq_net::Error::Cancel);
856 }
857
858 #[tokio::test]
859 #[cfg(feature = "openh264")]
860 async fn h264_roundtrip_publishes_avc3() {
861 let (name, config) = roundtrip_rendition(Codec::H264, encoder::Kind::Software).await;
864 assert!(name.ends_with(".avc3"));
865 assert_eq!(config.coded_width, Some(320));
866 assert_eq!(config.coded_height, Some(240));
867 }
868
869 #[cfg(target_os = "macos")]
872 #[tokio::test]
873 async fn h265_roundtrip_publishes_hev1() {
874 let (name, config) = roundtrip_rendition(Codec::H265, encoder::Kind::Hardware).await;
875 assert!(name.ends_with(".hev1"));
876 assert_eq!(config.coded_width, Some(320));
877 assert_eq!(config.coded_height, Some(240));
878 }
879
880 #[cfg(all(feature = "capture", feature = "openh264"))]
887 mod clock {
888 use std::time::{Duration, Instant, SystemTime};
889
890 use super::*;
891 use crate::capture::Synthetic;
892
893 const SAMPLING: Duration = Duration::from_millis(250);
895 const ROUNDING: u64 = 2;
897 const RETAIN: Duration = Duration::from_secs(600);
899
900 struct Opens(tokio::sync::mpsc::UnboundedReceiver<capture::Stream>);
902
903 impl CaptureSource for Opens {
904 async fn open(&mut self, _config: &capture::Config) -> Result<capture::Stream, Error> {
905 self.0
906 .recv()
907 .await
908 .ok_or_else(|| Error::SourceUnavailable("the fixture stopped opening cameras".to_string()))
909 }
910 }
911
912 struct Fixture {
913 epoch: Instant,
914 clock: moq_mux::Clock,
915 catalog: moq_mux::catalog::Producer,
916 consumer: moq_net::broadcast::Consumer,
917 _broadcast: moq_net::broadcast::Producer,
918 opens: tokio::sync::mpsc::UnboundedSender<capture::Stream>,
919 stop: Option<tokio::sync::oneshot::Sender<()>>,
920 task: tokio::task::JoinHandle<Result<(), Error>>,
921 }
922
923 impl Fixture {
924 async fn start(behind: Duration, wall: SystemTime) -> Self {
926 let epoch = Instant::now()
927 .checked_sub(behind)
928 .expect("a monotonic clock that far back");
929 let clock = moq_mux::Clock::at(epoch, wall).unwrap();
930 let mut broadcast = moq_net::broadcast::Info::new().produce();
931 let consumer = broadcast.consume();
932 let config = moq_mux::catalog::Config::default()
933 .with_clock(clock)
934 .with_max_age(RETAIN);
935 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config).unwrap();
936 let track = broadcast
937 .create_track(
938 "video",
939 catalog.track_info(hang::catalog::PRIORITY.video).with_max_age(RETAIN),
940 )
941 .unwrap();
942
943 let mut probe = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
944 probe.kind = encoder::Kind::Software;
945 let mut producer = Producer::with_track(track, catalog.clone(), probe.probe().await.unwrap()).unwrap();
946 let demand = producer.demand();
947
948 let (opens, rx) = tokio::sync::mpsc::unbounded_channel();
949 let (stop, stopped) = tokio::sync::oneshot::channel::<()>();
950 let task = tokio::task::spawn_local(async move {
952 let mut source = Opens(rx);
953 let options = Options {
954 kind: encoder::Kind::Software,
955 ..Options::default()
956 };
957 let config = capture::Config::default();
958 tokio::select! {
959 res = capture_loop(&mut producer, &demand, &mut source, &config, &options, &clock) => res?,
960 _ = stopped => {}
961 }
962 producer.finish()
963 });
964
965 Self {
966 epoch,
967 clock,
968 catalog,
969 consumer,
970 _broadcast: broadcast,
971 opens,
972 stop: Some(stop),
973 task,
974 }
975 }
976
977 async fn subscribe(&self) -> moq_mux::container::Consumer<moq_mux::catalog::hang::Container> {
979 let snapshot = self.catalog.snapshot();
980 let (name, rendition) = snapshot.video.renditions.iter().next().expect("the probed rendition");
981 let container = moq_mux::catalog::hang::Container::try_from(rendition).unwrap();
982 let track = self
983 .consumer
984 .track(name)
985 .unwrap()
986 .subscribe(moq_net::track::Subscription::default().with_max_age(RETAIN))
987 .await
988 .unwrap();
989 moq_mux::container::Consumer::new(track, container)
990 }
991
992 fn camera(&self) -> Synthetic {
994 let (camera, stream) = Synthetic::open(crate::Size::new(320, 240), crate::Rate::new(30, 1).unwrap());
995 self.opens.send(stream).unwrap();
996 camera
997 }
998
999 fn at(&self, instant: Instant) -> u64 {
1001 u64::try_from(instant.duration_since(self.epoch).as_micros()).unwrap()
1002 }
1003
1004 async fn finish(mut self) -> (moq_mux::catalog::Producer, moq_net::broadcast::Consumer) {
1006 let _ = self.stop.take().expect("finished once").send(());
1007 self.task.await.unwrap().unwrap();
1008 (self.catalog, self.consumer)
1009 }
1010
1011 fn assert_acquired(&self, published: u64, captured: Instant) {
1013 let exact = self.at(captured);
1014 let early = u64::try_from(SAMPLING.as_micros()).unwrap();
1015 assert!(
1016 published + early >= exact && published <= exact + ROUNDING,
1017 "published {published}us, acquired at {exact}us on the broadcast clock"
1018 );
1019 }
1020 }
1021
1022 fn surface() -> crate::frame::Surface {
1023 crate::frame::Surface::I420(crate::frame::I420 {
1024 width: 320,
1025 height: 240,
1026 data: vec![0x80; 320 * 240 * 3 / 2],
1027 color: None,
1028 })
1029 }
1030
1031 fn us(micros: u64) -> Timestamp {
1032 Timestamp::from_micros(micros).unwrap()
1033 }
1034
1035 async fn read(track: &mut moq_mux::container::Consumer<moq_mux::catalog::hang::Container>) -> u64 {
1036 let frame = track.read().await.unwrap().expect("a published frame");
1037 u64::try_from(frame.timestamp.as_micros()).unwrap()
1038 }
1039
1040 async fn read_new(
1042 track: &mut moq_mux::container::Consumer<moq_mux::catalog::hang::Container>,
1043 seen: &[u64],
1044 ) -> u64 {
1045 loop {
1046 let timestamp = read(track).await;
1047 if !seen.contains(×tamp) {
1048 return timestamp;
1049 }
1050 }
1051 }
1052
1053 #[tokio::test]
1056 async fn a_late_first_frame_publishes_its_acquisition() {
1057 tokio::task::LocalSet::new()
1058 .run_until(async {
1059 let fixture = Fixture::start(Duration::from_secs(5), SystemTime::now()).await;
1060 let mut track = fixture.subscribe().await;
1061 let camera = fixture.camera();
1062
1063 let captured = Instant::now();
1064 tokio::time::sleep(Duration::from_millis(50)).await;
1066 camera.push_at(surface(), captured);
1067 let published = read(&mut track).await;
1068
1069 assert!(published >= 4_000_000, "{published}us restarted the broadcast at zero");
1070 fixture.assert_acquired(published, captured);
1071 fixture.finish().await;
1072 })
1073 .await
1074 }
1075
1076 #[tokio::test]
1079 async fn a_device_clock_restart_continues_forward() {
1080 tokio::task::LocalSet::new()
1081 .run_until(async {
1082 let fixture = Fixture::start(Duration::from_secs(1), SystemTime::now()).await;
1083 let mut track = fixture.subscribe().await;
1084 let camera = fixture.camera();
1085
1086 camera.push_native(surface(), us(0));
1088 let first = read(&mut track).await;
1089 tokio::time::sleep(Duration::from_millis(40)).await;
1090 camera.push_native(surface(), us(40_000));
1091 let second = read(&mut track).await;
1092 assert_eq!(second - first, 40_000, "the device's spacing survives");
1093
1094 camera.push_native(surface(), us(0));
1096 let restarted = read(&mut track).await;
1097 assert!(restarted >= second, "{restarted}us rewound behind {second}us");
1098 tokio::time::sleep(Duration::from_millis(40)).await;
1099 camera.push_native(surface(), us(40_000));
1100 let resumed = read(&mut track).await;
1101 assert_eq!(resumed - restarted, 40_000, "the device's spacing resumes");
1102
1103 camera.close();
1105 let camera = fixture.camera();
1106 let pushed = Instant::now();
1107 camera.push_native(surface(), us(0));
1108 let reopened = read(&mut track).await;
1109 let arrived = fixture.at(Instant::now());
1110 assert!(reopened >= resumed, "{reopened}us rewound across the reopen");
1111 let early = u64::try_from(SAMPLING.as_micros()).unwrap();
1112 assert!(reopened + early >= fixture.at(pushed) && reopened <= arrived + ROUNDING);
1113 fixture.finish().await;
1114 })
1115 .await
1116 }
1117
1118 #[tokio::test]
1121 async fn a_restart_after_idle_keeps_the_gap() {
1122 tokio::task::LocalSet::new()
1123 .run_until(async {
1124 let idle = Duration::from_millis(300);
1125 let fixture = Fixture::start(Duration::from_secs(1), SystemTime::now()).await;
1126
1127 let mut track = fixture.subscribe().await;
1128 let camera = fixture.camera();
1129 let captured = Instant::now();
1130 camera.push_at(surface(), captured);
1131 let before = read(&mut track).await;
1132 fixture.assert_acquired(before, captured);
1133 drop(track);
1134 drop(camera);
1135
1136 tokio::time::sleep(idle).await;
1137
1138 let mut track = fixture.subscribe().await;
1139 let camera = fixture.camera();
1140 let captured = Instant::now();
1141 camera.push_at(surface(), captured);
1142 let after = read_new(&mut track, &[before]).await;
1143 fixture.assert_acquired(after, captured);
1144 assert!(
1145 after - before >= u64::try_from(idle.as_micros()).unwrap(),
1146 "the {idle:?} idle gap collapsed to {}us",
1147 after - before
1148 );
1149 fixture.finish().await;
1150 })
1151 .await
1152 }
1153
1154 #[tokio::test]
1157 async fn a_system_wall_adjustment_retimes_nothing() {
1158 tokio::task::LocalSet::new()
1159 .run_until(async {
1160 let now = SystemTime::now().duration_since(SystemTime::UNIX_EPOCH).unwrap();
1162 let wall = SystemTime::UNIX_EPOCH + Duration::from_secs(now.as_secs() - 3600);
1163 let fixture = Fixture::start(Duration::from_secs(1), wall).await;
1164 let advertised = fixture.catalog.snapshot().clock;
1165 assert_eq!(advertised, Some(fixture.clock.wall()));
1166
1167 let mut track = fixture.subscribe().await;
1168 let camera = fixture.camera();
1169 let captured = Instant::now();
1170 camera.push_at(surface(), captured);
1171 let published = read(&mut track).await;
1172
1173 fixture.assert_acquired(published, captured);
1175 let mapped = advertised.unwrap().wall_clock(us(published)).unwrap();
1176 assert_eq!(mapped, wall + Duration::from_millis(published / 1000));
1178 assert_eq!(fixture.catalog.snapshot().clock, advertised);
1179 fixture.finish().await;
1180 })
1181 .await
1182 }
1183
1184 #[tokio::test]
1187 async fn retained_archive_playback_keeps_the_live_timestamps() {
1188 tokio::task::LocalSet::new()
1189 .run_until(async {
1190 let fixture = Fixture::start(Duration::from_secs(1), SystemTime::now()).await;
1191 let section = fixture
1192 .catalog
1193 .snapshot()
1194 .archive
1195 .expect("the video track enrolls an archive");
1196 let mut timeline = moq_mux::timeline::Consumer::<()>::subscribe(&fixture.consumer, §ion)
1197 .await
1198 .unwrap();
1199
1200 let mut live = Vec::new();
1201 for _ in 0..2 {
1202 let mut track = fixture.subscribe().await;
1203 let camera = fixture.camera();
1204 let captured = Instant::now();
1205 camera.push_at(surface(), captured);
1206 let published = read_new(&mut track, &live).await;
1207 fixture.assert_acquired(published, captured);
1208 live.push(published);
1209 drop(track);
1210 tokio::time::sleep(moq_mux::timeline::DEFAULT_DURATION_MIN + Duration::from_millis(100)).await;
1212 }
1213
1214 let (catalog, _consumer) = fixture.finish().await;
1215 catalog.timeline().finish().unwrap();
1216 let mut archived = Vec::new();
1217 while let Some(event) = timeline.next().await.unwrap() {
1218 match event {
1219 moq_mux::timeline::Event::Push { entry, .. } => archived.push(entry),
1220 other => panic!("unexpected timeline event {other:?}"),
1221 }
1222 }
1223
1224 assert_eq!(archived.len(), live.len(), "one segment per capture run: {archived:?}");
1225 for (entry, live) in archived.iter().zip(&live) {
1226 assert_eq!(entry.pts.as_micros() / 1000, u128::from(*live / 1000), "{archived:?}");
1228 assert!(entry.tracks.contains_key("video"), "{archived:?}");
1229 }
1230 let first = &archived[0];
1231 assert!(
1232 archived[1].pts.as_micros() >= first.pts.as_micros() + first.duration.as_micros(),
1233 "the resumed segment overlaps the one before it: {archived:?}"
1234 );
1235 })
1236 .await
1237 }
1238 }
1239}