diff --git a/Cargo.lock b/Cargo.lock index 55b36549aa..bbd9c3a7a8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4905,6 +4905,7 @@ dependencies = [ "futures", "h264-reader", "hang", + "kio 0.6.0", "libc", "libloading 0.9.0", "linux-raw-sys 0.12.1", diff --git a/doc/lib/js/publish.md b/doc/lib/js/publish.md index 6c9c1dd0c5..d475c613da 100644 --- a/doc/lib/js/publish.md +++ b/doc/lib/js/publish.md @@ -54,6 +54,10 @@ framerate, and bitrate are tunable through `el.video.config`; the audio encoder exposes its codec and volume. For simulcast or several renditions, drop the element and register your own encoders on a `Publish.Broadcast`. +`el.video.cut()` asks for a keyframe on top of the `keyframeInterval` cadence, +for a resume, a recording cut, or a known tune-in moment. Requests coalesce into +the next keyframe, and forced keyframes land at least 500ms apart. + The video and audio encoders measure how far their output falls behind the media clock when they flush frames. Catalog jitter is the spread above each rendition's own recent minimum lateness, so a constant encoder delay is not jitter. diff --git a/doc/lib/rs/moq-audio.md b/doc/lib/rs/moq-audio.md index e218553075..e320d7d191 100644 --- a/doc/lib/rs/moq-audio.md +++ b/doc/lib/rs/moq-audio.md @@ -37,7 +37,7 @@ than falling back to the catalog's fields. Highlights: -- **`encode::Publication`** advertises the track and opens the microphone only while someone listens. Stop, swap devices, and restart without changing the track subscribers know; read a level meter for the UI. +- **`encode::Control`** advertises the track and opens the microphone only while someone listens. Stop, swap devices, and restart without changing the track subscribers know; read a level meter for the UI. - **A/V sync signal.** `Sink::buffered()` reports how far ahead the speaker is, which is what a video clock steers by. - **Activity per packet**, read off the Opus stream, so a call UI shows who is talking without a second voice detector. - **One Linux build dependency**: ALSA headers, and only when `capture` or `playback` is enabled. diff --git a/doc/lib/rs/moq-video.md b/doc/lib/rs/moq-video.md index aa16470c41..1ef85a9139 100644 --- a/doc/lib/rs/moq-video.md +++ b/doc/lib/rs/moq-video.md @@ -21,7 +21,7 @@ ffmpeg, no GStreamer, no system codec to install. Highlights: - **Automatic backend selection**, hardware first. Linux GPU libraries are `dlopen`ed at runtime, so one binary starts anywhere and warns when it falls back to software. openh264 (the default-on `openh264` feature) is statically linked as the H.264 fallback; H.265 is hardware-only; AV1 decodes via NVDEC. The VAAPI encoder, decoder, and GPU resize share one render node: the first whose driver does all three, or the one the `MOQ_VAAPI_DEVICE` environment variable names (for example `/dev/dri/renderD129`). -- **Publish on demand.** `encode::publish_capture` advertises the track up front and opens the camera only while someone subscribes. +- **Publish on demand.** `encode::publish_capture` advertises the track up front and opens the camera only while someone subscribes. `encode::Control::new` is the same with a handle kept, the mirror of `moq-audio`'s: `Control::cut()` asks for a keyframe, requests coalesce, and forced keyframes land at least 500ms apart. - **GPU ownership where the platform allows.** Matching codec backends consume their native GPU surfaces directly. The renderer imports `CVPixelBuffer` and supported DMA-BUF formats. Linux/NVIDIA producers can import dedicated Vulkan RGBA8 slots into CUDA with timeline-semaphore ordering and completion-driven slot return. Vulkan/CUDA surfaces deliberately have no CPU pixel fallback; other surfaces use the typed `Surface::into_i420()` and configured `Surface::to_rgba(config)` when needed. - **Live bitrate control** where the selected backend supports it, without forcing a keyframe. An unsupported backend keeps its opening rate. - **Typed group structure.** `encode::Config::gop` is a `Gop` enum (`Keyframe { interval }` today), so a later mode adds a variant instead of replacing the field. `cut()` opens a group at the next frame on both `Encoder` and `Sink`, and refuses with `Error::CutUnsupported` on a backend that cannot force one rather than letting the boundary silently slip to the interval. diff --git a/js/publish/src/video/encoder.test.ts b/js/publish/src/video/encoder.test.ts index 7f1ce33f69..9ab5a3e905 100644 --- a/js/publish/src/video/encoder.test.ts +++ b/js/publish/src/video/encoder.test.ts @@ -537,3 +537,104 @@ test.each(["encoder lag", "quiet startup"])("marks a rendition stalled for %s", else Reflect.deleteProperty(globalThis, "VideoEncoder"); } }); + +test("cut forces a keyframe, coalescing requests and spacing them at least 500ms apart", async () => { + const keys: number[] = []; + class RecordingVideoEncoder { + state: CodecState = "unconfigured"; + + static async isConfigSupported(config: VideoEncoderConfig): Promise<{ supported: boolean }> { + return { supported: config.codec.startsWith("avc1") }; + } + + configure(): void { + this.state = "configured"; + } + + encode(frame: VideoFrame, options?: VideoEncoderEncodeOptions): void { + if (options?.keyFrame) keys.push(frame.timestamp / 1000); + } + + close(): void { + this.state = "closed"; + } + } + + const original = Object.getOwnPropertyDescriptor(globalThis, "VideoEncoder"); + Object.defineProperty(globalThis, "VideoEncoder", { + configurable: true, + value: RecordingVideoEncoder, + writable: true, + }); + + class Frame { + readonly timestamp: number; + constructor(timestamp: number) { + this.timestamp = timestamp; + } + clone(): Frame { + return new Frame(this.timestamp); + } + close(): void {} + } + + const { Fanout } = await import("../fanout"); + let controller!: ReadableStreamDefaultController; + const fanout = new Fanout( + new ReadableStream({ + start: (c) => { + controller = c; + }, + }), + { clone: (frame) => frame.clone(), release: (frame) => frame.close() }, + ); + + const track = new Moq.Track.Producer("video").accept(); + const rendition = { + config: new Signal(undefined), + track: new Signal(track), + close: () => track.close(), + }; + const capture = { + in: { source: new Signal({ getSettings: () => ({ frameRate: 30 }), getConstraints: () => ({}) } as never) }, + out: { display: new Signal({ width: 640, height: 480 }), frames: new Signal(fanout) }, + }; + const encoder = new Encoder("video", { + enabled: true, + broadcast: { video: () => rendition } as never, + capture: capture as never, + }); + + // One frame per 100ms of media time, after calling cut() when `cut` says so. + const send = async (millis: number, cut = false) => { + if (cut) encoder.cut(); + controller.enqueue(new Frame(millis * 1000) as unknown as VideoFrame); + await new Promise((resolve) => setTimeout(resolve, 5)); + }; + + try { + await settle(); + + // The opening keyframe serves a request made before it. + await send(0, true); + await send(100); + // Several requests before one frame produce one keyframe. + encoder.cut(); + encoder.cut(); + await send(600, true); + await send(700); + // Too soon after the last: deferred until 500ms have passed, not dropped. + await send(800, true); + await send(1000); + await send(1100); + await send(1200); + + expect(keys).toEqual([0, 600, 1100]); + } finally { + encoder.close(); + fanout.close(); + track.close(); + if (original) Object.defineProperty(globalThis, "VideoEncoder", original); + else Reflect.deleteProperty(globalThis, "VideoEncoder"); + } +}); diff --git a/js/publish/src/video/encoder.ts b/js/publish/src/video/encoder.ts index 9964f8db3d..b8d888c26d 100644 --- a/js/publish/src/video/encoder.ts +++ b/js/publish/src/video/encoder.ts @@ -147,6 +147,9 @@ export class Encoder { // doesn't re-probe the hardware. #codecFilter: Computed; + // A keyframe asked for by {@link cut} and not yet encoded. + #cut = false; + #signals = new Effect(); #stalled = new Catalog.Stalled.Detector(); #firstCaptured?: Time.Micro; @@ -326,11 +329,17 @@ export class Encoder { const interval = config?.keyframeInterval ?? Time.Milli.fromSecond(2 as Time.Second); - // Force a keyframe if this is the first frame (no group yet), or GOP elapsed. + // Force a keyframe if this is the first frame (no group yet), the GOP elapsed, or + // the caller asked for one and the last is old enough. + const since = lastKeyframe === undefined ? undefined : frame.timestamp - lastKeyframe; const keyFrame = - !lastKeyframe || lastKeyframe + Time.Micro.fromMilli(interval) <= frame.timestamp; + since === undefined || + since >= Time.Micro.fromMilli(interval) || + (this.#cut && since >= Time.Micro.fromMilli(MIN_CUT_INTERVAL)); if (keyFrame) { lastKeyframe = frame.timestamp as Time.Micro; + // Any keyframe serves an outstanding request. + this.#cut = false; } encoder.encode(frame, { keyFrame }); @@ -664,11 +673,27 @@ export class Encoder { throw new Error("no supported codec"); } + /** + * Request a keyframe, opening a new group at a frame no earlier than this call. + * + * For a resume, a recording cut, or a known tune-in moment; {@link Config.keyframeInterval} is the + * cadence. Requests coalesce into the next keyframe, and forced keyframes land at least 500ms + * apart so a caller in a loop cannot pin the encoder at all-keyframe. A request while not encoding + * is served by the keyframe every encode starts with. + */ + cut(): void { + this.#cut = true; + } + close() { this.#signals.close(); } } +// The closest two requested keyframes may land. A keyframe costs several times a predicted frame, +// and this stays well under the default two-second GOP so a request still beats the cadence. +const MIN_CUT_INTERVAL = 500 as Time.Milli; + // The source's nominal frame rate: what the capture device settled on, or what a frame stream // declared. Undefined when nothing reports one. function sourceFrameRate(source: Source): number | undefined { diff --git a/quest/m1/README.md b/quest/m1/README.md index f7b37eb577..c55193ebe2 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -59,7 +59,6 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. - [NVENC recovery](/quest/m1/nvenc-recovery.md) - partial initialization and rejected rate changes preserve valid state - [Transcode source](/quest/m1/transcode-source.md) - select a rendition the chosen backend can actually decode - [Egress rendition pick](/quest/m1/egress-rendition-pick.md) - WHEP and single-track RTMP/FLV serve the best rendition, not the first by name -- [Keyframe trigger](/quest/m1/keyframe-trigger.md) - an application can ask the built-in capture encoder for a keyframe - [QoS](/quest/m1/qos/README.md) - broadcast health: relay starvation and timeliness histograms, and client stats broadcasts from publishers and viewers - [Drain](/quest/m1/drain/README.md) - relay restarts drain sessions over GOAWAY instead of hard-dropping them - [Transport upgrade](/quest/m1/transport-upgrade/README.md) - a session that came up over WebSocket moves to QUIC once the QUIC dial lands, handing over at a group boundary diff --git a/quest/m1/keyframe-trigger.md b/quest/m1/keyframe-trigger.md deleted file mode 100644 index 67e430de51..0000000000 --- a/quest/m1/keyframe-trigger.md +++ /dev/null @@ -1,37 +0,0 @@ -# [M] On-demand keyframe trigger - -## Goal - -An application publishing through the built-in capture path can ask for a -keyframe. `Encoder::cut()`, `Sink::cut()`, and the ffi/libmoq `cut` already -force one, refusing with `CutUnsupported` when a backend cannot, but the -turnkey capture paths have no way in. - -## Plan - -`Backend::encode(frame, cut)` honors a cut (NVENC via the `FORCEIDR` picture -flag with `repeatSPSPPS` so the IDR carries its parameter sets, deliberately -not `pictureType` which `enablePTD` ignores; openh264, VAAPI, VideoToolbox and -Media Foundation the same way) and `can_cut()` answers at open whether it can. -What is missing is a caller-facing trigger on the turnkey paths: - -- `publish_capture` relies on every backend opening with a keyframe and - otherwise rides the GOP cadence; its `Options` carry no trigger. -- `js/publish`'s encode path already calls `encoder.encode(frame, { keyFrame })`, - but `lastKeyframe` is a closure-local `let` with no external trigger. - `Config.keyframeInterval` is cadence, not on demand. - -Give both a trigger the caller owns: a handle on the Rust capture path, and a -Signal the JS encode effect reads instead of its local variable. Coalesce -requests, so several arriving within one frame interval produce one IDR rather -than a run of them, and rate limit at the publisher so a caller in a loop -cannot pin the encoder at all-IDR. - -Additive on both sides, so it lands on `main`. Real callers exist regardless -of whether a wire-level request ever ships: a resume, a recording cut, a -rendition switch, and an application that knows its own tune-in moment. - -## Related - -- [GOP overhead](/quest/m2/gop-overhead.md) - whether a long GOP driven by a - keyframe request is worth designing at all diff --git a/quest/m1/qos/stats/encoder-feedback.md b/quest/m1/qos/stats/encoder-feedback.md index 11518d8cc7..55d5cbee83 100644 --- a/quest/m1/qos/stats/encoder-feedback.md +++ b/quest/m1/qos/stats/encoder-feedback.md @@ -46,5 +46,3 @@ prefix. Keyframe requests stay out. - [Ladder](/quest/m1/ladder/README.md) - the transcode ladder that adapts to its uplink today -- [Keyframe trigger](/quest/m1/keyframe-trigger.md) - the keyframe request - this loop does not send diff --git a/quest/m2/capture-clock-source.md b/quest/m2/capture-clock-source.md index ddd4d10888..a4979cf4fa 100644 --- a/quest/m2/capture-clock-source.md +++ b/quest/m2/capture-clock-source.md @@ -2,9 +2,9 @@ ## Goal -`moq_video::encode::publish_capture` and `moq_audio::encode::Publication` stamp +`moq_video::encode::publish_capture` and `moq_audio::encode::Control` stamp on the clock their catalog advertises, with no separate clock to pass. Today -each takes its own `moq_mux::Clock`, and `PublicationOptions::default()` builds +each takes its own `moq_mux::Clock`, and `CaptureOptions::default()` builds a fresh one, so a caller relying on the default publishes audio against a mapping the catalog never advertised. `moq import capture` passes `catalog.clock()` to both, which is the only correct value. @@ -12,7 +12,7 @@ mapping the catalog never advertised. `moq import capture` passes ## Plan Drop the `clock` parameter from video `publish_capture` and the `clock` field -from `PublicationOptions`, reading `catalog.clock()` instead. Update moq-cli and +from `CaptureOptions`, reading `catalog.clock()` instead. Update moq-cli and any binding that forwards a clock. The clock fixtures in both crates already pass the catalog's clock, so they keep grading the same path. diff --git a/quest/m2/gop-overhead.md b/quest/m2/gop-overhead.md index 67383329de..9c9a8649b3 100644 --- a/quest/m2/gop-overhead.md +++ b/quest/m2/gop-overhead.md @@ -33,11 +33,6 @@ a verdict on whether a long GOP plus a keyframe request is worth designing. A verdict of "2 seconds is fine" is a valid, expected outcome and completes this quest. -## Related - -- [Keyframe trigger](/quest/m1/keyframe-trigger.md) - the publisher-side half a - keyframe request would drive, useful on its own - ## Closes - [#2284](https://github.com/moq-dev/moq/issues/2284) - close this issue when the quest finishes diff --git a/rs/moq-audio/src/encode/capture.rs b/rs/moq-audio/src/encode/capture.rs index 33fc3f428f..0a067dabe8 100644 --- a/rs/moq-audio/src/encode/capture.rs +++ b/rs/moq-audio/src/encode/capture.rs @@ -40,7 +40,7 @@ pub enum Status { /// The post-processing level of the most recently captured buffer. /// -/// Zero while the input is closed. Read it with [`Publication::level`] at +/// Zero while the input is closed. Read it with [`Control::level`] at /// whatever rate the meter draws at. #[derive(Clone, Copy, Debug, Default, PartialEq)] pub struct Level { @@ -82,8 +82,8 @@ impl Level { /// A snapshot of a capture-backed publication's lifecycle. /// /// Levels are deliberately not here: they change every buffer, so they would -/// drown out the transitions [`Publication::changed`] exists to report. Read -/// them with [`Publication::level`] instead. +/// drown out the transitions [`Control::changed`] exists to report. Read +/// them with [`Control::level`] instead. #[derive(Clone, Debug)] pub struct State { status: Status, @@ -116,14 +116,14 @@ impl State { } } -/// Capture and encode settings for [`Publication`]. +/// Capture and encode settings for [`Control::new`] and [`publish_capture`]. /// -/// `#[non_exhaustive]`: construct via [`PublicationOptions::default`] and set +/// `#[non_exhaustive]`: construct via [`CaptureOptions::default`] and set /// fields, so new publication settings can be added without changing -/// [`Publication::new`]. +/// [`Control::new`]. #[derive(Clone, Debug, Default)] #[non_exhaustive] -pub struct PublicationOptions { +pub struct CaptureOptions { /// The initial input and its capture processing. pub capture: capture::Config, /// The track's stable codec and encode settings. @@ -145,13 +145,13 @@ struct PublishedState { revision: u64, } -/// A retained handle for a capture-backed audio publication. +/// A handle for controlling a capture-backed audio publication. /// /// Clones control the same MoQ track. Stopping or replacing the input releases /// the device but leaves the track and catalog rendition intact, so restarting /// does not change the broadcast identity. Dropping the final clone stops the /// driver and releases the publication. -pub struct Publication { +pub struct Control { desired: kio::Producer, state: kio::Consumer, level: kio::Consumer, @@ -159,7 +159,7 @@ pub struct Publication { track_name: Arc, } -impl Clone for Publication { +impl Clone for Control { fn clone(&self) -> Self { Self { desired: self.desired.clone(), @@ -171,23 +171,23 @@ impl Clone for Publication { } } -impl fmt::Debug for Publication { +impl fmt::Debug for Control { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - f.debug_struct("Publication") + f.debug_struct("Control") .field("track_name", &self.track_name) .field("state", &self.state.read().state) .finish_non_exhaustive() } } -impl Publication { +impl Control { /// Register one stable audio track and return its control handle and driver. /// /// The initial source is enabled, but opens only while the track has a /// subscriber. On macOS the driver's future is `!Send`, because the permission /// prompt and ScreenCaptureKit hold ObjC handles across an await, so await /// [`Driver::run`] on a local task there; elsewhere it can be spawned. - /// [`Publication`] itself is `Send + Sync`, so the controls can live anywhere. + /// [`Control`] itself is `Send + Sync`, so the controls can live anywhere. /// /// The track is registered here, but its catalog rendition describes the /// source's PCM layout, so the driver probes for that first and registers the @@ -197,7 +197,7 @@ impl Publication { pub fn new( broadcast: moq_net::broadcast::Producer, catalog: moq_mux::catalog::Producer, - options: PublicationOptions, + options: CaptureOptions, ) -> Result<(Self, Driver), Error> { Self::build(broadcast, catalog, options, Supervisor::default()) } @@ -205,7 +205,7 @@ impl Publication { fn build( mut broadcast: moq_net::broadcast::Producer, catalog: moq_mux::catalog::Producer, - options: PublicationOptions, + options: CaptureOptions, supervisor: Supervisor, ) -> Result<(Self, Driver), Error> { let reserved = Reserved::new(&mut broadcast, catalog, &options.encode)?; @@ -227,7 +227,7 @@ impl Publication { let desired_tx = kio::Producer::new(desired); let state_tx = kio::Producer::new(initial); let level_tx = kio::Producer::new(Level::default()); - let publication = Self { + let control = Self { desired: desired_tx.clone(), state: state_tx.consume(), level: level_tx.consume(), @@ -246,7 +246,7 @@ impl Publication { level: level_tx, park_on_failure: true, }; - Ok((publication, driver)) + Ok((control, driver)) } /// Enable capture using the selected source. @@ -332,7 +332,7 @@ impl Publication { /// The task that opens the selected input and publishes its samples. /// /// The driver owns the broadcast producer so its identity remains alive through -/// stop, failure, replacement, and restart. Dropping the final [`Publication`] +/// stop, failure, replacement, and restart. Dropping the final [`Control`] /// ends the driver and releases that identity. pub struct Driver { _broadcast: moq_net::broadcast::Producer, @@ -705,15 +705,16 @@ fn publish_state( /// Capture audio on demand and publish it as an encoded MoQ track. /// -/// This convenience function runs a controllable [`Publication`] from -/// `options`. Use [`Publication::new`] directly to retain controls. +/// This convenience function runs a [`Control`] from +/// `options` without handing back the handle. Use [`Control::new`] directly to +/// retain controls. pub async fn publish_capture( broadcast: moq_net::broadcast::Producer, catalog: moq_mux::catalog::Producer, - options: PublicationOptions, + options: CaptureOptions, ) -> Result<(), Error> { // Held, not dropped: the driver ends as soon as the last control handle goes. - let (_publication, driver) = Publication::new(broadcast, catalog, options)?; + let (_control, driver) = Control::new(broadcast, catalog, options)?; Driver { park_on_failure: false, ..driver @@ -731,7 +732,7 @@ pub async fn publish_capture( fn assert_publish_capture_send( broadcast: moq_net::broadcast::Producer, catalog: moq_mux::catalog::Producer, - options: PublicationOptions, + options: CaptureOptions, ) { fn is_send(_: &T) {} is_send(&publish_capture(broadcast, catalog, options)); @@ -876,7 +877,7 @@ impl Output for EncoderOutput<'_, E> { impl EncoderOutput<'_, E> { /// Levels ride their own channel: they change every buffer, so folding them - /// into [`State`] would wake every [`Publication::changed`] waiter at the + /// into [`State`] would wake every [`Control::changed`] waiter at the /// capture rate and rebuild the whole snapshot to do it. fn publish_level(&self, level: Level) { let Ok(mut published) = self.level.write() else { return }; @@ -1435,7 +1436,7 @@ mod tests { async fn setup_publication( opens: impl IntoIterator, ) -> ( - Publication, + Control, Driver, MockSource, moq_net::track::Subscriber, @@ -1444,13 +1445,12 @@ mod tests { let mut broadcast = moq_net::broadcast::Info::new().produce(); let consumer = broadcast.consume(); let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap(); - let mut options = PublicationOptions::default(); + let mut options = CaptureOptions::default(); options.capture.source = capture::Source::Microphone(Some("first".into())); options.encode.track = Some("audio".into()); - let (publication, driver) = - Publication::build(broadcast, catalog.clone(), options, Supervisor::exact()).unwrap(); + let (control, driver) = Control::build(broadcast, catalog.clone(), options, Supervisor::exact()).unwrap(); let subscription = consumer.track("audio").unwrap().subscribe(None).await.unwrap(); - (publication, driver, source(opens, false), subscription, catalog) + (control, driver, source(opens, false), subscription, catalog) } /// Same, with the caller's encode options, which `setup_publication` fixes. @@ -1458,7 +1458,7 @@ mod tests { channels: u32, opens: impl IntoIterator, ) -> ( - Publication, + Control, Driver, MockSource, moq_net::track::Subscriber, @@ -1467,24 +1467,23 @@ mod tests { let mut broadcast = moq_net::broadcast::Info::new().produce(); let consumer = broadcast.consume(); let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap(); - let mut options = PublicationOptions::default(); + let mut options = CaptureOptions::default(); options.capture.source = capture::Source::Microphone(Some("first".into())); options.encode.track = Some("audio".into()); options.encode.settings.layout = PcmLayout::from_channels(channels).unwrap(); - let (publication, driver) = - Publication::build(broadcast, catalog.clone(), options, Supervisor::exact()).unwrap(); + let (control, driver) = Control::build(broadcast, catalog.clone(), options, Supervisor::exact()).unwrap(); let subscription = consumer.track("audio").unwrap().subscribe(None).await.unwrap(); - (publication, driver, source(opens, false), subscription, catalog) + (control, driver, source(opens, false), subscription, catalog) } - async fn wait_for(publication: &mut Publication, status: Status) -> State { + async fn wait_for(control: &mut Control, status: Status) -> State { tokio::time::timeout(Duration::from_secs(1), async { loop { - let state = publication.state(); + let state = control.state(); if state.status() == status { return state; } - publication.changed().await.expect("driver still running"); + control.changed().await.expect("driver still running"); } }) .await @@ -1495,25 +1494,25 @@ mod tests { async fn failed_publication_retries_but_duplicate_start_is_idempotent() { let (_events, recovered) = stream(None); let recovered = with_device(recovered, "allowed"); - let (mut publication, driver, source, _subscription, _catalog) = + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Fatal("permission denied"), Open::Stream(recovered)]).await; let attempts = source.attempts.clone(); let task = tokio::spawn(driver.run_with(source)); - let failed = wait_for(&mut publication, Status::Failed).await; + let failed = wait_for(&mut control, Status::Failed).await; assert!(matches!(failed.failure(), Some(Error::Capture(message)) if message == "permission denied")); assert_eq!(attempts.load(Ordering::SeqCst), 1); - publication.start(); - let live = wait_for(&mut publication, Status::Live).await; + control.start(); + let live = wait_for(&mut control, Status::Live).await; assert_eq!(live.device().map(|device| device.id.as_str()), Some("allowed")); assert_eq!(attempts.load(Ordering::SeqCst), 2); - publication.start(); + control.start(); tokio::task::yield_now().await; assert_eq!(attempts.load(Ordering::SeqCst), 2); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1524,40 +1523,40 @@ mod tests { let first = with_device(first, "first"); let (_second_events, second) = stream(None); let second = with_device(second, "second"); - let (mut publication, driver, source, _subscription, _catalog) = + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(first), Open::Stream(second)]).await; let task = tokio::spawn(driver.run_with(source)); - wait_for(&mut publication, Status::Live).await; - publication.stop(); - wait_for(&mut publication, Status::Stopped).await; + wait_for(&mut control, Status::Live).await; + control.stop(); + wait_for(&mut control, Status::Stopped).await; assert_eq!(first_drops.load(Ordering::SeqCst), 1); - let track = publication.track_name().to_string(); + let track = control.track_name().to_string(); - publication.replace(capture::Source::Microphone(Some("second".into()))); - publication.start(); - let live = wait_for(&mut publication, Status::Live).await; + control.replace(capture::Source::Microphone(Some("second".into()))); + control.start(); + let live = wait_for(&mut control, Status::Live).await; assert_eq!(live.device().map(|device| device.id.as_str()), Some("second")); - assert_eq!(publication.track_name(), track); + assert_eq!(control.track_name(), track); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } #[tokio::test] async fn reports_post_processing_level() { let (events, input) = stream(None); - let (mut publication, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; let task = tokio::spawn(driver.run_with(source)); - wait_for(&mut publication, Status::Live).await; - assert_eq!(publication.level(), Level::default()); + wait_for(&mut control, Status::Live).await; + assert_eq!(control.level(), Level::default()); events .try_push(Ok(capture::Samples::plain(vec![0.25, -0.5, 1.0, -1.0], false))) .unwrap(); let level = tokio::time::timeout(Duration::from_secs(1), async { loop { - let level = publication.level(); + let level = control.level(); if level != Level::default() { return level; } @@ -1569,7 +1568,7 @@ mod tests { assert!((level.rms() - 0.760_345_34).abs() < 0.000_001); assert_eq!(level.peak(), 1.0); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1577,20 +1576,20 @@ mod tests { #[tokio::test] async fn level_updates_do_not_wake_state_waiters() { let (events, input) = stream(None); - let (mut publication, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; let task = tokio::spawn(driver.run_with(source)); - wait_for(&mut publication, Status::Live).await; + wait_for(&mut control, Status::Live).await; for _ in 0..4 { events .try_push(Ok(capture::Samples::plain(vec![0.25, -0.5, 1.0, -1.0], false))) .unwrap(); } - tokio::time::timeout(Duration::from_millis(100), publication.changed()) + tokio::time::timeout(Duration::from_millis(100), control.changed()) .await .expect_err("a level change woke a state waiter"); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1598,7 +1597,7 @@ mod tests { /// failure would hang `moq import capture` instead of reporting the denial. #[tokio::test] async fn a_publication_without_controls_returns_its_terminal_failure() { - let (mut publication, driver, source, _subscription, _catalog) = + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Fatal("permission denied")]).await; let driver = Driver { park_on_failure: false, @@ -1611,22 +1610,22 @@ mod tests { .expect_err("the terminal failure was swallowed"); assert!(matches!(&err, Error::Capture(message) if message == "permission denied")); - assert!(publication.is_finished()); - assert_eq!(publication.changed().await.unwrap().status(), Status::Failed); - assert!(publication.changed().await.is_none()); + assert!(control.is_finished()); + assert_eq!(control.changed().await.unwrap().status(), Status::Failed); + assert!(control.changed().await.is_none()); } /// A retained publication keeps the track and waits to be told what to do. #[tokio::test] async fn a_terminal_failure_parks_a_retained_publication() { - let (mut publication, driver, source, _subscription, _catalog) = + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Fatal("permission denied")]).await; let task = tokio::spawn(driver.run_with(source)); - wait_for(&mut publication, Status::Failed).await; - assert!(!publication.is_finished()); + wait_for(&mut control, Status::Failed).await; + assert!(!control.is_finished()); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1634,26 +1633,26 @@ mod tests { #[tokio::test] async fn stopping_zeroes_the_level() { let (events, input) = stream(None); - let (mut publication, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; let task = tokio::spawn(driver.run_with(source)); - wait_for(&mut publication, Status::Live).await; + wait_for(&mut control, Status::Live).await; events .try_push(Ok(capture::Samples::plain(vec![0.5, -0.5], false))) .unwrap(); tokio::time::timeout(Duration::from_secs(1), async { - while publication.level() == Level::default() { + while control.level() == Level::default() { tokio::task::yield_now().await; } }) .await .expect("no level was measured"); - publication.stop(); - wait_for(&mut publication, Status::Stopped).await; - assert_eq!(publication.level(), Level::default()); + control.stop(); + wait_for(&mut control, Status::Stopped).await; + assert_eq!(control.level(), Level::default()); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1663,20 +1662,20 @@ mod tests { async fn a_terminal_live_failure_keeps_the_device_that_failed() { let (events, input) = stream(None); let input = with_device(input, "wired"); - let (mut publication, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; + let (mut control, driver, source, _subscription, _catalog) = setup_publication([Open::Stream(input)]).await; let task = tokio::spawn(driver.run_with(source)); - wait_for(&mut publication, Status::Live).await; + wait_for(&mut control, Status::Live).await; events .try_push(Err(capture::Failure::fatal(Error::Capture("device vanished".into())))) .unwrap(); - let failed = wait_for(&mut publication, Status::Failed).await; + let failed = wait_for(&mut control, Status::Failed).await; assert_eq!(failed.device().map(|device| device.id.as_str()), Some("wired")); assert!(matches!(failed.failure(), Some(Error::Capture(message)) if message == "device vanished")); // Retained controls park rather than end, so dropping them is the clean exit. - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1684,13 +1683,13 @@ mod tests { /// construction must not wait on one and the controls must work meanwhile. #[tokio::test] async fn controls_are_live_before_the_input_is_discovered() { - let (mut publication, driver, mut source, _subscription, catalog) = setup_publication([]).await; + let (mut control, driver, mut source, _subscription, catalog) = setup_publication([]).await; // Nothing ever answers a probe, which is a machine with no input device. source.formats.clear(); let attempts = source.format_attempts.clone(); let task = tokio::spawn(driver.run_with(source)); - assert_eq!(publication.track_name(), "audio"); + assert_eq!(control.track_name(), "audio"); tokio::time::timeout(Duration::from_secs(1), async { while attempts.load(Ordering::SeqCst) == 0 { tokio::task::yield_now().await; @@ -1699,16 +1698,16 @@ mod tests { .await .expect("the driver never probed the input"); - assert_eq!(publication.state().status(), Status::Starting); + assert_eq!(control.state().status(), Status::Starting); // The layout is unknown, so there is no rendition to advertise yet. assert!(catalog.snapshot().audio.renditions.is_empty()); // A probe nothing answers is still cancellable by the controls. - publication.stop(); - wait_for(&mut publication, Status::Stopped).await; + control.stop(); + wait_for(&mut control, Status::Stopped).await; assert_eq!(attempts.load(Ordering::SeqCst), 1); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1717,27 +1716,26 @@ mod tests { #[tokio::test] async fn discovery_registers_the_rendition() { let (_events, input) = stream(None); - let (mut publication, driver, mut source, _subscription, catalog) = - setup_publication([Open::Stream(input)]).await; + let (mut control, driver, mut source, _subscription, catalog) = setup_publication([Open::Stream(input)]).await; source.formats = [Discovery::Fatal("permission denied"), Discovery::Format(48_000, 2)] .into_iter() .collect(); let task = tokio::spawn(driver.run_with(source)); - let failed = wait_for(&mut publication, Status::Failed).await; + let failed = wait_for(&mut control, Status::Failed).await; assert!(matches!(failed.failure(), Some(Error::Capture(message)) if message == "permission denied")); assert!(catalog.snapshot().audio.renditions.is_empty()); // A terminal probe failure parks like a terminal open failure, so `start` // retries it rather than the driver ending with no track. - publication.start(); - wait_for(&mut publication, Status::Live).await; + control.start(); + wait_for(&mut control, Status::Live).await; let renditions = catalog.snapshot().audio.renditions; let config = renditions.get("audio").expect("the rendition was never registered"); assert_eq!(config.sample_rate, 48_000); assert_eq!(config.channel_count, 2); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1748,23 +1746,22 @@ mod tests { async fn a_rejected_layout_parks_the_publication() { let (_events, input) = stream(None); // A discrete source cannot be assigned stereo speaker positions implicitly. - let (mut publication, driver, mut source, _subscription, catalog) = - setup_encoding(2, [Open::Stream(input)]).await; + let (mut control, driver, mut source, _subscription, catalog) = setup_encoding(2, [Open::Stream(input)]).await; source.formats = [Discovery::Format(48_000, 6), Discovery::Format(48_000, 2)] .into_iter() .collect(); let task = tokio::spawn(driver.run_with(source)); - let failed = wait_for(&mut publication, Status::Failed).await; + let failed = wait_for(&mut control, Status::Failed).await; assert!(failed.failure().is_some()); - assert!(!publication.is_finished()); + assert!(!control.is_finished()); assert!(catalog.snapshot().audio.renditions.is_empty()); - publication.replace(capture::Source::Microphone(Some("second".into()))); - wait_for(&mut publication, Status::Live).await; + control.replace(capture::Source::Microphone(Some("second".into()))); + wait_for(&mut control, Status::Live).await; assert_eq!(catalog.snapshot().audio.renditions.len(), 1); - drop(publication); + drop(control); task.await.unwrap().unwrap(); } @@ -1772,7 +1769,7 @@ mod tests { #[test] fn publication_controls_cross_threads() { fn assert_send_sync() {} - assert_send_sync::(); + assert_send_sync::(); assert_send_sync::(); assert_send_sync::(); } @@ -2064,7 +2061,7 @@ mod tests { clock: moq_mux::Clock, catalog: moq_mux::catalog::Producer, consumer: moq_net::broadcast::Consumer, - publication: Publication, + publication: Control, task: tokio::task::JoinHandle>, } @@ -2081,12 +2078,12 @@ mod tests { .with_clock(clock) .with_max_age(RETAIN); let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config).unwrap(); - let mut options = PublicationOptions::default(); + let mut options = CaptureOptions::default(); options.encode.track = Some("audio".into()); // The broadcast's own clock, as `moq import capture` hands it. options.clock = catalog.clock(); let (publication, driver) = - Publication::build(broadcast, catalog.clone(), options, Supervisor::exact()).unwrap(); + Control::build(broadcast, catalog.clone(), options, Supervisor::exact()).unwrap(); let task = tokio::spawn(driver.run_with(source(opens, false))); Self { epoch, diff --git a/rs/moq-audio/src/encode/mod.rs b/rs/moq-audio/src/encode/mod.rs index 9c918b9744..7756686b35 100644 --- a/rs/moq-audio/src/encode/mod.rs +++ b/rs/moq-audio/src/encode/mod.rs @@ -27,4 +27,4 @@ pub use encoder::{Codec, Encoder, Finish, Input, Settings}; pub use producer::{Options, Producer}; #[cfg(feature = "capture")] -pub use capture::{Driver, Level, Publication, PublicationOptions, State, Status, publish_capture}; +pub use capture::{CaptureOptions, Control, Driver, Level, State, Status, publish_capture}; diff --git a/rs/moq-audio/src/lib.rs b/rs/moq-audio/src/lib.rs index b8c7aa9242..f6c3572268 100644 --- a/rs/moq-audio/src/lib.rs +++ b/rs/moq-audio/src/lib.rs @@ -14,10 +14,10 @@ //! they don't exist in a default build. //! - [`encode`] encodes PCM and publishes it through `moq_mux::container`, //! registering the rendition in the `hang` catalog. Two entry points: -//! - `encode::Publication` and `encode::Driver` capture a controllable -//! microphone publication. The retained publication starts, stops, and -//! replaces the input while preserving one track identity, and reports the -//! active device, failures, and post-processing level. +//! - `encode::Control::new` returns a `Control` handle and the `Driver` +//! that captures the microphone. The handle starts, stops, and replaces +//! the input while preserving one track identity, and reports the active +//! device, failures, and post-processing level. //! - `encode::publish_capture` is the turnkey shorthand. It encodes strictly //! on demand: the track and catalog are advertised up front, but the device //! opens only while a subscriber is listening and is released when the last diff --git a/rs/moq-cli/src/publish.rs b/rs/moq-cli/src/publish.rs index 306480c9ae..95722faa89 100644 --- a/rs/moq-cli/src/publish.rs +++ b/rs/moq-cli/src/publish.rs @@ -402,7 +402,11 @@ impl Publish { async move { match video { Some((config, encode)) => { - moq_video::encode::publish_capture(broadcast, catalog, config, encode, clock) + let mut options = moq_video::encode::CaptureOptions::default(); + options.capture = config; + options.encode = encode; + options.clock = clock; + moq_video::encode::publish_capture(broadcast, catalog, options) .await .map_err(anyhow::Error::from) } @@ -415,7 +419,7 @@ impl Publish { async move { match audio { Some((config, encode)) => { - let mut options = moq_audio::encode::PublicationOptions::default(); + let mut options = moq_audio::encode::CaptureOptions::default(); options.capture = config; options.encode = encode; options.clock = clock; diff --git a/rs/moq-video/Cargo.toml b/rs/moq-video/Cargo.toml index 968a7c1fcb..30ff9ebd6d 100644 --- a/rs/moq-video/Cargo.toml +++ b/rs/moq-video/Cargo.toml @@ -31,7 +31,7 @@ default = ["mediacodec", "nvidia", "openh264"] # Linux goes through moq-v4l (bindings checked in), macOS and Windows use OS # frameworks. Still off by default until every distribution ships it, which is # a packaging decision rather than a build cost. moq-cli's `capture` turns it on. -capture = ["dep:libc", "dep:moq-v4l", "dep:x11rb", "dep:zune-jpeg"] +capture = ["dep:kio", "dep:libc", "dep:moq-v4l", "dep:x11rb", "dep:zune-jpeg"] # Android MediaCodec encode/decode and its AHardwareBuffer frame surface. Kept # behind a feature because these NDK entry points require API 26, while the # language-binding artifacts still support API 24 and disable default features. @@ -106,6 +106,8 @@ bytes = { workspace = true } fast_image_resize = "6" # Catalog types (VideoConfig / VideoCodec) the decode consumer reads. hang = { workspace = true } +# The capture publish's control handle. +kio = { workspace = true, optional = true } moq-mux = { workspace = true } moq-net = { workspace = true } # Vendored, statically linked software H.264 fallback (no system dependency). diff --git a/rs/moq-video/README.md b/rs/moq-video/README.md index ea4edda4d6..f3b123a522 100644 --- a/rs/moq-video/README.md +++ b/rs/moq-video/README.md @@ -91,11 +91,14 @@ driver without the force-keyframe control), and queues nothing then: groups keep falling where `Config::gop` puts them. `encode::Sink` answers the same way, awaited. -Two public entry points: +Public entry points: - `encode::publish_capture(...)` captures a webcam, encodes it, and publishes on demand: the track and catalog are advertised up front, but the camera opens only while a subscriber is watching and is released when the last one leaves. +- `encode::Control::new(...)` does the same but returns a `Control` handle with + the `Driver` that runs it, like `moq-audio`'s. `Control::cut()` asks for a + keyframe: requests coalesce, and forced keyframes land at least 500ms apart. - `encode::Producer` publishes frames you encoded yourself (`publish(&[Encoded])`), handling the catalog and framing. Each is published at its own timestamp. diff --git a/rs/moq-video/src/encode/cuts.rs b/rs/moq-video/src/encode/cuts.rs new file mode 100644 index 0000000000..9eeb2cd7e1 --- /dev/null +++ b/rs/moq-video/src/encode/cuts.rs @@ -0,0 +1,110 @@ +//! [`Cuts`]: which captured frames to force a keyframe on. + +use std::time::Duration; + +use moq_net::Timestamp; + +/// The closest two forced keyframes may land, in media time. +/// +/// A keyframe costs several times a predicted frame, so a caller asking in a loop +/// would otherwise pin the encoder at all-IDR and starve the rest of the uplink. +/// Well under the default two-second GOP, so a request still beats the cadence. +const MIN_INTERVAL: Duration = Duration::from_millis(500); + +/// One encoder's view of the keyframe requests: decides which frames to cut, +/// coalescing and rate limiting. Built fresh for every encoder the capture opens. +/// +/// Requests arrive as a running count, so any number between two frames read as one. +pub(super) struct Cuts { + /// The request count already accounted for. + seen: u64, + /// A request not yet honored because the last keyframe was too recent. + pending: bool, + /// The last keyframe this encoder was asked for, or `None` before its first frame. + last: Option, +} + +impl Cuts { + /// Start from `requests`, the count already served by earlier encoders. + pub fn new(requests: u64) -> Self { + Self { + seen: requests, + pending: false, + last: None, + } + } + + /// Whether the frame at `timestamp` should be cut, given the running request + /// count, recording it if so. + pub fn due(&mut self, requests: u64, timestamp: Timestamp) -> bool { + self.pending |= requests != self.seen; + self.seen = requests; + + // A fresh encoder opens with a keyframe on every backend, which serves anything + // requested before it. + let Some(last) = self.last else { + self.last = Some(timestamp); + self.pending = false; + return false; + }; + + if !self.pending || Duration::from(timestamp).saturating_sub(Duration::from(last)) < MIN_INTERVAL { + return false; + } + + self.pending = false; + self.last = Some(timestamp); + true + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn at(millis: u64) -> Timestamp { + Timestamp::from_millis(millis).unwrap() + } + + #[test] + fn the_opening_keyframe_serves_earlier_requests() { + let mut cuts = Cuts::new(0); + assert!(!cuts.due(1, at(0))); + assert!(!cuts.due(1, at(1_000)), "the request was already served"); + } + + #[test] + fn requests_before_a_frame_coalesce() { + let mut cuts = Cuts::new(0); + assert!(!cuts.due(0, at(0))); + assert!(cuts.due(10, at(1_000))); + assert!(!cuts.due(10, at(2_000)), "ten requests before one frame cut once"); + } + + #[test] + fn a_request_too_soon_is_deferred_not_dropped() { + let mut cuts = Cuts::new(0); + assert!(!cuts.due(0, at(0))); + assert!(!cuts.due(1, at(100))); + assert!(!cuts.due(1, at(499))); + assert!(cuts.due(1, at(500)), "held until the interval elapsed"); + } + + #[test] + fn a_caller_in_a_loop_cannot_force_all_idr() { + let mut cuts = Cuts::new(0); + + // Ten seconds at 30 fps with a request before every frame. + let cut = (0..300u64) + .filter(|frame| cuts.due(frame + 1, at(frame * 1000 / 30))) + .count(); + assert_eq!(cut, 19, "one forced keyframe per half second after the opening one"); + } + + #[test] + fn a_reopened_encoder_ignores_requests_already_counted() { + let mut cuts = Cuts::new(5); + assert!(!cuts.due(5, at(0))); + assert!(!cuts.due(5, at(1_000))); + } +} diff --git a/rs/moq-video/src/encode/mod.rs b/rs/moq-video/src/encode/mod.rs index 585d9fe68e..005b122a52 100644 --- a/rs/moq-video/src/encode/mod.rs +++ b/rs/moq-video/src/encode/mod.rs @@ -5,7 +5,9 @@ //! //! Entry points, high to low level: //! - `publish_capture` captures and publishes a webcam (turnkey). Requires the -//! `capture` feature. +//! `capture` feature. `Control::new` is the same with a handle kept for +//! controlling it while it runs (asking for a keyframe), plus the `Driver` +//! that runs it. //! - [`Encoder`] encodes raw [`Frame`](crate::Frame)s you supply into //! [`Encoded`] access units, and [`Producer`] publishes those (bring your own //! frames). Build both for the same [`Codec`]. @@ -19,7 +21,7 @@ //! discover a track nothing has encoded yet. That's what makes on-demand //! encoding possible at all. //! -//! `Options` (with `capture`) / [`Kind`] / [`Config`] configure them. The decode/consume +//! `CaptureOptions` / `Options` (with `capture`) / [`Kind`] / [`Config`] configure them. The decode/consume //! counterpart (mirror of `moq-audio`'s consumer) lives in the sibling //! [`decode`](crate::decode) module. //! @@ -31,13 +33,16 @@ mod encoded; mod encoder; mod producer; mod sink; +// Compiled without `capture` so its tests stay in the default merge gate. +#[cfg_attr(not(feature = "capture"), allow(dead_code))] +mod cuts; pub use backend::NAMES; pub use encoded::Encoded; pub use encoder::{Codec, Config, Encoder, Gop, Kind}; pub use producer::Producer; #[cfg(feature = "capture")] -pub use producer::{Options, publish_capture}; +pub use producer::{CaptureOptions, Control, Driver, Options, publish_capture}; pub use sink::Sink; #[cfg(test)] diff --git a/rs/moq-video/src/encode/producer.rs b/rs/moq-video/src/encode/producer.rs index 86353d907b..e81c365465 100644 --- a/rs/moq-video/src/encode/producer.rs +++ b/rs/moq-video/src/encode/producer.rs @@ -13,7 +13,7 @@ use std::time::Instant; use moq_mux::catalog::hang::CatalogExt; #[cfg(feature = "capture")] -use moq_mux::rate::{Control, Policy}; +use moq_mux::rate; #[cfg(any(feature = "capture", test))] use moq_net::Timestamp; @@ -26,6 +26,8 @@ use crate::capture; use super::Encoded; #[cfg(feature = "capture")] use super::Sink; +#[cfg(feature = "capture")] +use super::cuts::Cuts; #[cfg(any(feature = "capture", test))] use super::encoder; #[cfg(feature = "capture")] @@ -241,7 +243,7 @@ impl Producer { } } -/// Source-agnostic encode knobs for [`publish_capture`], where the geometry +/// Source-agnostic encode knobs for a capture publish, where the geometry /// (width / height / framerate) comes from the capture source, not the caller. /// For the bring-your-own-frames [`Encoder`](super::Encoder) path, where you /// must specify geometry, use [`Config`](super::Config) instead. @@ -279,65 +281,193 @@ pub struct Options { pub bandwidth: moq_net::bandwidth::Allocator, } -/// Capture a webcam and publish it as an on-demand video track. +/// Capture and encode settings for [`Control::new`] and [`publish_capture`]. /// -/// Returns when the broadcast is dropped (the track stops being announced) -/// or the capture loop fails. Frames are stamped from `clock`, so passing the -/// same [`Clock`](moq_mux::Clock) to a concurrent audio publish keeps the two -/// tracks aligned. +/// `#[non_exhaustive]`: construct via [`CaptureOptions::default`] and set +/// fields, so new settings can be added without changing [`Control::new`]. +#[derive(Clone, Debug, Default)] +#[non_exhaustive] +#[cfg(feature = "capture")] +pub struct CaptureOptions { + /// The source to capture. + pub capture: capture::Config, + /// The track's codec and encode settings. + pub encode: Options, + /// The shared clock used to align this track with concurrent media. + pub clock: moq_mux::Clock, +} + +/// A handle for controlling a running capture publish. /// -/// The camera is opened once at startup to probe the mode it negotiates, then released until a -/// subscriber arrives and reopened for as long as one is watching. That one open is what lets the -/// catalog rendition be exact before a single frame is published, so a consumer can size itself -/// against it (and discover the track at all) without waiting for an encoder that may never run. +/// Clones control the same track. Dropping the final clone stops the [`Driver`] +/// and ends the track. +#[derive(Clone, Debug)] #[cfg(feature = "capture")] -pub async fn publish_capture( - broadcast: moq_net::broadcast::Producer, +pub struct Control { + /// A running count of keyframe requests, so the driver coalesces any number + /// between two frames into one. + cuts: kio::Producer, +} + +#[cfg(feature = "capture")] +impl Control { + /// Register one video track and return its control handle and driver. + /// + /// The track is registered here, but its catalog rendition describes the mode + /// the source negotiates, so the driver opens the source once to probe it + /// before publishing the rendition. Off macOS the driver's future is `Send` + /// and can be spawned; on macOS the capture session is `!Send`, so await + /// [`Driver::run`] on a local task there. [`Control`] itself is `Send + Sync`, + /// so it can live anywhere. + pub fn new( + broadcast: moq_net::broadcast::Producer, + catalog: moq_mux::catalog::Producer, + options: CaptureOptions, + ) -> Result<(Self, Driver), Error> { + let suffix = match options.encode.codec { + Codec::H264 => ".avc3", + Codec::H265 => ".hev1", + }; + let track = broadcast.unique_track(suffix, catalog.track_info(hang::catalog::PRIORITY.video))?; + let cuts = kio::Producer::new(0); + let driver = Driver { + track, + catalog, + options, + cuts: cuts.consume(), + }; + Ok((Self { cuts }, driver)) + } + + /// Request a keyframe, opening a new group at a frame no earlier than this call. + /// + /// For a resume, a recording cut, or a known tune-in moment; [`Config::gop`](super::Config::gop) + /// is the cadence. Requests coalesce into the next keyframe, and forced keyframes land at least + /// 500ms apart, so a caller in a loop cannot pin the encoder at all-IDR; a request that arrives + /// too soon waits rather than being dropped. A request while nothing is watching is served by + /// the keyframe every fresh encoder opens with. A backend that cannot force one logs a warning + /// and keeps its GOP cadence; see [`Error::CutUnsupported`]. + pub fn cut(&self) { + if let Ok(mut cuts) = self.cuts.write() { + *cuts = cuts.wrapping_add(1); + } + } +} + +/// The task that captures, encodes, and publishes the track. +/// +/// Runs until the track ends, the capture fails, or the final [`Control`] drops. +#[cfg(feature = "capture")] +pub struct Driver { + track: moq_net::track::Producer, catalog: moq_mux::catalog::Producer, - capture: capture::Config, - encode: Options, - clock: moq_mux::Clock, -) -> Result<(), Error> { - // Open the camera once to find out what it actually negotiated, since a requested size is only a - // hint (macOS ignores it outright) and the encoder is built from the mode, not the request. It - // closes again immediately: this costs one camera open at startup and buys a rendition that says - // exactly what the stream will carry, rather than one every consumer has to treat as provisional. - let rendition = { - let camera = capture::open(&capture).await?; - let mut probe_config = encoder::Config::new( - camera.width(), - camera.height(), - capture - .framerate - .or_else(|| camera.framerate()) - .unwrap_or(DEFAULT_FRAMERATE), - ); - probe_config.bitrate = encode.bitrate; - probe_config.codec = encode.codec; - probe_config.kind = encode.kind.clone(); - probe_config.color = camera.color(); - probe_config.probe().await? - }; + options: CaptureOptions, + cuts: kio::Consumer, +} - let mut producer = Producer::new(broadcast, catalog, rendition)?; - let demand = producer.demand(); +#[cfg(feature = "capture")] +impl std::fmt::Debug for Driver { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("Driver").finish_non_exhaustive() + } +} - let result = capture_loop(&mut producer, &demand, &mut DeviceSource, &capture, &encode, &clock).await; +#[cfg(feature = "capture")] +impl Driver { + /// Run capture until the final control handle drops or the MoQ track ends. + /// + /// The source is opened once at startup to probe the mode it negotiates, then released until a + /// subscriber arrives and reopened for as long as one is watching. That one open is what lets + /// the catalog rendition be exact before a single frame is published, so a consumer can size + /// itself against it without waiting for an encoder that may never run. + pub async fn run(self) -> Result<(), Error> { + let Self { + track, + catalog, + options, + cuts, + } = self; + let CaptureOptions { capture, encode, clock } = options; + + // Open the camera once to find out what it actually negotiated, since a requested size is + // only a hint (macOS ignores it outright) and the encoder is built from the mode, not the + // request. It closes again immediately: this costs one camera open at startup and buys a + // rendition that says exactly what the stream will carry, rather than one every consumer has + // to treat as provisional. + let rendition = async { + let camera = capture::open(&capture).await?; + let mut probe_config = encoder::Config::new( + camera.width(), + camera.height(), + capture + .framerate + .or_else(|| camera.framerate()) + .unwrap_or(DEFAULT_FRAMERATE), + ); + probe_config.bitrate = encode.bitrate; + probe_config.codec = encode.codec; + probe_config.kind = encode.kind.clone(); + probe_config.color = camera.color(); + probe_config.probe().await + }; + let rendition = match rendition.await { + Ok(rendition) => rendition, + Err(err) => { + // A track that already ended has nobody left to tell. + let _ = track.abort(moq_net::Error::Transport(err.to_string())); + return Err(err); + } + }; - // This runs only when the loop ends on its own (the track is usually already - // going away by then); a Ctrl+C cancels the future before this point, since - // async `Drop` can't finalize the track. - match &result { - // Clean end (the track was dropped): best-effort finish. - Ok(()) => { - if let Err(err) = producer.finish() { - tracing::debug!(error = %err, "video track finish after capture ended"); + let mut producer = Producer::with_track(track, catalog, rendition)?; + let demand = producer.demand(); + + let result = capture_loop( + &mut producer, + &demand, + &mut DeviceSource, + &capture, + &encode, + &clock, + &cuts, + ) + .await; + + // This runs only when the loop ends on its own (the track is usually already + // going away by then); a Ctrl+C cancels the future before this point, since + // async `Drop` can't finalize the track. + match &result { + // Clean end (the track was dropped): best-effort finish. + Ok(()) => { + if let Err(err) = producer.finish() { + tracing::debug!(error = %err, "video track finish after capture ended"); + } } + // The capture loop failed: abort with the real cause so subscribers see it. + Err(err) => producer.abort(moq_net::Error::Transport(err.to_string())), } - // The capture loop failed: abort with the real cause so subscribers see it. - Err(err) => producer.abort(moq_net::Error::Transport(err.to_string())), + result } - result +} + +/// Capture a webcam and publish it as an on-demand video track. +/// +/// This convenience function runs a [`Control`] from `options` without handing +/// back the handle. Use [`Control::new`] directly to retain controls. +/// +/// Returns when the broadcast is dropped (the track stops being announced) or +/// the capture fails. Frames are stamped from +/// [`CaptureOptions::clock`], so passing the same [`Clock`](moq_mux::Clock) to +/// a concurrent audio publish keeps the two tracks aligned. +#[cfg(feature = "capture")] +pub async fn publish_capture( + broadcast: moq_net::broadcast::Producer, + catalog: moq_mux::catalog::Producer, + options: CaptureOptions, +) -> Result<(), Error> { + // Held, not dropped: the driver ends as soon as the last control handle goes. + let (_control, driver) = Control::new(broadcast, catalog, options)?; + driver.run().await } /// Off macOS, [`publish_capture`]'s future must stay `Send` so a server can @@ -350,12 +480,10 @@ pub async fn publish_capture( fn assert_publish_capture_send( broadcast: moq_net::broadcast::Producer, catalog: moq_mux::catalog::Producer, - capture: capture::Config, - encode: Options, - clock: moq_mux::Clock, + options: CaptureOptions, ) { fn is_send(_: &T) {} - is_send(&publish_capture(broadcast, catalog, capture, encode, clock)); + is_send(&publish_capture(broadcast, catalog, options)); } /// Where the capture loop opens its camera. Kept apart from the device backends so @@ -381,7 +509,7 @@ impl CaptureSource for DeviceSource { /// `None` rather than being absent. Retiring stops the `select!` arm from spinning on /// a channel that is permanently ready. #[cfg(feature = "capture")] -type RateControl = Option<(moq_net::bandwidth::Consumer, Control)>; +type RateControl = Option<(moq_net::bandwidth::Consumer, rate::Control)>; /// Wait for the next bandwidth estimate, or forever when rate control is off or /// finished. Cancel-safe: [`Consumer::changed`](moq_net::bandwidth::Consumer::changed) @@ -497,6 +625,7 @@ async fn capture_loop( capture: &capture::Config, encode: &Options, clock: &moq_mux::Clock, + cuts: &kio::Consumer, ) -> Result<(), Error> { // This track's claim on the connection. Taken on the first open, because the // negotiated mode is what finally says how much this encoder can ever send, and @@ -507,9 +636,15 @@ async fn capture_loop( // Idle until a viewer subscribes; the track ending is a clean exit. The // catalog rendition was published when the track was created, so a // subscriber can get here without a frame ever having been encoded. - if let Err(err) = demand.used().await { - log_track_ended(err); - return Ok(()); + tokio::select! { + biased; + res = demand.used() => { + if let Err(err) = res { + log_track_ended(err); + return Ok(()); + } + } + () = cuts.closed() => return Ok(()), } // Open the camera and an encoder sized to its negotiated mode. @@ -553,7 +688,9 @@ async fn capture_loop( // so the policy's ceiling is that rate and the target starts there. A // reopened camera starts optimistic again rather than inheriting the // backed-off rate from whatever the link was doing last time. - let mut rate = Some((reservation.consumer(), Control::new(Policy::new(ceiling)))); + let mut rate = Some((reservation.consumer(), rate::Control::new(rate::Policy::new(ceiling)))); + // Per encoder, so a reopen forgets the old encoder's last keyframe along with it. + let mut forced = Some(Cuts::new(*cuts.read())); loop { // Race the next frame against the last viewer leaving so we release the @@ -569,6 +706,8 @@ async fn capture_loop( } break; // no viewers: release the camera, then wait for one } + // The final control handle dropped: stop publishing. + () = cuts.closed() => return Ok(()), // Retune between frames rather than mid-encode, and only when // the policy says the target actually moved. estimate = next_estimate(&mut rate) => { @@ -589,6 +728,21 @@ async fn capture_loop( let Some(mut frame) = frame else { break }; frame.timestamp = map_capture_timestamp(capture_epoch, frame.timestamp)?; + let requests = *cuts.read(); + if forced + .as_mut() + .is_some_and(|forced| forced.due(requests, frame.timestamp)) + { + match encoder.cut().await { + Ok(()) => {} + // Keep capturing on the GOP cadence and stop asking this encoder. + Err(Error::CutUnsupported(name)) => { + tracing::warn!(encoder = name, "encoder cannot force a keyframe on request"); + forced = None; + } + Err(err) => return Err(err), + } + } let started = Instant::now(); let Some(encoded) = wait_capture(producer, demand, encoder.encode(frame)).await? else { break; @@ -954,10 +1108,23 @@ mod tests { ..Options::default() }; let config = capture::Config::default(); + // Hold the producer: dropping it closes the consumer and the loop + // treats that as the end of capture, ahead of this fixture's stop. + let cuts = kio::Producer::new(0); + let requests = cuts.consume(); tokio::select! { - res = capture_loop(&mut producer, &demand, &mut source, &config, &options, &clock) => res?, + res = capture_loop( + &mut producer, + &demand, + &mut source, + &config, + &options, + &clock, + &requests, + ) => res?, _ = stopped => {} } + drop(cuts); producer.finish() }); diff --git a/rs/moq-video/src/lib.rs b/rs/moq-video/src/lib.rs index dcc98e206c..be252fcac6 100644 --- a/rs/moq-video/src/lib.rs +++ b/rs/moq-video/src/lib.rs @@ -24,7 +24,9 @@ //! the matching `moq_mux::codec` importer, which handles catalog registration //! and framing. The codec is chosen via [`encode::Codec`]: H.264 (openh264 / //! VideoToolbox / Media Foundation / NVENC / VAAPI / V4L2) or H.265 -//! (VideoToolbox / Media Foundation / NVENC). Two entry points: +//! (VideoToolbox / Media Foundation / NVENC). Entry points: +//! - `encode::Control::new` returns a `Control` handle (ask for a keyframe +//! with `cut`) and the `Driver` that captures and publishes the webcam. //! - `encode::publish_capture` captures a webcam and publishes it (turnkey). //! It encodes strictly on demand: the track and catalog are advertised up //! front (the camera opens once at startup so they can be exact), and the