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..a9ed24e869 100644
--- a/js/publish/src/element.ts
+++ b/js/publish/src/element.ts
@@ -6,11 +6,12 @@
* @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";
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;
@@ -89,8 +90,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,14 +101,15 @@ 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.
+ // consumers read them back via `video.capture` / `audio.capture` rather than here.
#videoSource = new Signal(undefined);
#audioSource = new Signal(undefined);
@@ -174,8 +175,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 +190,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 +228,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,
@@ -310,23 +311,27 @@ 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") {
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 +346,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 +375,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-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/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.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 4183b6b457..7ef00e4747 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);
@@ -132,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) {
@@ -180,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. */
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