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/publish.ts
Original file line number Diff line number Diff line change
Expand Up @@ -404,7 +404,7 @@ meta.run((effect) => {
if (!net) return;

// A day-long cache so a viewer joining long after the last edit still replays the value.
const track = net.createTrack(META_TRACK, { maxAge: 86_400_000 });
const track = net.createTrack(META_TRACK, { maxAge: Net.Time.Milli(86_400_000) });
effect.cleanup(() => track.close());

const producer = new Json.Snapshot.Producer<unknown>({ track });
Expand Down
2 changes: 1 addition & 1 deletion demo/web/src/stats.ts
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ function subscribeNode(effect: Signals.Effect, origin: Net.Origin.Table, path: N
if (!consumer) return;

const sub = <K extends keyof NodeStats>(trackName: string, key: K) => {
const track = consumer.subscribe(trackName).ordered();
const track = consumer.track(trackName).subscribe().ordered();
effect.cleanup(() => track.close());
effect.spawn(async () => {
for (;;) {
Expand Down
20 changes: 15 additions & 5 deletions doc/lib/js/net.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,11 @@ const url = new URL("https://cdn.moq.dev/anon?jwt=...");
// Publish. The origin is the routing table the connection announces and serves,
// so a broadcast survives a reconnect.
const origin = new Moq.Origin.Producer();
const connection = await Moq.Connection.connect(url, { publish: origin.consume() });
const connection = await Moq.Connection.connect({
url,
publish: origin.consume(),
consume: origin,
});

const broadcast = origin.createBroadcast(Moq.Path.from("chat.room"));
const track = broadcast.createTrack("messages");
Expand All @@ -29,7 +33,13 @@ group.close();
broadcast.announce();

// Subscribe
const consumer = connection.consume(Moq.Path.from("chat.room")).track("messages").subscribe({ priority: 0 });
const request = origin.request(Moq.Path.from("chat.room"));
let active = request.active.peek();
while (!active) {
await request.active.changed();
active = request.active.peek();
}
const consumer = active.track("messages").subscribe({ priority: 0 });
for (;;) {
const group = await consumer.recvGroup();
if (!group) break;
Expand All @@ -41,9 +51,9 @@ for (;;) {
- **Connections** race WebTransport against WebSocket. `new Connection({ url })` pools one connection per relay URL and reconnects with backoff, which the elements use. `closed` settles when the handle is released (`null` on a clean close); the failure that stopped retrying the current URL is `error`, and a new URL recovers the same handle. A connection owns one send-rate sampler and one `Bandwidth.Allocator`; publishers reserve against it so their encoder targets sum to the estimate instead of each matching it.
- **Bandwidth** (`Bandwidth.Allocator`) divides the connection's send-rate estimate by track priority, max-min fair within a tier. An idle track claims nothing. The receive side is untouched.
- **Discovery** by any pattern scope (`origin.announced(scope)`, such as `room/*/chat`; default everything). Each event's `path` is the covered prefix relative to the origin, `captures` reports what the scope's wildcards matched when the prefix pins them, and `kind` says whether it was announced, updated, or retracted. The consumer is an async iterable. `origin.dynamic(prefix, route)` advertises a prefix.
- **Subscriptions** carry a priority and max age; groups arrive out of order and are read frame by frame, with `Lagged` when a reader asks for a frame the group never held and `GroupTooLarge` when a write exceeds the cache budget and aborts the group.
- **Subscriptions** carry a priority, a `Time.Milli` max age, and optional `groups` bounds. Groups arrive out of order and are read frame by frame, with `Error.TooFarBehind` when a reader asks for a frame the group never held and `Error.GroupTooLarge` when a write exceeds the cache budget and aborts the group.
- **Datagrams** on moq-lite 05+ and fetch-by-sequence for history.
- **Errors** split by scope: a stream reset throws `StreamError` with a `StreamCode`, a session close gives `SessionError` with a `SessionCode`. The registries are disjoint, so the same number means different things in each, and 64+ is yours. Same on either transport. Named conditions like `Lagged` subclass `StreamError`, so one `code` check catches a gap whether it happened here or at the peer, and resetting a moq-lite stream with one sends that code rather than a bare internal error. IETF streams use their own mapping: cancellation sends CANCELLED, other local failures send INTERNAL\_ERROR, and received codes remain opaque.
- **Errors** live under one namespace: a stream reset throws `Error.Stream` with a `StreamCode`, while a session close gives `Error.Session` with a `SessionCode`. The registries are disjoint, so the same number means different things in each, and 64+ is yours. Named conditions such as `Error.TooFarBehind`, `Error.FrameTooLarge`, and `Error.GroupTooLarge` subclass `Error.Stream`, so one `code` check handles a condition raised here or reported by the peer. IETF streams use their own mapping: cancellation sends CANCELLED, other local failures send INTERNAL\_ERROR, and received codes remain opaque.
- **Paths** with `Path.relative` for the cross-broadcast catalog references hang uses. Path patterns (`Path.Pattern`, `Path.Patterns`) are re-exported from [`@moq/pattern`](https://www.npmjs.com/package/@moq/pattern). Literal `Path` stays a coordinate.

The [path pattern](/concept/moq-lite#path-patterns) grammar lives on the
Expand Down Expand Up @@ -83,7 +93,7 @@ Three operations, on an origin:
- `origin.dynamic(prefix, route)` claims `prefix` and every path beneath it
(`""` claims everything). Hold the returned `Origin.Dynamic` while the
claim should stay advertised; `close()` retracts it. A request beneath it
with no local broadcast is a `BroadcastRequest` to `accept` or `reject`;
with no local broadcast is an `Origin.Request` to `accept` or `reject`;
reject what you will not serve rather than narrowing the claim, since a
route is always a prefix on every wire.

Expand Down
2 changes: 1 addition & 1 deletion js/binary/src/snapshot/consumer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ export class Consumer {
// Anything else is the track's terminal error, which every later read would throw
// again; swallowing it would spin here instead of telling the caller the
// subscription died.
if (!(err instanceof Moq.Group.Lagged)) throw err;
if (!(err instanceof Moq.Error.TooFarBehind)) throw err;
continue;
}

Expand Down
2 changes: 1 addition & 1 deletion js/binary/src/snapshot/producer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ export class Producer {
// Check before opening a group. `appendGroup` publishes immediately, so letting `writeFrame`
// reject the frame would leave an empty newest group behind: a snapshot consumer jumps to the
// newest, so the previous value would be lost even though this update threw.
if (encoded.byteLength > Moq.Group.MAX_GROUP_CACHE_BYTES) throw new Moq.Group.FrameTooLarge();
if (encoded.byteLength > Moq.Group.MAX_GROUP_CACHE_BYTES) throw new Moq.Error.FrameTooLarge();

const group = this.#track.appendGroup();
try {
Expand Down
2 changes: 1 addition & 1 deletion js/binary/src/snapshot/snapshot.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ const bytes = (...values: number[]) => new Uint8Array(values);

// Walking a finished track's groups inspects a complete timeline, so request a replay window
// instead of the transport's live-edge default, which skips every superseded group.
const REPLAY_LATENCY = 30_000;
const REPLAY_LATENCY = Time.Milli(30_000);

// Drain every value currently available from a fresh consumer over the (finished) track.
async function drain(track: Track.Subscriber, compression: boolean): Promise<Uint8Array[]> {
Expand Down
2 changes: 1 addition & 1 deletion js/binary/src/stream/stream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import { Consumer, Producer, Rolled } from "./index.ts";

// Ask for a replay window, so the superseded first group is delivered rather than skipped by the
// subscriber's default max-age budget. A rolled log is exactly the case where both groups matter.
const REPLAY_LATENCY = 30_000;
const REPLAY_LATENCY = Time.Milli(30_000);

const payloads = (count: number) => Array.from({ length: count }, (_, n) => new Uint8Array(8).fill(n));

Expand Down
30 changes: 14 additions & 16 deletions js/clock/src/main.ts
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ ENVIRONMENT VARIABLES:
async function publish(config: Config) {
// The origin holds what we publish; the connection announces and serves it.
const origin = new Moq.Origin.Producer();
await Moq.Connection.connect(new URL(config.url), { publish: origin.consume() });
const connection = await Moq.Connection.connect({ url: new URL(config.url), publish: origin.consume() });
console.log("✅ Connected to relay:", config.url);

// Create a new "broadcast", which is a collection of tracks.
Expand All @@ -83,19 +83,8 @@ async function publish(config: Config) {

console.log("✅ Published broadcast:", config.broadcast);

// Wait until we get a subscription for the track
for (;;) {
const request = await broadcast.requested();
if (!request) break;

if (request.name === config.track) {
// Accept to commit the track's immutable properties (so a lite-05 TRACK
// request resolves) and obtain the Track to produce into.
void publishTrack(request.accept());
} else {
request.reject(new Error("not found"));
}
}
void publishTrack(broadcast.createTrack(config.track));
await connection.closed;
}

async function publishTrack(track: Moq.Track.Producer) {
Expand Down Expand Up @@ -141,10 +130,16 @@ async function publishTrack(track: Moq.Track.Producer) {
}

async function subscribe(config: Config) {
const connection = await Moq.Connection.connect(new URL(config.url));
const origin = new Moq.Origin.Producer();
const connection = await Moq.Connection.connect({ url: new URL(config.url), consume: origin });
console.log("✅ Connected to relay:", config.url);

const broadcast = connection.consume(Moq.Path.from(config.broadcast));
const request = origin.request(Moq.Path.from(config.broadcast));
let broadcast = request.active.peek();
while (!broadcast) {
await request.active.changed();
broadcast = request.active.peek();
}
const track = broadcast.track(config.track).subscribe({ priority: 0 });

console.log("✅ Subscribed to track:", config.track);
Expand Down Expand Up @@ -185,6 +180,9 @@ async function subscribe(config: Config) {
console.log(clockEmoji, base + seconds);
}
}

connection.close();
origin.close();
}

// Wait for the WebTransport polyfill to be ready
Expand Down
8 changes: 4 additions & 4 deletions js/hang/src/container/consumer.outoforder.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,8 +47,8 @@ async function drain(consumer: Consumer): Promise<[number, number | undefined][]
// audio writes into a timestamp-indexed ring and video drops a late frame at render, and the
// subscription's own max age already bounds how far back one can be.
test("out-of-order groups are delivered rather than dropped", async () => {
const track = new Track.Producer("test").accept({ maxAge: 30_000 });
const consumer = new Consumer(track.subscribe({ maxAge: 5000 }), {
const track = new Track.Producer("test").accept({ maxAge: Time.Milli(30_000) });
const consumer = new Consumer(track.subscribe({ maxAge: Time.Milli(5000) }), {
format: new LegacyFormat("data"),
maxAge: 5000 as Time.Milli,
});
Expand Down Expand Up @@ -77,8 +77,8 @@ test("out-of-order groups are delivered rather than dropped", async () => {
// drained (the decode loop consumes faster than the network delivers). Removing it at that
// instant silently truncates its tail, so removal must wait for the group to finish.
test("a below-cursor group still downloading is not truncated", async () => {
const track = new Track.Producer("test").accept({ maxAge: 30_000 });
const consumer = new Consumer(track.subscribe({ maxAge: 5000 }), {
const track = new Track.Producer("test").accept({ maxAge: Time.Milli(30_000) });
const consumer = new Consumer(track.subscribe({ maxAge: Time.Milli(5000) }), {
format: new LegacyFormat("data"),
maxAge: 5000 as Time.Milli,
});
Expand Down
28 changes: 16 additions & 12 deletions js/hang/src/container/consumer.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { expect, spyOn, test } from "bun:test";
import { Format as LocFormat, Producer as LocProducer } from "@moq/loc";
import { Group, SessionCode, SessionError, StreamCode, StreamError, Time, Track, Varint } from "@moq/net";
import { Group, Error as NetError, SessionCode, StreamCode, Time, Track, Varint } from "@moq/net";
import { AudioConfigSchema } from "../catalog/audio.ts";
import { decodeInitSegment, type InitSegment } from "./cmaf/decode.ts";
import { createAudioInitSegment, encodeDataSegment } from "./cmaf/encode.ts";
Expand Down Expand Up @@ -50,7 +50,7 @@ function settle(ms = 20): Promise<void> {
// These tests write every group up front and only then read, so they ask for history
// rather than the live edge.
function replay(track: Track.Producer): Track.Subscriber {
return track.subscribe({ maxAge: 30_000 });
return track.subscribe({ maxAge: Time.Milli(30_000) });
}

// --- LegacyFormat ---
Expand Down Expand Up @@ -130,7 +130,7 @@ test("Legacy Producer accepts open-GOP leading pictures above the previous group

test("Legacy Producer writes a duration marker at the next keyframe", async () => {
const track = new Track.Producer("test");
const subscriber = track.subscribe({ maxAge: 30_000 });
const subscriber = track.subscribe({ maxAge: Time.Milli(30_000) });
const producer = new LegacyProducer(track, new LegacyFormat("video"));
producer.encode(new Uint8Array([0xde, 0xad]), 0 as Time.Micro, true);
producer.encode(new Uint8Array([0xbe, 0xef]), 10_000 as Time.Micro, false);
Expand All @@ -155,7 +155,7 @@ test("Legacy Producer writes a duration marker at the next keyframe", async () =

test("Legacy Producer omits a reordered group's presentation endpoint marker", async () => {
const track = new Track.Producer("test");
const subscriber = track.subscribe({ maxAge: 30_000 });
const subscriber = track.subscribe({ maxAge: Time.Milli(30_000) });
const producer = new LegacyProducer(track, new LegacyFormat("video"));
for (const [index, timestamp] of [0, 120_000, 40_000, 80_000].entries()) {
producer.encode(new Uint8Array([1]), timestamp as Time.Micro, index === 0);
Expand All @@ -180,7 +180,7 @@ test("Legacy Producer omits a reordered group's presentation endpoint marker", a

test("Legacy Producer estimates the tail from the current cadence", async () => {
const track = new Track.Producer("test");
const subscriber = track.subscribe({ maxAge: 30_000 });
const subscriber = track.subscribe({ maxAge: Time.Milli(30_000) });
const producer = new LegacyProducer(track, new LegacyFormat("video"));
for (const [index, timestamp] of [0, 16_000, 32_000, 65_000, 98_000].entries()) {
producer.encode(new Uint8Array([1]), timestamp as Time.Micro, index === 0);
Expand All @@ -200,7 +200,7 @@ test("Legacy Producer estimates the tail from the current cadence", async () =>

test("Legacy Producer rejects a backwards cut without closing the group", async () => {
const track = new Track.Producer("test");
const subscriber = track.subscribe({ maxAge: 30_000 });
const subscriber = track.subscribe({ maxAge: Time.Milli(30_000) });
const producer = new LegacyProducer(track, new LegacyFormat("video"));
producer.encode(new Uint8Array([1]), 20_000 as Time.Micro, true);
expect(() => producer.cut(10_000 as Time.Micro)).toThrow();
Expand Down Expand Up @@ -1445,7 +1445,7 @@ for (const code of [StreamCode.Cancel, StreamCode.Internal, StreamCode.Old, Stre
timestamp: Time.Timestamp.now(),
});
await settle();
group.close(new StreamError(code));
group.close(new NetError.Stream(code));
await settle();
expect((await consumer.next())?.frame?.payload).toEqual(new Uint8Array([1]));
expect((await consumer.next())?.frame).toBeUndefined();
Expand All @@ -1470,9 +1470,9 @@ for (const code of [StreamCode.Cancel, StreamCode.Internal, StreamCode.Old, Stre

for (const end of [
null,
new StreamError(StreamCode.Cancel),
new StreamError(StreamCode.Internal),
new SessionError(SessionCode.ProtocolViolation),
new NetError.Stream(StreamCode.Cancel),
new NetError.Stream(StreamCode.Internal),
new NetError.Session(SessionCode.ProtocolViolation),
]) {
test(`Consumer settles a pending read when the track ends: ${end}`, async () => {
const track = new Track.Producer("test");
Expand All @@ -1495,7 +1495,7 @@ test("a group reset does not hide a later container decode failure", async () =>
try {
const reset = track.appendGroup();
await settle();
reset.close(new StreamError(StreamCode.Cancel));
reset.close(new NetError.Stream(StreamCode.Cancel));
await settle();
const malformed = track.appendGroup();
malformed.writeFrame({ payload: new Uint8Array(), timestamp: Time.Timestamp.now() });
Expand All @@ -1509,7 +1509,11 @@ test("a group reset does not hide a later container decode failure", async () =>
}
});

for (const end of [null, new StreamError(StreamCode.Internal), new SessionError(SessionCode.ProtocolViolation)]) {
for (const end of [
null,
new NetError.Stream(StreamCode.Internal),
new NetError.Session(SessionCode.ProtocolViolation),
]) {
test(`Consumer drains a permanent buffered gap after track termination: ${end}`, async () => {
const track = new Track.Producer("test");
const consumer = new Consumer(replay(track), {
Expand Down
2 changes: 1 addition & 1 deletion js/hang/src/container/consumer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -249,7 +249,7 @@ export class Consumer {
// frames decoded before the error are still valid and playable.
// The tail is gone though, so the next group does not continue this one.
group.truncated = true;
if (!(err instanceof Moq.StreamError)) throw err;
if (!(err instanceof Moq.Error.Stream)) throw err;
} finally {
group.done = true;

Expand Down
4 changes: 2 additions & 2 deletions js/hang/src/container/track.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import { Time, type Track } from "@moq/net";
// How long a media track asks its publisher (and, through TRACK_INFO, every relay) to keep a
// non-latest group fetchable. Must match `hang::container::MAX_AGE` in
// rs/hang/src/container/frame.rs.
const MAX_AGE_MS = 30_000;
const MAX_AGE_MS = Time.Milli(30_000);

/**
* Track properties for a track carrying media frames, for `Request.accept`.
Expand Down Expand Up @@ -53,5 +53,5 @@ export type TrackInfoOptions = {
* A RETENTION budget, not a delivery one, so it never makes anyone play further behind live
* and lowering it does not reduce latency: it only shortens how far back a fetch can reach.
*/
maxAge?: number;
maxAge?: Time.Milli;
};
4 changes: 2 additions & 2 deletions js/json/src/snapshot/compression.test.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,14 @@
import { expect, test } from "bun:test";
import { Decoder } from "@moq/flate";
import { Track } from "@moq/net";
import { Time, Track } from "@moq/net";
import { Consumer } from "./consumer.ts";
import { Producer } from "./producer.ts";

type Value = Record<string, unknown>;

const enc = new TextEncoder();
const dec = new TextDecoder();
const REPLAY_LATENCY = 30_000;
const REPLAY_LATENCY = Time.Milli(30_000);

// Reconstruct every value a compressed consumer yields, in order.
async function drainCompressed(track: Track.Subscriber): Promise<Value[]> {
Expand Down
2 changes: 1 addition & 1 deletion js/json/src/snapshot/producer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ export class Producer<T> {
// replacement before the frame is written, so discovering the limit inside `writeFrame` would
// leave an empty newest group behind: a snapshot consumer jumps to the newest, so the last
// good value would vanish even though this update reported an error.
if (encoded.payload.byteLength > Moq.Group.MAX_GROUP_CACHE_BYTES) throw new Moq.Group.FrameTooLarge();
if (encoded.payload.byteLength > Moq.Group.MAX_GROUP_CACHE_BYTES) throw new Moq.Error.FrameTooLarge();

if (encoded.keyframe) {
// The previous group is complete; no more frames will be appended to it. Drop the handle
Expand Down
6 changes: 3 additions & 3 deletions js/json/src/snapshot/snapshot.test.ts
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
import { expect, test } from "bun:test";
import { Group, Track } from "@moq/net";
import { Group, Error as NetError, Time, Track } from "@moq/net";
import { Consumer } from "./consumer.ts";
import { Producer } from "./producer.ts";

type Value = Record<string, unknown>;

// These tests inspect complete finished timelines, so request a replay window
// instead of the transport's live-edge default.
const REPLAY_LATENCY = 30_000;
const REPLAY_LATENCY = Time.Milli(30_000);

// Reconstruct every value a consumer yields, in order.
async function drain(track: Track.Subscriber): Promise<Value[]> {
Expand Down Expand Up @@ -323,7 +323,7 @@ test("a rejected update leaves the previous value readable", async () => {

// Serializes past the group cache limit, so the frame cannot be published.
const oversized = { big: "x".repeat(Group.MAX_GROUP_CACHE_BYTES + 1) };
expect(() => producer.update(oversized)).toThrow(Group.FrameTooLarge);
expect(() => producer.update(oversized)).toThrow(NetError.FrameTooLarge);
producer.finish();

// A reader arriving now still finds the last good value, not an empty superseding group.
Expand Down
Loading
Loading