Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion demo/web/src/meet.ts
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ function join(): void {
enabled: true,
});
local = new Local({
origin: connection.origin,
connection,
identity,
enabled: true,
user: { id: name, name },
Expand Down
3 changes: 2 additions & 1 deletion demo/web/src/publish.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
8 changes: 5 additions & 3 deletions doc/lib/js/publish.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<Meta>({ track });
Expand Down Expand Up @@ -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 });
```

Expand Down
2 changes: 1 addition & 1 deletion doc/lib/js/room.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ const connection = new Connection({
});

const local = new Local({
origin: connection.origin,
connection,
identity: Path.from("alice"),
user: { name: "Alice" },
});
Expand Down
6 changes: 3 additions & 3 deletions doc/lib/js/watch.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.

Expand Down
23 changes: 13 additions & 10 deletions js/moq-boy/src/game.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Moq.Time.Milli | undefined>(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.
Expand All @@ -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,
Expand Down
6 changes: 4 additions & 2 deletions js/publish/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 } });
Expand All @@ -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);
```
Expand Down
10 changes: 3 additions & 7 deletions js/publish/src/audio/encoder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,6 @@ export type EncoderInput = {
/** Constructor options: the wired inputs plus the live-editable tuning knobs. */
export type EncoderProps = Inputs<EncoderInput> & {
// User tuning knobs. Seed a value or wire a Signal; also live-editable via the matching field.
muted?: boolean | Signal<boolean>;
volume?: number | Signal<number>;

// Codec selection plus encoder settings. Defaults to "opus".
Expand Down Expand Up @@ -123,8 +122,6 @@ export class Encoder {

readonly in: Readonlys<EncoderInput>;

/** Silence the encoded audio without tearing down the capture graph. */
muted: Signal<boolean>;
/** Linear gain applied before encoding, where 1 is unity. */
volume: Signal<number>;
/** The live-editable codec selection plus its encoder settings. */
Expand Down Expand Up @@ -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<Codec>(props?.codec ?? "opus");

Expand All @@ -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;
Expand All @@ -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 (;;) {
Expand All @@ -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
Expand Down
73 changes: 38 additions & 35 deletions js/publish/src/element.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand All @@ -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<Source.Camera | Source.Screen | undefined>(undefined),
audio: new Signal<Source.Microphone | Source.Screen | undefined>(undefined),
file: new Signal<Source.File | undefined>(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<Video.Source | undefined>(undefined);
#audioSource = new Signal<Audio.Source | undefined>(undefined);

Expand Down Expand Up @@ -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({
Expand All @@ -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,
});
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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();
Expand All @@ -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();
Expand All @@ -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(() => {
Expand Down
Loading
Loading