From c100e530938bbe264ea62df60fc90b2f50f2038f Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sun, 20 Sep 2026 08:19:52 -0700 Subject: [PATCH 1/3] chore: claim api watch publish quest Co-Authored-By: GPT-6 From 50bb3de4ce70b9559ad77103240fe26526e6bef6 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sun, 20 Sep 2026 08:46:28 -0700 Subject: [PATCH 2/3] refactor(js)!: unify watch and publish APIs Co-Authored-By: GPT-6 --- demo/web/src/meet.ts | 2 +- demo/web/src/publish.ts | 3 +- doc/lib/js/publish.md | 8 +- doc/lib/js/room.md | 2 +- doc/lib/js/watch.md | 6 +- js/moq-boy/src/game.ts | 23 ++-- js/publish/README.md | 6 +- js/publish/src/audio/encoder.ts | 10 +- js/publish/src/element.ts | 61 +++++----- js/publish/src/source/camera.ts | 7 +- js/publish/src/source/file.ts | 7 +- js/publish/src/source/index.ts | 1 + js/publish/src/source/microphone.ts | 7 +- js/publish/src/source/retry.test.ts | 6 +- js/publish/src/source/screen.ts | 5 +- js/publish/src/source/types.ts | 8 ++ js/publish/src/ui/components/stats-tab.ts | 6 +- js/publish/src/ui/components/status-badge.ts | 3 +- js/publish/src/ui/element.ts | 3 +- js/room/README.md | 2 +- js/room/src/local.test.ts | 15 ++- js/room/src/local.ts | 111 +++++++++--------- js/room/src/remote.ts | 16 ++- js/watch/README.md | 14 +-- js/watch/src/audio/decoder.ts | 25 +++- js/watch/src/audio/emitter.ts | 10 +- js/watch/src/audio/latency.test.ts | 6 +- js/watch/src/audio/latency.ts | 10 +- js/watch/src/broadcast.test.ts | 14 ++- js/watch/src/broadcast.ts | 29 +++-- js/watch/src/element.ts | 115 ++++++++----------- js/watch/src/enabled.test.ts | 6 +- js/watch/src/sync.test.ts | 33 ++++++ js/watch/src/sync.ts | 35 ++++-- js/watch/src/text/renderer.ts | 14 ++- js/watch/src/video/decoder.ts | 15 ++- js/watch/src/video/renderer.test.ts | 2 +- js/watch/src/video/renderer.ts | 10 +- js/watch/src/video/source.test.ts | 4 +- js/watch/src/video/source.ts | 14 ++- quest/m1/README.md | 1 - quest/m1/api-review-gate.md | 1 - quest/m1/api-watch-publish.md | 61 ---------- quest/m2/watch-player.md | 1 - 44 files changed, 386 insertions(+), 352 deletions(-) create mode 100644 js/publish/src/source/types.ts delete mode 100644 quest/m1/api-watch-publish.md diff --git a/demo/web/src/meet.ts b/demo/web/src/meet.ts index 6a161b20e1..274c75c192 100644 --- a/demo/web/src/meet.ts +++ b/demo/web/src/meet.ts @@ -127,7 +127,7 @@ function join(): void { enabled: true, }); local = new Local({ - origin: connection.origin, + connection, identity, enabled: true, user: { id: name, name }, diff --git a/demo/web/src/publish.ts b/demo/web/src/publish.ts index 6b45369f61..5b5ee56369 100644 --- a/demo/web/src/publish.ts +++ b/demo/web/src/publish.ts @@ -456,7 +456,8 @@ $("publish-graphs").append(captureGraph.el, uploadGraph.el, rttGraph.el); // a signal coalesces a burst into one notification, which undercounts the rate. let frames = 0; viz.run((effect) => { - const fanout = effect.get(publish.capture.out.frames); + const capture = effect.get(publish.video.in.capture); + const fanout = capture ? effect.get(capture.out.frames) : undefined; if (!fanout) return; const reader = fanout.subscribe(effect).getReader(); diff --git a/doc/lib/js/publish.md b/doc/lib/js/publish.md index 8dc1b9b71a..3b64e4092c 100644 --- a/doc/lib/js/publish.md +++ b/doc/lib/js/publish.md @@ -60,7 +60,7 @@ signals.run((effect) => { if (!net) return; // A day-long retention so a late viewer still replays the last value. - const track = net.createTrack("meta.json", { latencyMax: 86_400_000 }); + const track = net.createTrack("meta.json", { maxAge: 86_400_000 }); effect.cleanup(() => track.close()); const meta = new Json.Snapshot.Producer({ track }); @@ -97,13 +97,15 @@ const broadcast = new Publish.Broadcast({ const camera = new Publish.Source.Camera({ enabled: true }); const microphone = new Publish.Source.Microphone({ enabled: true }); -const capture = new Publish.Video.Capture({ source: camera.out.source }); +const video = new Publish.Signals.Computed((effect) => effect.get(camera.out.source)?.video); +const capture = new Publish.Video.Capture({ source: video }); // Each encoder registers a rendition on the broadcast (`broadcast.video(name)`) and // encodes only while someone is subscribed. new Publish.Video.Encoder("video/hd", { broadcast, capture, enabled: true }); new Publish.Video.Encoder("video/sd", { broadcast, capture, enabled: true, config: { maxScale: 0.25 } }); -const audioCapture = new Publish.Audio.Capture({ source: microphone.out.source }); +const audioSource = new Publish.Signals.Computed((effect) => effect.get(microphone.out.source)?.audio); +const audioCapture = new Publish.Audio.Capture({ source: audioSource }); new Publish.Audio.Encoder("audio", { broadcast, capture: audioCapture, enabled: true }); ``` diff --git a/doc/lib/js/room.md b/doc/lib/js/room.md index c88f767b6f..7ab1dbf43f 100644 --- a/doc/lib/js/room.md +++ b/doc/lib/js/room.md @@ -26,7 +26,7 @@ const connection = new Connection({ }); const local = new Local({ - origin: connection.origin, + connection, identity: Path.from("alice"), user: { name: "Alice" }, }); diff --git a/doc/lib/js/watch.md b/doc/lib/js/watch.md index 1d0da55e96..b812d81f3f 100644 --- a/doc/lib/js/watch.md +++ b/doc/lib/js/watch.md @@ -33,9 +33,8 @@ in sync at the latency you ask for. | `delay` | How far playback trails the live edge: `"auto"` (derived from RTT, the default), a duration like `"300ms"`, or `"instant"` to paint frames as they decode with no pacing at all. | | `buffer` | Future-dated media held beyond the live edge before playback skips ahead, e.g. `"30s"`. Defaults to none. | | `captions` | The caption track to show, or absent for off. `el.text.out.available` lists the renditions for a picker. | -| `jitter` | The jitter buffer in ms. | | `visible` | Only subscribe to video while the element is on screen: a margin (`"20%"` default, `"200px"`), `"always"`, or `"never"`. | -| `reload` | Wait for the broadcast to be announced before subscribing (default on), so a player can be mounted before the stream exists. | +| `announced` | Wait for the broadcast to be announced before subscribing (default on), so a player can be mounted before the stream exists. | | `catalog-format` | `hang` (default, from the `.hang` suffix), `hangz` (compressed), `msf`, or `manual` to supply the catalog yourself. | The overlay adds play/pause, volume, fullscreen, a quality selector, a @@ -143,7 +142,8 @@ const broadcast = new Watch.Broadcast({ origin: connection.origin, name: Moq.Pat ``` `Watch.Broadcast`, `Video.Decoder`, `Video.Renderer`, `Audio.Decoder`, and -`Audio.Emitter` are the pieces the element assembles; every input and output +`Audio.Emitter` are the pieces the element assembles. Their constructors take +one properties object, and every input and output is a signal from [`@moq/signals`](/lib/js/signals). Load from a CDN (`https://esm.sh/@moq/watch/element`) for a no-build embed. diff --git a/js/moq-boy/src/game.ts b/js/moq-boy/src/game.ts index 6f4e37a72d..c9bec3f4d5 100644 --- a/js/moq-boy/src/game.ts +++ b/js/moq-boy/src/game.ts @@ -136,26 +136,24 @@ export class Game { }); this.#signals.cleanup(() => this.audioSource.close()); - // The decoder owns rendition handoffs but also needs Sync. Bridge its jitter output through a - // local signal so Sync can be constructed first. - const videoJitter = new Moq.Signals.Signal(undefined); this.sync = new Watch.Sync({ delay: this.delay, probe: connection.probe, - video: videoJitter, - audio: this.audioSource.out.jitter, }); this.#signals.cleanup(() => this.sync.close()); this.#signals.run(this.#runPixelBudget.bind(this)); const videoEnabled = new Moq.Signals.Signal(true); - this.videoDecoder = new Watch.Video.Decoder(this.videoSource, this.sync, { enabled: videoEnabled }); - this.#signals.proxy(videoJitter, this.videoDecoder.out.jitter); + this.videoDecoder = new Watch.Video.Decoder({ + source: this.videoSource, + sync: this.sync, + enabled: videoEnabled, + }); this.#signals.cleanup(() => this.videoDecoder.close()); // Renderer needs a canvas created by the UI layer, set via `canvas`. - this.videoRenderer = new Watch.Video.Renderer(this.videoDecoder, { canvas: this.canvas }); + this.videoRenderer = new Watch.Video.Renderer({ decoder: this.videoDecoder, canvas: this.canvas }); this.#signals.cleanup(() => this.videoRenderer.close()); // Download on the grid or when expanded, but only while the tile is on-screen. @@ -164,13 +162,18 @@ export class Game { // Audio pipeline. The emitter stops the download when muted or paused. const audioEnabled = new Moq.Signals.Signal(false); - this.audioDecoder = new Watch.Audio.Decoder(this.audioSource, this.sync, { enabled: audioEnabled }); + this.audioDecoder = new Watch.Audio.Decoder({ + source: this.audioSource, + sync: this.sync, + enabled: audioEnabled, + }); this.#signals.cleanup(() => this.audioDecoder.close()); const audioPaused = new Moq.Signals.Signal(true); this.#signals.run(this.#runAudioPaused.bind(this, audioPaused)); - this.audioEmitter = new Watch.Audio.Emitter(this.audioDecoder, { + this.audioEmitter = new Watch.Audio.Emitter({ + source: this.audioDecoder, volume: this.volume, muted: this.userMuted, paused: audioPaused, diff --git a/js/publish/README.md b/js/publish/README.md index 2731345139..b5b13f9930 100644 --- a/js/publish/README.md +++ b/js/publish/README.md @@ -100,7 +100,8 @@ const broadcast = new Publish.Broadcast({ // Capture, then encode. Each encoder registers its rendition on the broadcast // and encodes only while someone is subscribed. const camera = new Publish.Source.Camera({ enabled: true }); -const capture = new Publish.Video.Capture({ source: camera.out.source }); +const video = new Publish.Signals.Computed((effect) => effect.get(camera.out.source)?.video); +const capture = new Publish.Video.Capture({ source: video }); const hd = new Publish.Video.Encoder("video/hd", { broadcast, capture, enabled: true }); const sd = new Publish.Video.Encoder("video/sd", { broadcast, capture, enabled: true, config: { maxScale: 0.25 } }); @@ -109,7 +110,8 @@ const sd = new Publish.Video.Encoder("video/sd", { broadcast, capture, enabled: hd.config.set({ codec: "vp09.00.10.08", maxBitrate: 4_000_000 }); const microphone = new Publish.Source.Microphone({ enabled: true }); -const audioCapture = new Publish.Audio.Capture({ source: microphone.out.source }); +const audioSource = new Publish.Signals.Computed((effect) => effect.get(microphone.out.source)?.audio); +const audioCapture = new Publish.Audio.Capture({ source: audioSource }); const audio = new Publish.Audio.Encoder("audio", { broadcast, capture: audioCapture, enabled: true }); audio.volume.set(0.8); ``` diff --git a/js/publish/src/audio/encoder.ts b/js/publish/src/audio/encoder.ts index 315057ff0f..21ceed0ec4 100644 --- a/js/publish/src/audio/encoder.ts +++ b/js/publish/src/audio/encoder.ts @@ -84,7 +84,6 @@ export type EncoderInput = { /** Constructor options: the wired inputs plus the live-editable tuning knobs. */ export type EncoderProps = Inputs & { // User tuning knobs. Seed a value or wire a Signal; also live-editable via the matching field. - muted?: boolean | Signal; volume?: number | Signal; // Codec selection plus encoder settings. Defaults to "opus". @@ -123,8 +122,6 @@ export class Encoder { readonly in: Readonlys; - /** Silence the encoded audio without tearing down the capture graph. */ - muted: Signal; /** Linear gain applied before encoding, where 1 is unity. */ volume: Signal; /** The live-editable codec selection plus its encoder settings. */ @@ -179,7 +176,6 @@ export class Encoder { capture: getter(props?.capture), bandwidth: getter(props?.bandwidth), }; - this.muted = Signal.from(props?.muted ?? false); this.volume = Signal.from(props?.volume ?? 1); this.codec = Signal.from(props?.codec ?? "opus"); @@ -198,7 +194,7 @@ export class Encoder { // Pump PCM off the capture into whatever is currently publishing, applying the volume knobs on // the way through. Tied to the capture's lifetime rather than the encoder's, so reconfiguring - // (or muting) never has to reacquire the stream, which for a decoded file would be fatal. + // never has to reacquire the stream, which for a decoded file would be fatal. #runCapture(effect: Effect): void { const capture = effect.get(this.in.capture); if (!capture) return; @@ -213,7 +209,7 @@ export class Encoder { reader.cancel().catch(() => {}); }); - const gain = new Gain(this.muted.peek() ? 0 : this.volume.peek()); + const gain = new Gain(this.volume.peek()); effect.spawn(async () => { for (;;) { @@ -225,7 +221,7 @@ export class Encoder { // Every rendition shares the captured frame, so gain returns a copy rather than // scaling in place; muting one rendition must not silence the rest. - gain.set(this.muted.peek() ? 0 : this.volume.peek()); + gain.set(this.volume.peek()); const frame = gain.apply(next.value, format.sampleRate); // The config rebuilds when the channel count moves, so skip anything that arrives diff --git a/js/publish/src/element.ts b/js/publish/src/element.ts index f1839dc456..28df9edac9 100644 --- a/js/publish/src/element.ts +++ b/js/publish/src/element.ts @@ -6,7 +6,7 @@ * @module */ import * as Moq from "@moq/net"; -import { Effect, Signal } from "@moq/signals"; +import { Effect, readonlys, Signal } from "@moq/signals"; import * as Audio from "./audio"; import { Broadcast } from "./broadcast"; import * as Preview from "./preview"; @@ -89,8 +89,7 @@ export default class MoqPublish extends HTMLElement { * page and URL resolves it locally with no round trip. */ connection: Moq.Connection; - /** The video capture, shared by every video rendition. Also reachable as `video.capture`. */ - capture: Video.Capture; + #capture: Video.Capture; broadcast: Broadcast; // The single video and audio encoders. For multiple renditions (e.g. simulcast), drop the element @@ -101,11 +100,12 @@ export default class MoqPublish extends HTMLElement { // The selected input sources: the Camera/Screen, Microphone/Screen, and File holders driving capture. // Read by the UI (device pickers) and written by #runSource. - sources = { + readonly #sources = { video: new Signal(undefined), audio: new Signal(undefined), file: new Signal(undefined), }; + readonly sources = readonlys(this.#sources); // The captured media tracks, written by #runSource. Fed to the video and audio captures, so // consumers read them back via `capture.in.source` / `audio.capture.in.source` rather than here. @@ -174,8 +174,8 @@ export default class MoqPublish extends HTMLElement { this.#announcing.set(announce === "always" || (announce === "source" && hasMedia)); }); - this.capture = new Video.Capture({ source: this.#videoSource }); - this.signals.cleanup(() => this.capture.close()); + this.#capture = new Video.Capture({ source: this.#videoSource }); + this.signals.cleanup(() => this.#capture.close()); // Reached as `audio.capture` rather than a field of its own, so audio and video read alike. const audioCapture = new Audio.Capture({ @@ -189,14 +189,14 @@ export default class MoqPublish extends HTMLElement { enabled: this.#enabled, announce: this.#announcing, name: this.#name, - display: this.capture.out.display, + display: this.#capture.out.display, flip: this.#flip, }); this.signals.cleanup(() => this.broadcast.close()); this.video = new Video.Encoder("video", { broadcast: this.broadcast, - capture: this.capture, + capture: this.#capture, enabled: this.#videoEnabled, bandwidth: this.connection.bandwidth, }); @@ -227,8 +227,8 @@ export default class MoqPublish extends HTMLElement { if (preview instanceof HTMLCanvasElement) { const renderer = new Preview.Renderer({ canvas: preview, - frames: this.capture.out.frames, - display: this.capture.out.display, + frames: this.#capture.out.frames, + display: this.#capture.out.display, flip: this.#flip, encoder: this.video, mode: this.controls.preview, @@ -314,19 +314,14 @@ export default class MoqPublish extends HTMLElement { if (source === "camera") { const video = new Source.Camera({ enabled: this.#videoEnabled }); - this.signals.run((effect) => { - const source = effect.get(video.out.source); - this.#videoSource.set(source); - }); - const audio = new Source.Microphone({ enabled: this.#audioEnabled }); - this.signals.run((effect) => { - const source = effect.get(audio.out.source); - this.#audioSource.set(source); - }); - effect.set(this.sources.video, video); - effect.set(this.sources.audio, audio); + effect.set(this.#sources.video, video); + effect.set(this.#sources.audio, audio); + effect.run((nested) => { + nested.set(this.#videoSource, nested.get(video.out.source)?.video); + nested.set(this.#audioSource, nested.get(audio.out.source)?.audio); + }); effect.cleanup(() => { video.close(); @@ -341,16 +336,14 @@ export default class MoqPublish extends HTMLElement { enabled: this.#eitherEnabled, }); - this.signals.run((effect) => { - const source = effect.get(screen.out.source); - if (!source) return; - - effect.set(this.#videoSource, source.video); - effect.set(this.#audioSource, source.audio); + effect.run((nested) => { + const media = nested.get(screen.out.source); + nested.set(this.#videoSource, media?.video); + nested.set(this.#audioSource, media?.audio); }); - effect.set(this.sources.video, screen); - effect.set(this.sources.audio, screen); + effect.set(this.#sources.video, screen); + effect.set(this.#sources.audio, screen); effect.cleanup(() => { screen.close(); @@ -372,12 +365,12 @@ export default class MoqPublish extends HTMLElement { fileSource.prompt(); } - effect.set(this.sources.file, fileSource); + effect.set(this.#sources.file, fileSource); - this.signals.run((effect) => { - const source = effect.get(fileSource.out.source); - this.#videoSource.set(source.video); - this.#audioSource.set(source.audio); + effect.run((nested) => { + const media = nested.get(fileSource.out.source); + nested.set(this.#videoSource, media?.video); + nested.set(this.#audioSource, media?.audio); }); effect.cleanup(() => { diff --git a/js/publish/src/source/camera.ts b/js/publish/src/source/camera.ts index a385cf9e13..ad1e829e47 100644 --- a/js/publish/src/source/camera.ts +++ b/js/publish/src/source/camera.ts @@ -2,6 +2,7 @@ import { Effect, type Getter, getter, type Inputs, type Readonlys, readonlys, Si import type * as Video from "../video"; import { Device, type DeviceProps } from "./device"; import { Retry } from "./retry"; +import type { Media } from "./types"; // Signals the camera reads. export type CameraInput = { @@ -19,7 +20,7 @@ export interface CameraProps extends Inputs { type CameraOutput = { // The live camera track, or undefined while disabled or denied. - source: Signal; + source: Signal; }; /** Captures video from a camera, tracking the available devices. */ @@ -45,7 +46,7 @@ export class Camera { constraints: Signal; readonly #out: CameraOutput = { - source: new Signal(undefined), + source: new Signal(undefined), }; readonly out = readonlys(this.#out); @@ -122,7 +123,7 @@ export class Camera { if (!source || source.readyState === "ended") return this.#retry.failed(); this.#retry.succeeded(effect, source); - effect.set(this.#out.source, source); + effect.set(this.#out.source, { video: source }); }); } diff --git a/js/publish/src/source/file.ts b/js/publish/src/source/file.ts index 4183b6b457..6b2672c763 100644 --- a/js/publish/src/source/file.ts +++ b/js/publish/src/source/file.ts @@ -14,6 +14,7 @@ import { import type * as Audio from "../audio"; import type * as Video from "../video"; import { Timeline } from "./timeline"; +import type { Media } from "./types"; // Signals the file source reads. export type FileInput = { @@ -28,8 +29,8 @@ export interface FileProps extends Inputs { } type FileOutput = { - // The sources decoded from the file, empty while disabled or undecodable. - source: Signal<{ video?: Video.Source; audio?: Audio.Source }>; + // The sources decoded from the file, undefined while disabled or undecodable. + source: Signal; }; // Image, video, and audio files we know how to decode (see #decode). @@ -56,7 +57,7 @@ export class File { file: Signal; readonly #out: FileOutput = { - source: new Signal<{ video?: Video.Source; audio?: Audio.Source }>({}), + source: new Signal(undefined), }; readonly out = readonlys(this.#out); diff --git a/js/publish/src/source/index.ts b/js/publish/src/source/index.ts index 2c6b097942..b19e662279 100644 --- a/js/publish/src/source/index.ts +++ b/js/publish/src/source/index.ts @@ -3,3 +3,4 @@ export * from "./device"; export * from "./file"; export * from "./microphone"; export * from "./screen"; +export * from "./types"; diff --git a/js/publish/src/source/microphone.ts b/js/publish/src/source/microphone.ts index 9c20c77b98..658eb85d9d 100644 --- a/js/publish/src/source/microphone.ts +++ b/js/publish/src/source/microphone.ts @@ -2,6 +2,7 @@ import { Effect, type Getter, getter, type Inputs, type Readonlys, readonlys, Si import type * as Audio from "../audio"; import { Device, type DeviceProps } from "./device"; import { Retry } from "./retry"; +import type { Media } from "./types"; // Signals the microphone reads. export type MicrophoneInput = { @@ -19,7 +20,7 @@ export interface MicrophoneProps extends Inputs { type MicrophoneOutput = { // The live microphone track, or undefined while disabled or denied. - source: Signal; + source: Signal; }; /** Captures audio from a microphone, tracking the available devices. */ @@ -33,7 +34,7 @@ export class Microphone { constraints: Signal; readonly #out: MicrophoneOutput = { - source: new Signal(undefined), + source: new Signal(undefined), }; readonly out = readonlys(this.#out); @@ -109,7 +110,7 @@ export class Microphone { if (!track || track.readyState === "ended") return this.#retry.failed(); this.#retry.succeeded(effect, track); - effect.set(this.#out.source, { track, kind: "voice" }); + effect.set(this.#out.source, { audio: { track, kind: "voice" } }); }); } diff --git a/js/publish/src/source/retry.test.ts b/js/publish/src/source/retry.test.ts index 585885aa96..70a7ab3e18 100644 --- a/js/publish/src/source/retry.test.ts +++ b/js/publish/src/source/retry.test.ts @@ -1,11 +1,10 @@ import { expect, spyOn, test } from "bun:test"; import { Signal } from "@moq/signals"; -import type * as Audio from "../audio"; -import type * as Video from "../video"; import { Camera } from "./camera"; import { Microphone } from "./microphone"; import { Retry } from "./retry"; import { Screen } from "./screen"; +import type { Media } from "./types"; // A MediaStreamTrack ends on its own when the device disappears or the OS revokes it. Only the bits // the sources touch, plus end() to fire it. @@ -170,7 +169,8 @@ const QUIET_MARGIN = 100; const SPENT_TIMEOUT = 30_000; /** The track a source published, or undefined. */ -function published(source: Audio.Source | Video.Source | undefined): unknown { +function published(media: Media | undefined): unknown { + const source = media?.audio ?? media?.video; if (!source) return undefined; return "track" in source ? source.track : source; } diff --git a/js/publish/src/source/screen.ts b/js/publish/src/source/screen.ts index 06f5e441f2..e885db1680 100644 --- a/js/publish/src/source/screen.ts +++ b/js/publish/src/source/screen.ts @@ -1,6 +1,7 @@ import { Effect, type Getter, getter, type Inputs, type Readonlys, readonlys, Signal } from "@moq/signals"; import type * as Audio from "../audio"; import type * as Video from "../video"; +import type { Media } from "./types"; // Signals the screen capture reads. export type ScreenInput = { @@ -19,7 +20,7 @@ export interface ScreenProps extends Inputs { type ScreenOutput = { // The captured surface, or undefined while disabled or dismissed. - source: Signal<{ audio?: Audio.Source; video?: Video.Source } | undefined>; + source: Signal; }; /** Captures a screen, window, or tab that the user picks. */ @@ -32,7 +33,7 @@ export class Screen { audio: Signal; readonly #out: ScreenOutput = { - source: new Signal<{ audio?: Audio.Source; video?: Video.Source } | undefined>(undefined), + source: new Signal(undefined), }; readonly out = readonlys(this.#out); diff --git a/js/publish/src/source/types.ts b/js/publish/src/source/types.ts new file mode 100644 index 0000000000..778f0d5a79 --- /dev/null +++ b/js/publish/src/source/types.ts @@ -0,0 +1,8 @@ +import type * as Audio from "../audio"; +import type * as Video from "../video"; + +/** Audio and video captured by a publish source. */ +export type Media = { + video?: Video.Source; + audio?: Audio.Source; +}; diff --git a/js/publish/src/ui/components/stats-tab.ts b/js/publish/src/ui/components/stats-tab.ts index 7a7a8a4d53..bdb9ee4943 100644 --- a/js/publish/src/ui/components/stats-tab.ts +++ b/js/publish/src/ui/components/stats-tab.ts @@ -70,7 +70,8 @@ export function statsTab(parent: Effect, publish: MoqPublish): HTMLElement { // Resolution/codec from the live capture (display) + catalog; card hides when not capturing video. parent.run((effect) => { - const display = effect.get(publish.capture.out.display); + const capture = effect.get(publish.video.in.capture); + const display = capture ? effect.get(capture.out.display) : undefined; const cfg = effect.get(publish.video.out.catalog); videoCard.el.style.display = display ? "" : "none"; vRes.textContent = display ? `${display.width}×${display.height}` : "—"; @@ -103,7 +104,8 @@ export function statsTab(parent: Effect, publish: MoqPublish): HTMLElement { // one notification and undercount the rate. let frames = 0; parent.run((effect) => { - const fanout = effect.get(publish.capture.out.frames); + const capture = effect.get(publish.video.in.capture); + const fanout = capture ? effect.get(capture.out.frames) : undefined; if (!fanout) return; const reader = fanout.subscribe(effect).getReader(); diff --git a/js/publish/src/ui/components/status-badge.ts b/js/publish/src/ui/components/status-badge.ts index 948f154ba7..e7debd1c27 100644 --- a/js/publish/src/ui/components/status-badge.ts +++ b/js/publish/src/ui/components/status-badge.ts @@ -34,7 +34,8 @@ export function statusBadge(parent: Effect, publish: MoqPublish): HTMLElement { const status = effect.get(publish.connection.status); const audioCapture = effect.get(publish.audio.in.capture); const audioSource = audioCapture ? effect.get(audioCapture.in.source) : undefined; - const videoSource = effect.get(publish.capture.in.source); + const videoCapture = effect.get(publish.video.in.capture); + const videoSource = videoCapture ? effect.get(videoCapture.in.source) : undefined; const muted = effect.get(publish.controls.muted); const invisible = effect.get(publish.controls.invisible); diff --git a/js/publish/src/ui/element.ts b/js/publish/src/ui/element.ts index 4943adecb6..6094039060 100644 --- a/js/publish/src/ui/element.ts +++ b/js/publish/src/ui/element.ts @@ -99,7 +99,8 @@ export default class MoqPublishUi extends HTMLElement { // last live preview's aspect ratio (defaults to 16:9) instead of snapping. let lastAspect = 16 / 9; effect.run((e) => { - const src = e.get(publish.capture.in.source); + const capture = e.get(publish.video.in.capture); + const src = capture ? e.get(capture.in.source) : undefined; if (src) { player.classList.remove("player--empty"); player.style.aspectRatio = ""; diff --git a/js/room/README.md b/js/room/README.md index 31360d9686..1973440cf2 100644 --- a/js/room/README.md +++ b/js/room/README.md @@ -49,7 +49,7 @@ const connection = new Connection({ const identity = Path.from("alice"); const local = new Local({ - origin: connection.origin, + connection, identity, user: { name: "Alice" }, }); diff --git a/js/room/src/local.test.ts b/js/room/src/local.test.ts index 598c67a2d5..63325dce7f 100644 --- a/js/room/src/local.test.ts +++ b/js/room/src/local.test.ts @@ -3,6 +3,7 @@ import { Path } from "@moq/net"; import { Signal } from "@moq/signals"; const sources: FakeSource[] = []; +const encoders: Record[] = []; class FakeSource { out = { source: new Signal(undefined) }; constructor() { @@ -12,6 +13,9 @@ class FakeSource { } class FakePipeline { out = { display: new Signal(undefined), frame: new Signal(undefined) }; + constructor(name?: unknown, props?: Record) { + if (typeof name === "string" && props) encoders.push(props); + } close() {} } class FakeBroadcast { @@ -32,9 +36,18 @@ async function flush() { } test("screen capture stays enabled while pending and resets after a live share ends", async () => { - const local = new Local({ origin: undefined, identity: Path.from("alice") }); + const bandwidth = new Signal(undefined); + const local = new Local({ + connection: { + origin: new Signal(undefined), + bandwidth, + } as never, + identity: Path.from("alice"), + }); try { await flush(); + expect(encoders).toHaveLength(6); + expect(encoders.every((props) => props.bandwidth === bandwidth)).toBe(true); local.screenEnabled.set(true); await flush(); expect(local.screenEnabled.peek()).toBe(true); diff --git a/js/room/src/local.ts b/js/room/src/local.ts index 7e81353329..05a814ea49 100644 --- a/js/room/src/local.ts +++ b/js/room/src/local.ts @@ -12,18 +12,18 @@ import { broadcastPath, KIND } from "./path.ts"; /** Constructor options for {@link Local}. */ export interface LocalProps { - /** Origin to publish into, usually a `Connection`'s `origin`. */ - origin: GetterInit; + /** Reconnecting connection to publish through. */ + connection: Moq.Connection; /** Participant identity; broadcast names are `{identity}/camera.hang` and `{identity}/screen.hang`. */ identity: GetterInit; /** When true, announce the camera broadcast (joining the room). Defaults to false. */ - enabled?: boolean | Signal; + enabled?: GetterInit; /** Capture the camera. Pass a Signal to share it with the app (hang.live Settings). */ - cameraEnabled?: boolean | Signal; + cameraEnabled?: GetterInit; /** Capture the microphone. Pass a Signal to share it with the app. */ - microphoneEnabled?: boolean | Signal; + microphoneEnabled?: GetterInit; /** Prompt for and capture a screen. Pass a Signal to share it with the app. */ - screenEnabled?: boolean | Signal; + screenEnabled?: GetterInit; /** Seed the published user.json fields. */ user?: UserProps; } @@ -73,43 +73,36 @@ export class Local { readonly cameraCapture: Publish.Video.Capture; /** Shared capture feeding the screen renditions. */ readonly screenCapture: Publish.Video.Capture; - /** Shared capture feeding the camera microphone encoder. */ - readonly cameraAudioCapture: Publish.Audio.Capture; - /** Shared capture feeding the screen audio encoder. */ - readonly screenAudioCapture: Publish.Audio.Capture; - - /** Camera HD encoder. */ - readonly cameraHd: Publish.Video.Encoder; - /** Camera SD encoder. */ - readonly cameraSd: Publish.Video.Encoder; - /** Camera microphone encoder. */ - readonly cameraAudio: Publish.Audio.Encoder; - - /** Screen HD encoder. */ - readonly screenHd: Publish.Video.Encoder; - /** Screen SD encoder. */ - readonly screenSd: Publish.Video.Encoder; - /** Screen audio encoder (tab/system audio when the share includes it). */ - readonly screenAudio: Publish.Audio.Encoder; - #preview = new Signal({}); + #cameraVideo = new Signal(undefined); + #cameraAudio = new Signal(undefined); #screenVideo = new Signal(undefined); #screenAudioSource = new Signal(undefined); #screenLive = new Signal(false); #signals = new Effect(); + #control(value: GetterInit | undefined): Signal { + const input = getter(value ?? false); + if (input instanceof Signal) return input; + + const output = new Signal(input.peek()); + this.#signals.proxy(output, input); + return output; + } + constructor(props: LocalProps) { this.identity = getter(props.identity); - this.enabled = Signal.from(props.enabled ?? false); - this.cameraEnabled = Signal.from(props.cameraEnabled ?? false); - this.microphoneEnabled = Signal.from(props.microphoneEnabled ?? false); - this.screenEnabled = Signal.from(props.screenEnabled ?? false); + this.enabled = this.#control(props.enabled); + this.cameraEnabled = this.#control(props.cameraEnabled); + this.microphoneEnabled = this.#control(props.microphoneEnabled); + this.screenEnabled = this.#control(props.screenEnabled); this.typing = new Signal(false); this.chatting = new Signal(false); this.user = userFields(props.user); - const origin = getter(props.origin); + const origin = props.connection.origin; + const bandwidth = props.connection.bandwidth; this.webcam = new Publish.Source.Camera({ enabled: this.cameraEnabled, @@ -148,8 +141,12 @@ export class Local { }, }); this.#signals.cleanup(() => this.share.close()); + this.#signals.run((effect) => { + effect.set(this.#cameraVideo, effect.get(this.webcam.out.source)?.video); + effect.set(this.#cameraAudio, effect.get(this.microphone.out.source)?.audio); + }); - this.cameraCapture = new Publish.Video.Capture({ source: this.webcam.out.source }); + this.cameraCapture = new Publish.Video.Capture({ source: this.#cameraVideo }); this.#signals.cleanup(() => this.cameraCapture.close()); this.screenCapture = new Publish.Video.Capture({ source: this.#screenVideo }); @@ -180,68 +177,74 @@ export class Local { }); this.#signals.cleanup(() => this.screen.close()); - this.cameraHd = new Publish.Video.Encoder("video/hd", { + const cameraHd = new Publish.Video.Encoder("video/hd", { broadcast: this.camera, capture: this.cameraCapture, enabled: this.cameraEnabled, + bandwidth, config: { maxPixels: 1280 * 720 }, }); - this.#signals.cleanup(() => this.cameraHd.close()); + this.#signals.cleanup(() => cameraHd.close()); - this.cameraSd = new Publish.Video.Encoder("video/sd", { + const cameraSd = new Publish.Video.Encoder("video/sd", { broadcast: this.camera, capture: this.cameraCapture, enabled: this.cameraEnabled, + bandwidth, config: { maxPixels: 640 * 360 }, }); - this.#signals.cleanup(() => this.cameraSd.close()); + this.#signals.cleanup(() => cameraSd.close()); - this.cameraAudioCapture = new Publish.Audio.Capture({ - source: this.microphone.out.source, + const cameraAudioCapture = new Publish.Audio.Capture({ + source: this.#cameraAudio, enabled: this.microphoneEnabled, }); - this.#signals.cleanup(() => this.cameraAudioCapture.close()); + this.#signals.cleanup(() => cameraAudioCapture.close()); - this.cameraAudio = new Publish.Audio.Encoder("audio", { + const cameraAudio = new Publish.Audio.Encoder("audio", { broadcast: this.camera, - capture: this.cameraAudioCapture, + capture: cameraAudioCapture, enabled: this.microphoneEnabled, + bandwidth, }); - this.#signals.cleanup(() => this.cameraAudio.close()); + this.#signals.cleanup(() => cameraAudio.close()); - this.screenHd = new Publish.Video.Encoder("video/hd", { + const screenHd = new Publish.Video.Encoder("video/hd", { broadcast: this.screen, capture: this.screenCapture, enabled: this.#screenLive, + bandwidth, config: { maxPixels: 1920 * 1080 }, }); - this.#signals.cleanup(() => this.screenHd.close()); + this.#signals.cleanup(() => screenHd.close()); - this.screenSd = new Publish.Video.Encoder("video/sd", { + const screenSd = new Publish.Video.Encoder("video/sd", { broadcast: this.screen, capture: this.screenCapture, enabled: this.#screenLive, + bandwidth, config: { maxPixels: 960 * 540 }, }); - this.#signals.cleanup(() => this.screenSd.close()); + this.#signals.cleanup(() => screenSd.close()); - this.screenAudioCapture = new Publish.Audio.Capture({ + const screenAudioCapture = new Publish.Audio.Capture({ source: this.#screenAudioSource, enabled: this.#screenLive, }); - this.#signals.cleanup(() => this.screenAudioCapture.close()); + this.#signals.cleanup(() => screenAudioCapture.close()); - this.screenAudio = new Publish.Audio.Encoder("audio", { + const screenAudio = new Publish.Audio.Encoder("audio", { broadcast: this.screen, - capture: this.screenAudioCapture, + capture: screenAudioCapture, enabled: this.#screenLive, + bandwidth, }); - this.#signals.cleanup(() => this.screenAudio.close()); + this.#signals.cleanup(() => screenAudio.close()); this.#signals.run((effect) => { const source = effect.get(this.share.out.source); - this.#screenVideo.set(source?.video); - this.#screenAudioSource.set(source?.audio); + effect.set(this.#screenVideo, source?.video); + effect.set(this.#screenAudioSource, source?.audio); const live = !!source?.video || !!source?.audio; const wasLive = this.#screenLive.peek(); this.#screenLive.set(live); @@ -252,8 +255,8 @@ export class Local { this.#signals.run((effect) => { this.#preview.set({ - video: !!effect.get(this.webcam.out.source), - audio: !!effect.get(this.microphone.out.source), + video: !!effect.get(this.webcam.out.source)?.video, + audio: !!effect.get(this.microphone.out.source)?.audio, screen: effect.get(this.#screenLive), name: effect.get(this.user.name), avatar: effect.get(this.user.avatar), diff --git a/js/room/src/remote.ts b/js/room/src/remote.ts index 639fc0f35e..45a9a09687 100644 --- a/js/room/src/remote.ts +++ b/js/room/src/remote.ts @@ -53,7 +53,7 @@ export class Member { origin: connection.origin, enabled: true, name: path, - reload: true, + announced: true, }); this.#signals.cleanup(() => this.broadcast.close()); @@ -71,27 +71,25 @@ export class Member { audioSource.close(); }); - const videoJitter = new Signal(undefined); const sync = new Watch.Sync({ delay: "auto", probe: connection.probe, - video: videoJitter, - audio: audioSource.out.jitter, }); this.#signals.cleanup(() => sync.close()); - this.video = new Watch.Video.Decoder(videoSource, sync, { enabled: this.#videoEnabled }); - this.#signals.proxy(videoJitter, this.video.out.jitter); - this.audio = new Watch.Audio.Decoder(audioSource, sync, { enabled: this.#audioEnabled }); + this.video = new Watch.Video.Decoder({ source: videoSource, sync, enabled: this.#videoEnabled }); + this.audio = new Watch.Audio.Decoder({ source: audioSource, sync, enabled: this.#audioEnabled }); this.#signals.cleanup(() => { this.video.close(); this.audio.close(); }); - this.renderer = new Watch.Video.Renderer(this.video, { + this.renderer = new Watch.Video.Renderer({ + decoder: this.video, canvas: this.canvas, }); - this.emitter = new Watch.Audio.Emitter(this.audio, { + this.emitter = new Watch.Audio.Emitter({ + source: this.audio, volume: this.volume, muted: this.muted, }); diff --git a/js/watch/README.md b/js/watch/README.md index a758b87e7f..8f35f001ab 100644 --- a/js/watch/README.md +++ b/js/watch/README.md @@ -69,10 +69,10 @@ The simplest way to watch a stream: | `muted` | boolean | false | Mute audio | | `visible` | never, distance, or always | `20%` | When to download video (see below) | | `volume` | number | 0.5 | Audio volume (0-1) | -| `reload` | boolean | true | Wait for (re)announcement before subscribing. Ignored when the relay does not support broadcast discovery. | -| `latency` | `real-time`, ms, `instant` | `real-time` | Target latency. `instant` paints frames as they decode and disables audio. | -| `latency-min` | `real-time` or ms | `real-time` | The latency floor, opening a range instead of a single target. | -| `latency-max` | `real-time` or ms | `real-time` | The latency ceiling: buffer freely below it, skip ahead past it. | +| `announced` | boolean | true | Wait for (re)announcement before subscribing. Ignored when the relay does not support broadcast discovery. | +| `delay` | `auto`, duration, `instant` | `auto` | Distance from the live edge. `instant` paints frames as they decode and disables audio. | +| `buffer` | duration | `0ms` | Future-dated media held before playback skips ahead. | +| `captions` | string | off | Text rendition to render. | | `catalog-format` | hang, hangz, msf, manual | auto-detected | The catalog format; detected from the name suffix unless set. `hangz` (compressed) is opt-in. | The `visible` attribute controls when the video track is downloaded, based on the canvas @@ -112,10 +112,10 @@ const broadcast = new Watch.Broadcast({ const source = new Watch.Video.Source({ broadcast, supported: Watch.Video.Decoder.supported, probe: connection.probe }); const sync = new Watch.Sync({ probe: connection.probe }); -const decoder = new Watch.Video.Decoder(source, sync, { enabled: true }); +const decoder = new Watch.Video.Decoder({ source, sync, enabled: true }); // Video renders to a ; there is no MediaStream to assign. -const renderer = new Watch.Video.Renderer(decoder, { canvas }); +const renderer = new Watch.Video.Renderer({ decoder, canvas }); ``` Audio is the same shape: `Audio.Source` into `Audio.Decoder` into @@ -144,7 +144,7 @@ The `` element automatically discovers the nested `` el - **WebCodecs decoding**: Hardware-accelerated video and audio decoding - **Reactive state**: All properties are signals from `@moq/signals` -- **Latency control**: A single target, or a range that buffers future-dated frames +- **Latency control**: A delay target plus optional buffering for future-dated frames - **Quality selection**: Switch between available renditions - **Custom tracks**: Unknown catalog sections pass through, and `broadcast.out.active` subscribes your own tracks diff --git a/js/watch/src/audio/decoder.ts b/js/watch/src/audio/decoder.ts index 2bd22adb13..54fa26490e 100644 --- a/js/watch/src/audio/decoder.ts +++ b/js/watch/src/audio/decoder.ts @@ -39,6 +39,14 @@ export type DecoderInput = { enabled: Getter; }; +/** Constructor properties for {@link Decoder}. */ +export type DecoderProps = Inputs & { + /** Rendition selector supplying encoded audio. */ + source: Source; + /** Shared playback clock. */ + sync: Sync; +}; + type DecoderOutput = { context: Signal; @@ -117,13 +125,14 @@ export class Decoder { // context, worklet, and ring alone. readonly #config: Computed; - constructor(source: Source, sync: Sync, props?: Inputs) { + constructor(props: DecoderProps) { this.in = { enabled: getter(props?.enabled ?? true), }; - this.source = source; - this.sync = sync; + this.source = props.source; + this.sync = props.sync; + this.#signals.cleanup(this.sync.register(this.source.out.jitter)); this.#identity = this.#signals.computed((effect) => { const config = effect.get(this.source.out.config); return config ? playbackIdentity(config) : undefined; @@ -219,6 +228,10 @@ export class Decoder { #runEnabled(effect: Effect): void { const enabled = effect.get(this.in.enabled); if (!enabled) return; + if (effect.get(this.sync.in.delay) === "instant") { + this.reset(); + return; + } const context = effect.get(this.#out.context); if (!context) return; @@ -251,10 +264,11 @@ export class Decoder { // RTT jitter, and debounce so a slider drag coalesces into one re-anchor. Decreases are left to // natural catch-up. #runLatencyReanchor(effect: Effect): void { + const delay = effect.get(this.sync.out.delay); + const jitter = effect.get(this.sync.out.jitter); const floor = reanchorFloor({ delay: effect.get(this.sync.in.delay), - audio: effect.get(this.sync.in.audio), - video: effect.get(this.sync.in.video), + media: Time.Milli.sub(delay, jitter), }); if (this.#prevFloor === undefined) { // Startup: the initial fill already builds the cushion; just record the baseline. @@ -273,6 +287,7 @@ export class Decoder { #runDecoder(effect: Effect): void { const enabled = effect.get(this.in.enabled); if (!enabled) return; + if (effect.get(this.sync.in.delay) === "instant") return; const broadcast = effect.get(this.source.in.broadcast); if (!broadcast) return; diff --git a/js/watch/src/audio/emitter.ts b/js/watch/src/audio/emitter.ts index 458399ea9e..b622ca2c3c 100644 --- a/js/watch/src/audio/emitter.ts +++ b/js/watch/src/audio/emitter.ts @@ -15,6 +15,12 @@ export type EmitterInput = { paused: Getter; }; +/** Constructor properties for {@link Emitter}. */ +export type EmitterProps = Inputs & { + /** Decoder supplying PCM. */ + source: Decoder; +}; + type EmitterOutput = { // Whether audio should be downloaded. Wired into the decoder's `enabled` input by the owner. enabled: Signal; @@ -36,8 +42,8 @@ export class Emitter { // The gain node used to adjust the volume. #gain = new Signal(undefined); - constructor(source: Decoder, props?: Inputs) { - this.source = source; + constructor(props: EmitterProps) { + this.source = props.source; this.in = { volume: getter(props?.volume ?? 0.5), muted: getter(props?.muted ?? false), diff --git a/js/watch/src/audio/latency.test.ts b/js/watch/src/audio/latency.test.ts index 39cade4ba0..ecc7c1f4d2 100644 --- a/js/watch/src/audio/latency.test.ts +++ b/js/watch/src/audio/latency.test.ts @@ -6,12 +6,12 @@ const ms = (value: number) => value as Time.Milli; describe("reanchorFloor", () => { it("includes the fixed delay and the largest media delay", () => { - expect(reanchorFloor({ delay: ms(100), audio: ms(20), video: ms(80) })).toBe(ms(180)); + expect(reanchorFloor({ delay: ms(100), media: ms(80) })).toBe(ms(180)); }); it("tracks rendition delay without adaptive RTT jitter", () => { - expect(reanchorFloor({ delay: "auto", audio: ms(20), video: ms(80) })).toBe(ms(80)); - expect(reanchorFloor({ delay: "auto", audio: ms(20), video: ms(200) })).toBe(ms(200)); + expect(reanchorFloor({ delay: "auto", media: ms(80) })).toBe(ms(80)); + expect(reanchorFloor({ delay: "auto", media: ms(200) })).toBe(ms(200)); }); }); diff --git a/js/watch/src/audio/latency.ts b/js/watch/src/audio/latency.ts index e45485bdc8..1766d2b8d5 100644 --- a/js/watch/src/audio/latency.ts +++ b/js/watch/src/audio/latency.ts @@ -6,11 +6,8 @@ export interface ReanchorFloor { /** How far playback trails the live edge. */ delay: Delay; - /** Additional delay required by the selected audio rendition. */ - audio?: Time.Milli; - - /** Additional delay required by the active or pending video rendition. */ - video?: Time.Milli; + /** Largest additional delay required by the registered media decoders. */ + media?: Time.Milli; } /** The stable delay floor whose increase requires the audio ring to refill. */ @@ -18,8 +15,7 @@ export function reanchorFloor(props: ReanchorFloor): Time.Milli { // "auto" and "instant" contribute nothing: the adaptive RTT component is deliberately excluded // so an RTT wiggle doesn't re-anchor, and "instant" holds nothing at all. const target = typeof props.delay === "number" ? props.delay : Time.Milli.zero; - const media = Time.Milli.max(props.audio ?? Time.Milli.zero, props.video ?? Time.Milli.zero); - return Time.Milli.add(target, media); + return Time.Milli.add(target, props.media ?? Time.Milli.zero); } // An AudioWorkletProcessor renders in fixed 128-sample quanta, so a ring shallower than one can diff --git a/js/watch/src/broadcast.test.ts b/js/watch/src/broadcast.test.ts index 4a90144068..81f76dc116 100644 --- a/js/watch/src/broadcast.test.ts +++ b/js/watch/src/broadcast.test.ts @@ -12,7 +12,7 @@ function publish(origin: Origin.Producer, path: Path.Valid) { } // A real origin with local broadcasts at the given paths. Resolution is proven by -// discrimination: `relativeBroadcast` resolves blind against the table (reload: false), so +// discrimination: `relativeBroadcast` resolves blind against the table (announced: false), so // a defined result means the reference resolved to a published path and nothing else. function origin(paths: string[]): Origin.Producer { const producer = new Origin.Producer(); @@ -26,7 +26,7 @@ function broadcast(name: string, paths: string[] = [name]): { source: Broadcast; origin: owner, name: Path.from(name), enabled: true, - reload: false, + announced: false, catalogFormat: "manual", }); return { source, owner }; @@ -52,6 +52,10 @@ const videoRenditions = (source: Broadcast): string[] => const video = (codec: string, broadcast?: string): Catalog.VideoConfig => ({ codec, container: { kind: "legacy" }, broadcast }) as Catalog.VideoConfig; +it("refuses the released reload input", () => { + expect(() => new Broadcast({ reload: true } as never)).toThrow("renamed to `announced`"); +}); + describe("relativeBroadcast", () => { it("resolves a legal reference against the origin", () => { const { source, owner } = broadcast("a/b", ["a/b", "a/source", "a/sub"]); @@ -96,7 +100,7 @@ describe("relativeBroadcast", () => { origin: owner, name: Path.from("a/b"), enabled: true, - reload: false, + announced: false, catalogFormat: "manual", catalog, }); @@ -193,7 +197,7 @@ describe("relativeBroadcast", () => { describe("blind resolution", () => { it("holds a resolved request steady instead of flapping", async () => { - // reload: false with nothing routed stands a request; when a session answers, the + // announced: false with nothing routed stands a request; when a session answers, the // effect that read `request.active` reruns. That rerun must re-acquire the same // answer, not close the request and re-dial forever. const owner = new Origin.Producer(); @@ -201,7 +205,7 @@ describe("blind resolution", () => { origin: owner, name: Path.from("blind.hang"), enabled: true, - reload: false, + announced: false, catalogFormat: "manual", }); diff --git a/js/watch/src/broadcast.ts b/js/watch/src/broadcast.ts index 5d448f7e16..15b5276713 100644 --- a/js/watch/src/broadcast.ts +++ b/js/watch/src/broadcast.ts @@ -78,11 +78,6 @@ function filterCatalog(catalog: Catalog.Root, usable: (rel: string | undefined) export const CATALOG_FORMATS = [...Catalog.FORMATS, "hangz", "manual"] as const; export type CatalogFormat = (typeof CATALOG_FORMATS)[number]; -export function parseCatalogFormat(value: string | null): CatalogFormat | undefined { - if (value === null) return undefined; - return CATALOG_FORMATS.find((f) => f === value); -} - type Status = "offline" | "loading" | "live"; // Signals the component reads. Whoever owns the backing Signal (the caller, or @@ -98,9 +93,9 @@ export type BroadcastInput = { // The broadcast name. name: Getter; - // Whether to reload the broadcast when it goes offline. + // Whether to wait for the broadcast to be announced before subscribing. // Defaults to true; pass false to subscribe immediately without waiting for an announcement. - reload: Getter; + announced: Getter; // Which catalog format to use. When `undefined` (the default), the format is // auto-detected from the broadcast name extension (`.hang`, `.msf`), falling @@ -124,7 +119,7 @@ type BroadcastOutput = { catalog: Signal; }; -// A catalog source that (optionally) reloads automatically when live/offline. +// A catalog source that can wait for announcement before subscribing. export class Broadcast { readonly in: Readonlys; @@ -153,14 +148,18 @@ export class Broadcast { // strands the filtered copy on the previous contents. readonly #raw = new Signal(undefined); - #signals = new Effect(); + #signals: Effect; constructor(props?: Inputs) { + if (props && "reload" in props) { + throw new Error("Watch.Broadcast: `reload` was renamed to `announced`"); + } + this.#signals = new Effect(); this.in = { origin: getter(props?.origin), name: getter(props?.name ?? Path.empty()), enabled: getter(props?.enabled ?? true), - reload: getter(props?.reload ?? true), + announced: getter(props?.announced ?? true), catalogFormat: getter(props?.catalogFormat), catalog: getter(props?.catalog), }; @@ -177,7 +176,7 @@ export class Broadcast { this.#announced.set(undefined); if (!effect.get(this.#wantAnnounced)) return; - if (!effect.get(this.in.reload)) return; + if (!effect.get(this.in.announced)) return; const origin = effect.get(this.in.origin); if (!origin) return; @@ -212,7 +211,7 @@ export class Broadcast { // Whether `path` is covered by an announced route, for `relativeBroadcast`'s // cross-broadcast refs. Announcements are prefix routes, so a route at "room/" covers // "room/alice/cam.hang" without naming it. Opens the announcement stream on first use. - // The blind cases (reload off, no discovery) never reach here; see `#relativeTarget`. + // The blind cases (announcement gate off, no discovery) never reach here; see `#relativeTarget`. #isPathAnnounced(effect: Effect, path: Moq.Path.Valid): boolean { this.#wantAnnounced.set(true); @@ -251,7 +250,7 @@ export class Broadcast { const name = effect.get(this.in.name); // No announcement gate: subscribe immediately. - if (!effect.get(this.in.reload)) { + if (!effect.get(this.in.announced)) { effect.set(this.#out.active, this.#requestBroadcast(effect, origin, name), undefined); return; } @@ -364,11 +363,11 @@ export class Broadcast { // keeps a reconnect from briefly hiding every cross-broadcast rendition from selection. if (!origin) return { local: false, path: resolved }; - // Without an announcement gate (reload off, or no session supports discovery), + // Without an announcement gate (disabled, or no session supports discovery), // resolve blind rather than waiting for an announcement that never comes. With the // gate, only report the path usable once it is announced: the request then resolves // from the table, never blind. - if (effect.get(this.in.reload) && effect.get(origin.discovery) !== false) { + if (effect.get(this.in.announced) && effect.get(origin.discovery) !== false) { if (!this.#isPathAnnounced(effect, resolved)) return undefined; } diff --git a/js/watch/src/element.ts b/js/watch/src/element.ts index bc4f14ca95..657a44153f 100644 --- a/js/watch/src/element.ts +++ b/js/watch/src/element.ts @@ -10,7 +10,7 @@ import type { Time } from "@moq/net"; import * as Moq from "@moq/net"; import { Effect, Signal } from "@moq/signals"; import * as Audio from "./audio"; -import { Broadcast, type CatalogFormat, parseCatalogFormat } from "./broadcast"; +import { Broadcast, CATALOG_FORMATS, type CatalogFormat } from "./broadcast"; import { formatDuration, parseDuration } from "./duration"; import { type Delay, Sync } from "./sync"; import * as Text from "./text"; @@ -23,12 +23,11 @@ const OBSERVED = [ "volume", "muted", "visible", - "reload", + "announced", "delay", "buffer", - // Released spellings, kept parsing but off the documented surface. `latency-max` is absent - // deliberately: its old ceiling included the floor, so translating it faithfully would mean - // tracking the resolved delay reactively, which is the coupling `buffer` exists to remove. + // Released spellings are observed only so assigning them can fail loudly instead of being ignored. + "reload", "latency", "latency-min", "jitter", @@ -71,32 +70,16 @@ function parseBuffer(value: string | null): Time.Milli { return Moq.Time.Milli.zero; } -// The released `latency` / `latencyMin` property spellings, translated onto `delay`. The range -// object is refused rather than half-applied: its ceiling included the floor, so it has no faithful -// `buffer` without tracking the resolved delay, which is the coupling `buffer` exists to remove. -function coerceLegacyDelay(value: unknown): Delay { - if (value === "instant") return "instant"; - if (value === undefined || value === "auto" || value === "real-time") return "auto"; - if (typeof value === "number" && Number.isFinite(value)) return Moq.Time.Milli(value); - throw new Error( - "moq-watch: the latency range is gone. Set `delay` (how far playback trails the live edge) and `buffer` (media held beyond it).", - ); -} - -// The released spellings of `delay`, in bare milliseconds. Kept unitless so pages still on them -// behave exactly as they did; `delay` is the current surface and does require a unit. -function parseLegacyDelay(value: string | null): Delay { - const trimmed = value?.trim(); - if (!trimmed || trimmed === "real-time") return "auto"; - if (trimmed === "instant") return "instant"; - const parsed = Number.parseFloat(trimmed); - return Moq.Time.Milli(Number.isFinite(parsed) ? parsed : 100); +/** Parse the element's catalog-format attribute. */ +export function parseCatalogFormat(value: string | null): CatalogFormat | undefined { + if (value === null) return undefined; + return CATALOG_FORMATS.find((format) => format === value); } /** * Parse a boolean attribute: absent uses `defaultValue`, bare presence is true, and an explicit * `"false"`/`"0"` is false. Presence alone can't express false, and attributes that default to - * true (`reload`) need to, so every boolean attribute accepts the explicit form. + * true (`announced`) need to, so every boolean attribute accepts the explicit form. */ function parseBoolean(value: string | null, defaultValue: boolean): boolean { if (value === null) return defaultValue; @@ -138,8 +121,8 @@ export default class MoqWatch extends HTMLElement { /** Selects the caption track. `text.out.available` lists the renditions for a picker. */ text: Text.Source; - /** Renders caption cues into an overlay above the canvas. */ - captionsRenderer: Text.Renderer; + /** Renders the selected text cues into an overlay above the canvas. */ + textRenderer: Text.Renderer; /** Keeps audio and video playing at the configured delay. */ sync: Sync; @@ -165,7 +148,7 @@ export default class MoqWatch extends HTMLElement { // Broadcast configuration owned here and wired into `broadcast` as inputs. #name = new Signal(Moq.Path.empty()); - #reload = new Signal(true); + #announced = new Signal(true); #catalogFormat = new Signal(undefined); #catalog = new Signal(undefined); @@ -210,7 +193,7 @@ export default class MoqWatch extends HTMLElement { origin: this.connection.origin, enabled: this.#enabled, name: this.#name, - reload: this.#reload, + announced: this.#announced, catalogFormat: this.#catalogFormat, catalog: this.#catalog, }); @@ -238,33 +221,28 @@ export default class MoqWatch extends HTMLElement { }); this.signals.cleanup(() => this.text.close()); - // The video decoder owns rendition handoffs but also needs Sync. Bridge its output through a - // parent-owned signal so Sync can be constructed first without exposing mutable wiring. - const videoJitter = new Signal(undefined); - this.sync = new Sync({ delay: this.controls.delay, buffer: this.controls.buffer, probe: this.connection.probe, - video: videoJitter, - audio: audioSource.out.jitter, }); this.signals.cleanup(() => this.sync.close()); - this.video = new Video.Decoder(videoSource, this.sync, { enabled: this.#videoEnabled }); - this.signals.proxy(videoJitter, this.video.out.jitter); - this.audio = new Audio.Decoder(audioSource, this.sync, { enabled: this.#audioEnabled }); + this.video = new Video.Decoder({ source: videoSource, sync: this.sync, enabled: this.#videoEnabled }); + this.audio = new Audio.Decoder({ source: audioSource, sync: this.sync, enabled: this.#audioEnabled }); this.signals.cleanup(() => { this.video.close(); this.audio.close(); }); - this.emitter = new Audio.Emitter(this.audio, { + this.emitter = new Audio.Emitter({ + source: this.audio, volume: this.controls.volume, muted: this.controls.muted, paused: this.controls.paused, }); - this.renderer = new Video.Renderer(this.video, { + this.renderer = new Video.Renderer({ + decoder: this.video, canvas: this.#canvas, visible: this.controls.visible, }); @@ -273,11 +251,13 @@ export default class MoqWatch extends HTMLElement { this.renderer.close(); }); - this.captionsRenderer = new Text.Renderer(this.text, this.sync, { + this.textRenderer = new Text.Renderer({ + source: this.text, + sync: this.sync, container: this.#captionsOverlay, enabled: this.#captionsEnabled, }); - this.signals.cleanup(() => this.captionsRenderer.close()); + this.signals.cleanup(() => this.textRenderer.close()); // Captions follow playback, like audio and video. The caption clock runs off wall time, so // leaving them on while paused scrolls text over a frozen frame. @@ -285,20 +265,10 @@ export default class MoqWatch extends HTMLElement { this.#captionsEnabled.set(effect.get(this.#enabled) && !effect.get(this.controls.paused)); }); - // Audio download follows the emitter's enable policy (paused/muted), except an instant - // delay turns it off outright: the ring needs a target depth to avoid underrunning, and - // unpaced video has nothing pulling it back toward the audio clock. + // Audio download follows the emitter's enable policy (paused/muted). The decoder itself + // refuses instant mode because an unpaced clock cannot keep an audio ring filled. this.signals.run((effect) => { - const enabled = effect.get(this.emitter.out.enabled); - this.#audioEnabled.set(enabled && effect.get(this.controls.delay) !== "instant"); - }); - - // Stopping the download leaves the ring holding a floor's worth of PCM, and the emitter - // stays connected to drain it. Flush on the way in so audio stops now instead of playing - // against video that just jumped to the live edge. - this.signals.run((effect) => { - if (effect.get(this.controls.delay) !== "instant") return; - this.audio.reset(); + this.#audioEnabled.set(effect.get(this.emitter.out.enabled)); }); // Video downloads while playing and on-screen. When paused, keep downloading only @@ -477,14 +447,16 @@ export default class MoqWatch extends HTMLElement { this.controls.muted.set(parseBoolean(newValue, false)); } else if (name === "visible") { this.controls.visible.set(parseVisible(newValue)); - } else if (name === "reload") { - this.#reload.set(parseBoolean(newValue, true)); + } else if (name === "announced") { + this.#announced.set(parseBoolean(newValue, true)); } else if (name === "delay") { this.controls.delay.set(parseDelay(newValue)); } else if (name === "buffer") { this.controls.buffer.set(parseBuffer(newValue)); + } else if (name === "reload") { + console.warn("moq-watch: `reload` was renamed to `announced`"); } else if (name === "latency" || name === "latency-min" || name === "jitter") { - this.controls.delay.set(parseLegacyDelay(newValue)); + console.warn(`moq-watch: \`${name}\` is gone; use \`delay\` and \`buffer\``); } else if (name === "catalog-format") { this.#catalogFormat.set(parseCatalogFormat(newValue)); } else if (name === "captions") { @@ -544,12 +516,17 @@ export default class MoqWatch extends HTMLElement { this.controls.visible.set(value); } - get reload(): boolean { - return this.#reload.peek(); + get announced(): boolean { + return this.#announced.peek(); + } + + set announced(value: boolean) { + this.#announced.set(value); } - set reload(value: boolean) { - this.#reload.set(value); + /** @internal */ + set reload(_value: unknown) { + throw new Error("moq-watch: `reload` was renamed to `announced`"); } /** @@ -587,8 +564,8 @@ export default class MoqWatch extends HTMLElement { return this.controls.delay.peek(); } - set latency(value: unknown) { - this.controls.delay.set(coerceLegacyDelay(value)); + set latency(_value: unknown) { + throw new Error("moq-watch: `latency` is gone; use `delay` and `buffer`"); } /** @internal */ @@ -596,8 +573,8 @@ export default class MoqWatch extends HTMLElement { return this.controls.delay.peek(); } - set latencyMin(value: unknown) { - this.controls.delay.set(coerceLegacyDelay(value)); + set latencyMin(_value: unknown) { + throw new Error("moq-watch: `latencyMin` is gone; use `delay` and `buffer`"); } /** @internal */ @@ -612,6 +589,10 @@ export default class MoqWatch extends HTMLElement { return this.sync.out.jitter.peek(); } + set jitter(_value: unknown) { + throw new Error("moq-watch: `jitter` is a readout; set `delay` instead"); + } + /** * Re-anchor playback at an utterance boundary in buffered mode: reset the sync reference * and flush the audio buffer so the next utterance plays from its own first frame. diff --git a/js/watch/src/enabled.test.ts b/js/watch/src/enabled.test.ts index ef465b61aa..d9fa97a9f0 100644 --- a/js/watch/src/enabled.test.ts +++ b/js/watch/src/enabled.test.ts @@ -20,9 +20,9 @@ function enabled(mode: Mode): boolean[] { const props = inputs(mode); const components = [ new Watch.Broadcast(props), - new Watch.Audio.Decoder(audioSource, sync, props), - new Watch.Video.Decoder(videoSource, sync, props), - new Watch.Text.Renderer(textSource, sync, props), + new Watch.Audio.Decoder({ source: audioSource, sync, ...props }), + new Watch.Video.Decoder({ source: videoSource, sync, ...props }), + new Watch.Text.Renderer({ source: textSource, sync, ...props }), ]; try { diff --git a/js/watch/src/sync.test.ts b/js/watch/src/sync.test.ts index ca8c8209b1..ed2411a8ec 100644 --- a/js/watch/src/sync.test.ts +++ b/js/watch/src/sync.test.ts @@ -62,4 +62,37 @@ describe("delay and buffer", () => { expect(sync.out.maxAge.peek()).toBe(0 as Time.Milli); sync.close(); }); + + it("includes registered decoder jitter until the decoder unregisters", async () => { + const media = new Signal(20 as Time.Milli); + const sync = new Sync({ delay: 100 as Time.Milli }); + const unregister = sync.register(media); + await flush(); + expect(sync.out.delay.peek()).toBe(120 as Time.Milli); + + media.set(80 as Time.Milli); + await flush(); + expect(sync.out.delay.peek()).toBe(180 as Time.Milli); + + unregister(); + await flush(); + expect(sync.out.delay.peek()).toBe(100 as Time.Milli); + sync.close(); + }); + + it("unregisters duplicate jitter inputs independently", async () => { + const media = new Signal(20 as Time.Milli); + const sync = new Sync({ delay: 100 as Time.Milli }); + const unregisterFirst = sync.register(media); + const unregisterSecond = sync.register(media); + + unregisterFirst(); + await flush(); + expect(sync.out.delay.peek()).toBe(120 as Time.Milli); + + unregisterSecond(); + await flush(); + expect(sync.out.delay.peek()).toBe(100 as Time.Milli); + sync.close(); + }); }); diff --git a/js/watch/src/sync.ts b/js/watch/src/sync.ts index f27c1fb361..07a5da3d12 100644 --- a/js/watch/src/sync.ts +++ b/js/watch/src/sync.ts @@ -1,6 +1,15 @@ import type * as Moq from "@moq/net"; import { Time } from "@moq/net"; -import { Effect, type Getter, getter, type Inputs, type Readonlys, readonlys, Signal } from "@moq/signals"; +import { + type Dispose, + Effect, + type Getter, + getter, + type Inputs, + type Readonlys, + readonlys, + Signal, +} from "@moq/signals"; /** * How far playback trails the live edge. @@ -37,12 +46,6 @@ export type SyncInput = { * from a `Connection`'s `probe`. */ probe: Getter; - - /** Any additional delay required for audio (wired from the per-rendition source). */ - audio: Getter; - - /** Any additional delay required for video (wired from the per-rendition source). */ - video: Getter; }; type SyncOutput = { @@ -94,6 +97,7 @@ export class Sync { // Minimum RTT seen, used as the baseline for jitter calculation. // Avoids inflating jitter due to bufferbloat. #minRtt: number | undefined; + #media = new Signal<{ jitter: Getter }[]>([]); #signals = new Effect(); @@ -102,8 +106,6 @@ export class Sync { delay: getter(props?.delay ?? ("auto" as Delay)), buffer: getter(props?.buffer ?? Time.Milli.zero), probe: getter(props?.probe), - audio: getter(props?.audio), - video: getter(props?.video), }; this.#update = Promise.withResolvers(); @@ -113,6 +115,13 @@ export class Sync { this.#signals.run(this.#runMaxAge.bind(this)); } + /** Include a decoder's rendition delay in the shared playback clock until disposed. */ + register(jitter: Getter): Dispose { + const registered = { jitter }; + this.#media.update((media) => [...media, registered]); + return () => this.#media.update((media) => media.filter((candidate) => candidate !== registered)); + } + // Derive `buffered` / `maxAge` from the resolved delay and the configured lookahead. #runMaxAge(effect: Effect): void { const delay = effect.get(this.#out.delay); @@ -159,13 +168,15 @@ export class Sync { #runDelay(effect: Effect): void { const jitter = effect.get(this.#out.jitter); - const video = effect.get(this.in.video) ?? Time.Milli.zero; - const audio = effect.get(this.in.audio) ?? Time.Milli.zero; + let media = Time.Milli.zero; + for (const registered of effect.get(this.#media)) { + media = Time.Milli.max(media, effect.get(registered.jitter) ?? Time.Milli.zero); + } // A zero delay still holds the rendition's own delay, which is a frame interval at 60fps. // "instant" holds nothing at all. const instant = effect.get(this.in.delay) === "instant"; - const delay = instant ? Time.Milli.zero : Time.Milli.add(Time.Milli.max(video, audio), jitter); + const delay = instant ? Time.Milli.zero : Time.Milli.add(media, jitter); this.#out.delay.set(delay); this.#update.resolve(); diff --git a/js/watch/src/text/renderer.ts b/js/watch/src/text/renderer.ts index 6669c178fd..bd7253c943 100644 --- a/js/watch/src/text/renderer.ts +++ b/js/watch/src/text/renderer.ts @@ -165,6 +165,14 @@ export type RendererInput = { enabled: Getter; }; +/** Constructor properties for {@link Renderer}. */ +export type RendererProps = Inputs & { + /** Caption rendition selector. */ + source: Source; + /** Shared playback clock. */ + sync: Sync; +}; + /** * Subscribes to the selected caption track, parses each cue, and renders it into an overlay element * via [media-captions](https://github.com/vidstack/captions). @@ -184,9 +192,9 @@ export class Renderer { // publisher costs one line instead of one per cue. #skewWarned = false; - constructor(source: Source, sync: Sync, props?: Inputs) { - this.source = source; - this.sync = sync; + constructor(props: RendererProps) { + this.source = props.source; + this.sync = props.sync; this.in = { container: getter(props?.container), enabled: getter(props?.enabled ?? true), diff --git a/js/watch/src/video/decoder.ts b/js/watch/src/video/decoder.ts index ebfa264e8f..d8bab4511a 100644 --- a/js/watch/src/video/decoder.ts +++ b/js/watch/src/video/decoder.ts @@ -36,6 +36,14 @@ export type DecoderInput = { enabled: Getter; }; +/** Constructor properties for {@link Decoder}. */ +export type DecoderProps = Inputs & { + /** Rendition selector supplying encoded video. */ + source: Source; + /** Shared playback clock. */ + sync: Sync; +}; + /** Cumulative video statistics since the decoder started. */ export interface Stats { /** Number of decoded frames. */ @@ -97,13 +105,14 @@ export class Decoder { this.#out.timestamp.set(undefined); } - constructor(source: Source, sync: Sync, props?: Inputs) { + constructor(props: DecoderProps) { this.in = { enabled: getter(props?.enabled ?? true), }; - this.source = source; - this.sync = sync; + this.source = props.source; + this.sync = props.sync; + this.#signals.cleanup(this.sync.register(this.out.jitter)); this.#identity = this.#signals.computed((effect) => { const config = effect.get(this.source.out.config); return config ? playbackIdentity(config) : undefined; diff --git a/js/watch/src/video/renderer.test.ts b/js/watch/src/video/renderer.test.ts index e93df7e509..5c54c82c87 100644 --- a/js/watch/src/video/renderer.test.ts +++ b/js/watch/src/video/renderer.test.ts @@ -88,7 +88,7 @@ describe("Renderer", () => { }, source: { out: { catalog } }, } as unknown as Decoder; - const renderer = new Renderer(decoder, { canvas, visible: "never" }); + const renderer = new Renderer({ decoder, canvas, visible: "never" }); try { await settle(); diff --git a/js/watch/src/video/renderer.ts b/js/watch/src/video/renderer.ts index 0846a582f4..d5120be11a 100644 --- a/js/watch/src/video/renderer.ts +++ b/js/watch/src/video/renderer.ts @@ -26,6 +26,12 @@ export type RendererInput = { visible: Getter; }; +/** Constructor properties for {@link Renderer}. */ +export type RendererProps = Inputs & { + /** Decoder supplying video frames. */ + decoder: Decoder; +}; + type RendererOutput = { // The most recently rendered frame, updated after each rAF paint. frame: Signal; @@ -54,8 +60,8 @@ export class Renderer { #ctx = new Signal(undefined); #signals = new Effect(); - constructor(decoder: Decoder, props?: Inputs) { - this.decoder = decoder; + constructor(props: RendererProps) { + this.decoder = props.decoder; this.in = { canvas: getter(props?.canvas), visible: getter(props?.visible ?? "20%"), diff --git a/js/watch/src/video/source.test.ts b/js/watch/src/video/source.test.ts index 0a4d9ae6ec..b268792ddd 100644 --- a/js/watch/src/video/source.test.ts +++ b/js/watch/src/video/source.test.ts @@ -200,7 +200,7 @@ describe("Source stalled rendition selection", () => { source.close(); }); - it("skips a stalled manual target while an unstalled rendition exists", async () => { + it("keeps a stalled manual target while an unstalled rendition exists", async () => { const source = new Source({ broadcast: mockBroadcast({ low: config("avc1.64001e", { bitrate: 1_000_000 }), @@ -211,7 +211,7 @@ describe("Source stalled rendition selection", () => { }); await settle(); - expect(source.out.track.peek()).toBe("low"); + expect(source.out.track.peek()).toBe("high"); expect(Object.keys(source.out.available.peek())).toEqual(["low", "high"]); source.close(); }); diff --git a/js/watch/src/video/source.ts b/js/watch/src/video/source.ts index 582dfe2d61..047060df11 100644 --- a/js/watch/src/video/source.ts +++ b/js/watch/src/video/source.ts @@ -338,19 +338,21 @@ export class Source { } #runSelected(effect: Effect): void { - const available = selectableRenditions(effect.get(this.#out.available)); - if (Object.keys(available).length === 0) return; - + const supported = effect.get(this.#out.available); const target = effect.get(this.in.target); - // Manual selection by name skips all ABR logic. - if (target?.name && target.name in available) { - const config = available[target.name]; + // A manual choice stays selected while stalled. `stalled` steers automatic adaptation; it + // must not silently override an explicit user selection. + if (target?.name && target.name in supported) { + const config = supported[target.name]; effect.set(this.#out.track, target.name); effect.set(this.#out.config, config); return; } + const available = selectableRenditions(supported); + if (Object.keys(available).length === 0) return; + // Auto-select: use recv bandwidth if no explicit bitrate target. let effectiveTarget = target; if (!target?.bitrate) { diff --git a/quest/m1/README.md b/quest/m1/README.md index f917cce56c..89831c2894 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -36,7 +36,6 @@ the transport line in m2 assumes a single stack. - [moq-tokio shapes](/quest/m1/api-tokio-shapes.md) - a `Drop` on `Listener`, a worker `Member` that cannot be cross-wired, `std::time::Duration` fields, one construction idiom, no six-argument merge - [Catalog types](/quest/m1/api-hang-catalog.md) - `hang::Catalog` is the one section list, `Clock` holds a `Timestamp`, `Timeline` folds into `Archive` - [Rendition ownership](/quest/m1/api-mux-rendition.md) - one handle publishes a media track and reports its estimate, instead of five -- [Watch and publish shapes](/quest/m1/api-watch-publish.md) - props objects everywhere, silent `latency`/`jitter` aliases refuse, `Sync` stops needing a jitter bridge, rooms get bandwidth - [Gateway types](/quest/m1/api-gateways.md) - no `anyhow` in a gateway `Error`, `PathOwned` prefixes, `Duration` segments, `moq_rtc::Server::new(config)`, an SRT reject with a reason - [libmoq units](/quest/m1/api-libmoq-units.md) - `moq_client_config` is all microseconds, the header declares every enum and error code, NULL callbacks are refused - [Cluster -01](/quest/m1/cluster-01/README.md) - rs/moq-net and js/net speak the revised cluster extension (HOP_ID, REQUEST_UPDATE repricing) and -01 is published diff --git a/quest/m1/api-review-gate.md b/quest/m1/api-review-gate.md index 27cac419e9..fdc40a3872 100644 --- a/quest/m1/api-review-gate.md +++ b/quest/m1/api-review-gate.md @@ -22,7 +22,6 @@ The list: [Announce event](/quest/m1/api-net-announce.md), [moq-tokio shapes](/quest/m1/api-tokio-shapes.md), [Catalog types](/quest/m1/api-hang-catalog.md), [Rendition ownership](/quest/m1/api-mux-rendition.md), -[Watch and publish shapes](/quest/m1/api-watch-publish.md), [Gateway types](/quest/m1/api-gateways.md), [libmoq units](/quest/m1/api-libmoq-units.md). diff --git a/quest/m1/api-watch-publish.md b/quest/m1/api-watch-publish.md deleted file mode 100644 index dba467c3b4..0000000000 --- a/quest/m1/api-watch-publish.md +++ /dev/null @@ -1,61 +0,0 @@ -# [M] Watch and publish take one shape and refuse old spellings - -## Goal - -`@moq/watch`, `@moq/publish`, and `@moq/room` construct every component the -same way, mean one thing by each name, and refuse the attributes and props -the last release spelled differently instead of aliasing them silently. - -## Plan - -- `Sync` stops taking `video`/`audio` jitter inputs; the decoders already - hold `sync` and register their jitter. Today `element.ts`, room's - `Member`, and moq.pro each bridge a `Signal` between `new Sync({ video })` - and `new Video.Decoder(source, sync)` in four lines. -- Delete the silent aliases: the `latency`, `latency-min`, and `jitter` - attributes and the `latency`/`latencyMin` props in `js/watch/src/element.ts` - write `delay`/`buffer` today while only `latencyMax` throws, and the - `jitter` attribute is a knob where the `jitter` getter is a readout. - Setters throw like `latencyMax`; `get jitter()` stays a readout. -- Delete `Audio.Encoder.muted` in `js/publish`; `volume = 0` covers it and - `` keeps its capture-gating meaning like - ``. -- Watch constructors take a props object like publish: - `new Video.Decoder({ source, sync, enabled })`, - `new Video.Renderer({ decoder, canvas, visible })`, same for audio and - text; positional arguments stay for identity only. -- `room.Local({ connection })` like `Room({ connection })`, wiring - `connection.bandwidth` into its six encoders; today `Local` takes only - `origin`, so no room publisher gets a reservation and the send-estimate - split is silently off. `LocalProps` uses `GetterInit` like `RoomProps`. -- `Publish.Broadcast.in.maxAge` is `Getter`, not a - bare `number`. -- `room.Local` exports the sources, captures, and broadcasts the demo - touches, not seventeen readonly fields; `parseCatalogFormat` moves from - the watch index to the element. -- Every publish source exposes `out.source: { video?, audio? } | undefined`; - camera, screen, and file each have a different shape today and the - element unwraps three ways. -- ``'s `text` field is a source where `video`/`audio` are - decoders; name the pair `text`/`textRenderer` (or `captions`/ - `captionsRenderer`). `` wraps `sources` in `readonlys()` and - drops `el.capture`, a duplicate of `el.video.capture`. -- `reload` becomes `announced` on `Broadcast`, the element attribute, - room's `Member`, and the docs; it gates on the announcement and moq.pro - reads the old name as retry. The released spelling refuses. `stalled` - stays on both the catalog (publisher lagging) and the watch decoders - (viewer waiting); the maintainer likes one word for both sides. -- Bugs on the same files: a stalled rendition silently overrides a manual - quality pick (`video/source.ts`); `delay = "instant"` underruns audio - unless the owner remembers to disable it, so `Audio.Decoder` gates it. -- Docs: `js/watch/README.md` and `doc/lib/js/watch.md` still list - `latency*` and `jitter` attributes; `doc/lib/js/publish.md` uses - `latencyMax` where `Track.Info` has `maxAge`. - -Public API: breaking on @moq/watch, @moq/publish, and @moq/room, so on -dev. Wire: none. Consumers: `demo/web`, `doc/lib/js`, moq.pro's app. - -## Related - -- [Headless player](/quest/m2/watch-player.md) - the pipeline assembly every owner repeats, extracted after this settles the pieces -- [Audio time stretch](/quest/m2/watch-audio-time-stretch.md) - the audio ring behaviour this leaves alone diff --git a/quest/m2/watch-player.md b/quest/m2/watch-player.md index 2cc941bb97..141de48bbc 100644 --- a/quest/m2/watch-player.md +++ b/quest/m2/watch-player.md @@ -22,7 +22,6 @@ Public API: additive on @moq/watch and @moq/room. Wire: none. ## Required -- [Watch and publish shapes](/quest/m1/api-watch-publish.md) - the pieces settle before the assembly is named - [Merge dev](/quest/m1/merge-dev.md) - starts on main ## Related From 1b9715f37ce25778f80c40d096154aa0f618c955 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sun, 20 Sep 2026 13:35:26 -0700 Subject: [PATCH 3/3] fix(publish): clear inactive source state Co-Authored-By: OpenAI Codex --- js/publish/src/element.ts | 12 ++++- js/publish/src/source-state.test.ts | 71 +++++++++++++++++++++++++++++ js/publish/src/source-state.ts | 24 ++++++++++ js/publish/src/source/file.test.ts | 43 +++++++++++++++++ js/publish/src/source/file.ts | 4 +- 5 files changed, 151 insertions(+), 3 deletions(-) create mode 100644 js/publish/src/source-state.test.ts create mode 100644 js/publish/src/source-state.ts create mode 100644 js/publish/src/source/file.test.ts diff --git a/js/publish/src/element.ts b/js/publish/src/element.ts index 28df9edac9..a9ed24e869 100644 --- a/js/publish/src/element.ts +++ b/js/publish/src/element.ts @@ -11,6 +11,7 @@ import * as Audio from "./audio"; import { Broadcast } from "./broadcast"; import * as Preview from "./preview"; import * as Source from "./source"; +import { clearSourceState } from "./source-state"; import * as Video from "./video"; const OBSERVED = ["url", "name", "muted", "invisible", "source", "preview", "announce"] as const; @@ -108,7 +109,7 @@ export default class MoqPublish extends HTMLElement { readonly sources = readonlys(this.#sources); // The captured media tracks, written by #runSource. Fed to the video and audio captures, so - // consumers read them back via `capture.in.source` / `audio.capture.in.source` rather than here. + // consumers read them back via `video.capture` / `audio.capture` rather than here. #videoSource = new Signal(undefined); #audioSource = new Signal(undefined); @@ -310,6 +311,15 @@ export default class MoqPublish extends HTMLElement { #runSource(effect: Effect) { const source = effect.get(this.controls.source); + + // Every selection owns the complete source state. Explicitly clear every inactive slot so a + // switch cannot leave a holder or captured track from the previous selection observable. + clearSourceState(effect, { + holders: this.#sources, + video: this.#videoSource, + audio: this.#audioSource, + }); + if (!source) return; if (source === "camera") { diff --git a/js/publish/src/source-state.test.ts b/js/publish/src/source-state.test.ts new file mode 100644 index 0000000000..8304565df2 --- /dev/null +++ b/js/publish/src/source-state.test.ts @@ -0,0 +1,71 @@ +import { expect, test } from "bun:test"; +import { Effect, Signal } from "@moq/signals"; +import type * as Audio from "./audio"; +import type * as Source from "./source"; +import { clearSourceState, type SourceState } from "./source-state"; +import type * as Video from "./video"; + +const flush = () => new Promise((resolve) => queueMicrotask(resolve)); + +async function settle(times = 5): Promise { + for (let i = 0; i < times; i++) await flush(); +} + +test("switching and clearing selections removes every inactive source", async () => { + const selected = new Signal<"camera" | "file" | undefined>("camera"); + const camera = { kind: "camera" } as unknown as Source.Camera; + const microphone = { kind: "microphone" } as unknown as Source.Microphone; + const file = { kind: "file" } as unknown as Source.File; + const video = { kind: "video" } as unknown as Video.Source; + const audio = { kind: "audio" } as unknown as Audio.Source; + + const state: SourceState = { + holders: { + video: new Signal(undefined), + audio: new Signal(undefined), + file: new Signal(undefined), + }, + video: new Signal(undefined), + audio: new Signal(undefined), + }; + + const effect = new Effect((effect) => { + const source = effect.get(selected); + clearSourceState(effect, state); + + if (source === "camera") { + effect.set(state.holders.video, camera); + effect.set(state.holders.audio, microphone); + effect.set(state.video, video); + effect.set(state.audio, audio); + } else if (source === "file") { + effect.set(state.holders.file, file); + } + }); + + try { + await settle(); + expect(state.holders.video.peek()).toBe(camera); + expect(state.holders.audio.peek()).toBe(microphone); + expect(state.video.peek()).toBe(video); + expect(state.audio.peek()).toBe(audio); + + selected.set("file"); + await settle(); + expect(state.holders.video.peek()).toBeUndefined(); + expect(state.holders.audio.peek()).toBeUndefined(); + expect(state.holders.file.peek()).toBe(file); + expect(state.video.peek()).toBeUndefined(); + expect(state.audio.peek()).toBeUndefined(); + + selected.set(undefined); + await settle(); + expect(state.holders.video.peek()).toBeUndefined(); + expect(state.holders.audio.peek()).toBeUndefined(); + expect(state.holders.file.peek()).toBeUndefined(); + expect(state.video.peek()).toBeUndefined(); + expect(state.audio.peek()).toBeUndefined(); + } finally { + effect.close(); + } +}); diff --git a/js/publish/src/source-state.ts b/js/publish/src/source-state.ts new file mode 100644 index 0000000000..fe8872ea9e --- /dev/null +++ b/js/publish/src/source-state.ts @@ -0,0 +1,24 @@ +import type { Effect, Signal } from "@moq/signals"; +import type * as Audio from "./audio"; +import type * as Source from "./source"; +import type * as Video from "./video"; + +/** The mutable source state owned by the publish element. */ +export interface SourceState { + holders: { + video: Signal; + audio: Signal; + file: Signal; + }; + video: Signal; + audio: Signal; +} + +/** Clears every holder and captured source for the lifetime of this effect run. */ +export function clearSourceState(effect: Effect, state: SourceState): void { + effect.set(state.holders.video, undefined); + effect.set(state.holders.audio, undefined); + effect.set(state.holders.file, undefined); + effect.set(state.video, undefined); + effect.set(state.audio, undefined); +} diff --git a/js/publish/src/source/file.test.ts b/js/publish/src/source/file.test.ts new file mode 100644 index 0000000000..412fc56b10 --- /dev/null +++ b/js/publish/src/source/file.test.ts @@ -0,0 +1,43 @@ +import { expect, test } from "bun:test"; +import { File as FileSource } from "./file"; + +const flush = () => new Promise((resolve) => queueMicrotask(resolve)); + +async function settle(times = 5): Promise { + for (let i = 0; i < times; i++) await flush(); +} + +test("clearing a decoded file clears its published media", async () => { + const createImageBitmap = Object.getOwnPropertyDescriptor(globalThis, "createImageBitmap"); + const videoFrame = Object.getOwnPropertyDescriptor(globalThis, "VideoFrame"); + + Object.defineProperty(globalThis, "createImageBitmap", { + configurable: true, + value: async () => ({ close() {} }), + }); + Object.defineProperty(globalThis, "VideoFrame", { + configurable: true, + value: class { + close() {} + }, + }); + + const source = new FileSource({ + file: new File([new Uint8Array([0])], "still.png", { type: "image/png" }), + }); + + try { + await settle(); + expect(source.out.source.peek()?.video).toBeDefined(); + + source.file.set(undefined); + await settle(); + expect(source.out.source.peek()).toBeUndefined(); + } finally { + source.close(); + if (createImageBitmap) Object.defineProperty(globalThis, "createImageBitmap", createImageBitmap); + else Reflect.deleteProperty(globalThis, "createImageBitmap"); + if (videoFrame) Object.defineProperty(globalThis, "VideoFrame", videoFrame); + else Reflect.deleteProperty(globalThis, "VideoFrame"); + } +}); diff --git a/js/publish/src/source/file.ts b/js/publish/src/source/file.ts index 6b2672c763..7ef00e4747 100644 --- a/js/publish/src/source/file.ts +++ b/js/publish/src/source/file.ts @@ -133,7 +133,7 @@ export class File { }, }); - effect.set(this.#out.source, { video: { frames, frameRate: IMAGE_FRAME_RATE } }, {}); + effect.set(this.#out.source, { video: { frames, frameRate: IMAGE_FRAME_RATE } }, undefined); } async #decodeMedia(file: globalThis.File, effect: Effect) { @@ -181,7 +181,7 @@ export class File { }; if (signal.aborted) return; - effect.set(this.#out.source, source, {}); + effect.set(this.#out.source, source, undefined); } /** Stop decoding and release the file. */