Skip to main content

moq_video/encode/
producer.rs

1//! Publish encoded video frames as a moq video track, with optional capture.
2//!
3//! Encoding is strictly on demand: the track and its catalog rendition are
4//! advertised immediately (the rendition is probed from the encoder, since
5//! nothing has been encoded yet), and the encoder itself only runs while a
6//! subscriber is watching. Capture opens its camera once at startup to learn
7//! the mode it negotiates, then keeps it closed between viewers. This mirrors
8//! `moq-boy`, which pauses its emulator on `track::Producer::used()` /
9//! `unused()`.
10
11#[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/// Last-resort framerate when neither the caller nor the camera reports one.
35#[cfg(feature = "capture")]
36const DEFAULT_FRAMERATE: Rate = Rate::integer(30);
37
38/// Convert the probed rendition into the importer hint published before the first frame.
39fn 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	// Authoritative for both the catalog entry and the wire, so dropping it would silently
50	// downgrade a caller's selection to the default.
51	hint.container = rendition.container;
52	hint
53}
54
55/// Per-codec splitter + importer pair. Each codec frames its packets and resolves
56/// its catalog rendition differently, so the producer holds one of these.
57enum 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
68/// Publishes encoded video frames as a moq track (avc3 / hev1 depending on the
69/// codec).
70///
71/// Built on the async side so the track is advertised (and the catalog
72/// registered) before the camera opens; this is what lets a subscriber
73/// trigger capture on demand. The `moq_mux::codec` importer for the codec
74/// handles catalog registration and framing.
75/// `E` is the catalog's application extension, defaulting to none. A host
76/// carrying its own catalog sections (the FFI bindings use `hang::Extra`)
77/// publishes into a catalog of the same shape.
78pub struct Producer<E: CatalogExt = ()> {
79	codecs: Codecs,
80	_ext: std::marker::PhantomData<fn() -> E>,
81}
82
83impl<E: CatalogExt> Producer<E> {
84	/// Publish a track carrying `rendition` into `broadcast`, registering it in
85	/// `catalog`. The frames fed to [`publish`](Self::publish) must be in that
86	/// codec's framing, which is what the [`Encoder`](super::Encoder) the
87	/// rendition was probed from emits.
88	///
89	/// `rendition` comes from [`Config::probe`](super::Config::probe), so it is
90	/// what the encoder will actually emit rather than a guess. It is published
91	/// immediately, before anything is encoded, which is what lets a subscriber
92	/// discover a track an on-demand encoder has not run for yet; because it
93	/// already says what the first keyframe says, that keyframe confirms the
94	/// catalog instead of correcting it.
95	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	/// Publish `rendition` on an existing track, registering it in `catalog`.
114	///
115	/// Use this when the caller owns the track name. [`new`](Self::new) derives a
116	/// unique name from the codec instead.
117	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			// Unreachable via `Config::probe`, which only encodes what `Codec` covers.
132			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	/// A watch-only handle to the track's subscriber demand, created eagerly so
145	/// subscription state is observable before any frames arrive. Watch it via
146	/// [`used`](moq_net::track::Demand::used) / [`unused`](moq_net::track::Demand::unused).
147	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	/// Publish already-encoded frames, each at its own timestamp. Each frame is one
155	/// whole access unit in the producer's codec framing.
156	pub fn publish(&mut self, encoded: &[Encoded]) -> Result<(), Error> {
157		for frame in encoded {
158			let timestamp = Some(frame.timestamp);
159			// The encoder emits one whole access unit per frame, so flush to emit it.
160			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	/// Record the encode duration before publishing its frames so the catalog can report a stall.
179	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	/// Re-evaluate stall from source silence while waiting for the next frame.
188	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	/// The camera is released; this rendition is never stalled while idle.
197	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	/// Mark a break in the published timeline: whatever is published next does not continue
206	/// what came before.
207	///
208	/// Call this when the encoder stops rather than merely pausing between frames -- a
209	/// capture that goes idle, a source switch, anything that will resume on a re-anchored
210	/// clock. See [`Producer::discontinuity`](moq_mux::container::Producer::discontinuity)
211	/// for what the marker buys a consumer.
212	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	/// Finalize the track.
221	///
222	/// Borrows rather than consumes, so a later [`abort`](Self::abort) can still
223	/// run after a successful finish.
224	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	/// Abort the track with `err` instead of finishing it cleanly, so subscribers
233	/// see the real cause rather than [`moq_net::Error::Dropped`].
234	///
235	/// Consumes the producer. Still callable after [`finish`](Self::finish).
236	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/// Source-agnostic encode knobs for [`publish_capture`], where the geometry
245/// (width / height / framerate) comes from the capture source, not the caller.
246/// For the bring-your-own-frames [`Encoder`](super::Encoder) path, where you
247/// must specify geometry, use [`Config`](super::Config) instead.
248///
249/// `#[non_exhaustive]`: construct via [`Options::default`] and set fields, so
250/// new knobs can be added without breaking callers.
251#[derive(Clone, Debug, Default)]
252#[non_exhaustive]
253#[cfg(feature = "capture")]
254pub struct Options {
255	/// Target bitrate; `None` derives one from the resolution.
256	///
257	/// This is a ceiling, not a fixed rate: with [`bandwidth`](Self::bandwidth)
258	/// set, the encoder backs off below it while the uplink is congested and
259	/// climbs back afterwards, but never exceeds it.
260	pub bitrate: Option<moq_net::bandwidth::Rate>,
261	/// Output codec. Defaults to [`Codec::H264`].
262	pub codec: Codec,
263	/// Encoder implementation preference.
264	pub kind: encoder::Kind,
265	/// The connection's bandwidth, as an allocator over
266	/// [`Session::send_bandwidth`](moq_net::Session::send_bandwidth) (or
267	/// `moq_tokio::Connection::send_bandwidth`, which survives reconnects).
268	///
269	/// Set it and the encoder reserves this track's ceiling, then tracks its share of
270	/// the estimate per the default [`moq_mux::rate::Policy`], so a closing
271	/// uplink gets a softer picture instead of a stalled one. Pass the same allocator
272	/// to every sender on the connection, including the audio side: that's what keeps
273	/// their bitrates summing to the uplink instead of each matching it.
274	///
275	/// Defaults to [`Allocator::unlimited`](moq_net::bandwidth::Allocator::unlimited),
276	/// which holds [`bitrate`](Self::bitrate) regardless of congestion. That's what you
277	/// want when the estimate isn't meaningful (a local file, a test harness) or
278	/// unavailable (a publisher that only accepts inbound sessions).
279	pub bandwidth: moq_net::bandwidth::Allocator,
280}
281
282/// Capture a webcam and publish it as an on-demand video track.
283///
284/// Returns when the broadcast is dropped (the track stops being announced)
285/// or the capture loop fails. Frames are stamped from `clock`, so passing the
286/// same [`Clock`](moq_mux::Clock) to a concurrent audio publish keeps the two
287/// tracks aligned.
288///
289/// The camera is opened once at startup to probe the mode it negotiates, then released until a
290/// subscriber arrives and reopened for as long as one is watching. That one open is what lets the
291/// catalog rendition be exact before a single frame is published, so a consumer can size itself
292/// against it (and discover the track at all) without waiting for an encoder that may never run.
293#[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	// Open the camera once to find out what it actually negotiated, since a requested size is only a
302	// hint (macOS ignores it outright) and the encoder is built from the mode, not the request. It
303	// closes again immediately: this costs one camera open at startup and buys a rendition that says
304	// exactly what the stream will carry, rather than one every consumer has to treat as provisional.
305	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	// This runs only when the loop ends on its own (the track is usually already
328	// going away by then); a Ctrl+C cancels the future before this point, since
329	// async `Drop` can't finalize the track.
330	match &result {
331		// Clean end (the track was dropped): best-effort finish.
332		Ok(()) => {
333			if let Err(err) = producer.finish() {
334				tracing::debug!(error = %err, "video track finish after capture ended");
335			}
336		}
337		// The capture loop failed: abort with the real cause so subscribers see it.
338		Err(err) => producer.abort(moq_net::Error::Transport(err.to_string())),
339	}
340	result
341}
342
343/// Off macOS, [`publish_capture`]'s future must stay `Send` so a server can
344/// `tokio::spawn` it: the encoder runs on its own thread and the capture guard
345/// is `Send` there. This is never called; it exists only to fail compilation if
346/// the future ever regains a `!Send` component. macOS is exempt (the objc
347/// capture session is `!Send`).
348#[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/// Where the capture loop opens its camera. Kept apart from the device backends so
362/// the clock fixtures can drive the real loop from a synthetic source.
363#[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/// The live rate control state: the estimate source paired with the policy tracking
379/// it. `None` once it has *retired*, which is the only thing absence means now that
380/// every encoder has a share to read: an allocator with nothing to divide grants
381/// `None` rather than being absent. Retiring stops the `select!` arm from spinning on
382/// a channel that is permanently ready.
383#[cfg(feature = "capture")]
384type RateControl = Option<(moq_net::bandwidth::Consumer, Control)>;
385
386/// Wait for the next bandwidth estimate, or forever when rate control is off or
387/// finished. Cancel-safe: [`Consumer::changed`](moq_net::bandwidth::Consumer::changed)
388/// only reads shared state, so losing this race to a frame drops no estimate,
389/// it just re-reads the latest one next time round.
390#[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		// Retired: park this arm forever so `select!` ignores it.
395		None => std::future::pending().await,
396	}
397}
398
399/// Feed an estimate through the policy and retune the encoder if it moved.
400///
401/// `None` means the producer is gone (the session ended for good), so rate
402/// control retires; a `Some(None)` estimate means the value is merely
403/// unavailable right now, which the policy holds through.
404#[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		// The encoder can't retune, so keep encoding at the rate it opened with
425		// and stop asking. Dropping the source also stops the estimate arm, which
426		// would otherwise wake this loop for nothing on every change.
427		Err(Error::BitrateUnsupported(name)) => {
428			tracing::warn!(encoder = name, "encoder cannot follow the bandwidth estimate");
429			*rate = None;
430		}
431		// A transient failure: keep the policy running so the next change retries.
432		// The policy already moved its target, so a persistent failure just means
433		// the encoder trails it; that's better than giving up on the first blip.
434		Err(err) => tracing::warn!(error = %err, bitrate = bitrate.as_bps(), "failed to adjust encoder bitrate"),
435	}
436}
437
438/// A dropped or closed track is the normal end of a publish; any other cause is
439/// a real abort (e.g. a transport reset) worth surfacing rather than treating as
440/// a clean exit.
441#[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	// The shared clock keeps advancing while capture is stopped. Mark the break before waiting
453	// for demand again so the next timestamp does not stretch the previous frame across the gap.
454	producer.discontinuity()
455}
456
457// Keep observing silence while source setup or an encode is pending. The work
458// future stays pinned across ticks, so a slow operation is never restarted.
459#[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/// Async capture/encode loop. Opens the camera while at least one viewer is
484/// watching and releases it when the last one leaves.
485///
486/// Cancel safety: every wait here is a real `.await` (a frame read, a demand
487/// transition, or an encode), so dropping this future (e.g. on Ctrl+C) drops
488/// `camera` and `encoder`, which release the device (LED off) and join the
489/// encode thread. Both the capture and encode threads sit idle between frames,
490/// so their joins return promptly unless the underlying device or encoder is
491/// itself wedged.
492#[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	// This track's claim on the connection. Taken on the first open, because the
502	// negotiated mode is what finally says how much this encoder can ever send, and
503	// held across reopens so the claim doesn't lapse while the camera is closed.
504	let mut reservation: Option<moq_net::bandwidth::Reservation> = None;
505
506	loop {
507		// Idle until a viewer subscribes; the track ending is a clean exit. The
508		// catalog rendition was published when the track was created, so a
509		// subscriber can get here without a frame ever having been encoded.
510		if let Err(err) = demand.used().await {
511			log_track_ended(err);
512			return Ok(());
513		}
514
515		// Open the camera and an encoder sized to its negotiated mode.
516		let Some(mut camera) = wait_capture(producer, demand, source.open(capture)).await? else {
517			continue;
518		};
519		// Capture timestamps use a private monotonic timeline. Sample both clocks
520		// once at open so every queued frame maps to the shared broadcast epoch
521		// without mistaking dequeue time for acquisition time.
522		let capture_epoch =
523			u64::try_from(clock.now().as_micros().saturating_sub(camera.now().as_micros())).unwrap_or(u64::MAX);
524		// Prefer an explicit --fps, otherwise the camera's reported rate, falling
525		// back only if the backend doesn't expose one.
526		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		// Off macOS this opens the encoder on a dedicated thread; see `sink`.
536		// No cut on reopen: a fresh encoder opens with a keyframe on every backend,
537		// so the viewer whose subscription reopened the camera can decode from the
538		// first frame regardless, and a backend that cannot cut still captures.
539		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		// A reopen can negotiate a different mode (a display resized while nothing was
545		// subscribed), and the claim follows it: the old ceiling would otherwise cap a
546		// larger mode below what it can send, or keep claiming room a smaller one no
547		// longer needs.
548		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		// Rate control is per encoder: this one opened at the configured bitrate,
553		// so the policy's ceiling is that rate and the target starts there. A
554		// reopened camera starts optimistic again rather than inheriting the
555		// backed-off rate from whatever the link was doing last time.
556		let mut rate = Some((reservation.consumer(), Control::new(Policy::new(ceiling))));
557
558		loop {
559			// Race the next frame against the last viewer leaving so we release the
560			// camera promptly when demand drops. `biased` checks demand first so an
561			// unwatched track stops before reading another frame.
562			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; // no viewers: release the camera, then wait for one
571				}
572				// Retune between frames rather than mid-encode, and only when
573				// the policy says the target actually moved.
574				estimate = next_estimate(&mut rate) => {
575					apply_estimate(&mut encoder, &mut rate, estimate).await;
576					continue;
577				}
578				// A read error is terminal for this selection (the source is gone
579				// or was refused); `None` just ends the stream, so reopen below.
580				// Timing out is a quiet camera: mark the rendition stalled and wait again.
581				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 the camera (LED off) and encoder before waiting for the next viewer.
602		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	/// Encode a handful of synthetic frames for `codec` and publish them through a real
635	/// [`Producer`], returning the catalog rendition's track name and config.
636	///
637	/// Asserts the property the whole design rests on: the rendition published before anything is
638	/// encoded is the one the first keyframe resolves. A guessed codec string would be corrected
639	/// here; a probed one is confirmed, so the catalog is written once.
640	///
641	/// `kind` is explicit so the test picks a deterministic encoder rather than `Auto`, which on
642	/// Linux CI would try the NVENC backend and panic in cudarc on a GPU-less runner.
643	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		// Jitter and delay aside, which are measured from the frames rather than declared by either.
667		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	/// The catalog's single video rendition, if it has one yet.
678	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	/// An on-demand capture resumes on the same wall clock after releasing its camera and encoder,
697	/// so the idle transition must publish a marker group between the two runs. This uses synthetic
698	/// frames and the software encoder to exercise the transition without capture hardware.
699	#[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		// The synthetic clock jumps ten seconds. Keep every fixture group readable until the
705		// assertion instead of letting the default five-second publisher window evict the marker.
706		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	/// Regression: a caller's container selection has to survive the config -> hint conversion.
765	///
766	/// [`VideoHint::container`](moq_mux::catalog::VideoHint::container) is authoritative for both the
767	/// track writer and the published rendition, so a conversion that drops it silently downgrades
768	/// the caller's selection to Legacy while the catalog still claims whatever it defaulted to.
769	#[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		// Software (openh264) so the test is deterministic and never touches a hardware backend.
777		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	/// Regression: the rendition has to reach the wire before anything is encoded.
788	///
789	/// A catalog reservation is held until the rendition resolves, and an unresolved one withholds
790	/// the whole catalog from the broadcast. An encoder that runs only while watched then closes a
791	/// cycle: the catalog waits on a keyframe, the keyframe waits on a subscriber, and the
792	/// subscriber waits on the catalog. Nothing errors on either side; the publisher simply serves
793	/// nothing, forever.
794	#[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		// Software (openh264) so the test is deterministic and never touches a hardware backend.
804		config.kind = encoder::Kind::Software;
805		let _producer = Producer::new(broadcast, catalog, config.probe().await.unwrap()).unwrap();
806
807		// Published, not merely staged: this reads the catalog track a subscriber would.
808		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		// Read out of the encoder rather than guessed: the avc3 shape (parameter sets in band) and
822		// the geometry it was opened at, which is what its first keyframe will carry.
823		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		// Neither is in the bitstream, so both come from the config that was probed.
830		assert_eq!(rendition.framerate, Some(30.0));
831		assert_eq!(rendition.bitrate, Some(6_000_000));
832	}
833
834	/// Finish leaves the handle, so abort can still run.
835	#[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		// Software (openh264) so the test is deterministic and never touches a
862		// hardware backend.
863		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	/// H.265 has no software encoder, so this only runs where a hardware one
870	/// exists (VideoToolbox on macOS, the only hardware backend on this target).
871	#[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	/// Clock fixtures: the real capture loop, fed by a synthetic camera against a pinned
881	/// broadcast clock, graded on the timestamps a subscriber reads back.
882	///
883	/// Each expectation is the acquisition instant measured on the broadcast clock. The loop
884	/// samples the broadcast clock and then the camera's timeline when it opens a camera, so a
885	/// published timestamp may land up to `SAMPLING` early, never late.
886	#[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		/// How early a mapped timestamp may land: the gap between the loop's two clock samples.
894		const SAMPLING: Duration = Duration::from_millis(250);
895		/// Rounding slack on the late side: each clock reading truncates to a microsecond.
896		const ROUNDING: u64 = 2;
897		/// Retain every fixture group, so a slow runner never evicts one before it is read.
898		const RETAIN: Duration = Duration::from_secs(600);
899
900		/// Hands the loop one fixture-supplied stream per camera open.
901		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			/// Start the capture loop on a broadcast whose clock began `behind` ago, at `wall`.
925			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				// Local: a capture stream is `!Send` on macOS, so each test runs on a `LocalSet`.
951				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			/// Subscribe to the video track, which is what opens the camera.
978			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			/// Plug in the camera the loop opens next.
993			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			/// `instant` on the broadcast clock, in microseconds.
1000			fn at(&self, instant: Instant) -> u64 {
1001				u64::try_from(instant.duration_since(self.epoch).as_micros()).unwrap()
1002			}
1003
1004			/// Stop the loop and finalize the track, as a clean end of capture does.
1005			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			/// A frame acquired at `captured` publishes at that instant on the broadcast clock.
1012			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		/// Read the next frame not already in `seen`: a resubscription replays retained groups first.
1041		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(&timestamp) {
1048					return timestamp;
1049				}
1050			}
1051		}
1052
1053		/// A camera whose first frame arrives long after the broadcast began stamps it at its
1054		/// acquisition: not zero, and not the later instant the loop dequeued it.
1055		#[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					// Delivered well after acquisition: dequeue time must not leak into the timestamp.
1065					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		/// A device clock that restarts at zero, mid-stream or across a reopen, continues the
1077		/// broadcast forward with the device's spacing instead of rewinding it.
1078		#[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					// The device numbers from zero, and real time keeps pace with it.
1087					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					// The device restarts its clock without the stream ending.
1095					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					// The device goes away and comes back numbering from zero again.
1104					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		/// Releasing the camera while nobody watches keeps the broadcast clock running: the
1119		/// frame after a resume lands after the real idle gap, at its own acquisition.
1120		#[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		/// The wall mapping is pinned when the broadcast clock is built. A system clock stepped
1155		/// an hour since then retimes neither the published timestamps nor the advertised mapping.
1156		#[tokio::test]
1157		async fn a_system_wall_adjustment_retimes_nothing() {
1158			tokio::task::LocalSet::new()
1159				.run_until(async {
1160					// Whole seconds, so the advertised mapping holds it exactly.
1161					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					// Timestamps follow the monotonic epoch and map to walls under the pinned mapping.
1174					fixture.assert_acquired(published, captured);
1175					let mapped = advertised.unwrap().wall_clock(us(published)).unwrap();
1176					// The catalog maps to walls at millisecond precision.
1177					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		/// A recording replays what the live edge published: the archive's segment records
1185		/// carry the live timestamps across an idle restart, with the idle gap left in.
1186		#[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, &section)
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						// Idle past the minimum segment, so each run is archived as its own segment.
1211						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						// The archive keeps millisecond precision.
1227						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}