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
18 changes: 12 additions & 6 deletions doc/bin/relay/cluster.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,10 @@ same way so the cluster converges instead of flapping. Both wire protocols
carry it: natively on moq-lite, and via the [cluster extension](/draft/moq-cluster)
on moq-transport 17+.

Failover routes must carry copies of the same broadcast. For each track, the
Failover routes must carry copies of the same broadcast. A relay moves a
subscription only between sources from the same origin: on moq-lite-07 the one a
source's SUBSCRIBE\_OK or FETCH\_OK names, otherwise the first hop of its route.
A change of origin ends the subscription and the viewer re-subscribes. For each track, the
relay requires matching timescale, retention window, publisher priority, and
group ordering. A source with different properties is refused before its groups
are spliced in. If no compatible source remains, the track fails with
Expand Down Expand Up @@ -50,11 +53,14 @@ prefixes (`grant/**`); an over-wide pattern is refused rather than clamped.

Routing prefers the most specific pattern, then a fully identified hop list
over one that holds a 0 (an anonymous hop) at any depth, then the lowest cost,
then the shortest hop list, breaking any remaining tie toward the newest
announcement so a reconnecting publisher isn't outranked by the session it
replaced. An assigned identity for an anonymous peer is local selection state
and is never written into the hop list. Resolving a non-prefix pattern into a
subscription is not implemented yet.
then the shortest hop list, then a hash of the requested path and the hop list,
breaking any remaining tie toward the newest announcement so a reconnecting
publisher isn't outranked by the session it replaced. Hashing the requested
path spreads equal-cost advertisers of one prefix, such as a transcode pool,
across its paths instead of sending every path to one of them, and every relay
picks the same one for a given path. An assigned identity for an anonymous peer
is local selection state and is never written into the hop list. Resolving a
non-prefix pattern into a subscription is not implemented yet.

```toml
[cluster]
Expand Down
54 changes: 43 additions & 11 deletions drafts/draft-lcurley-moq-lite.md
Original file line number Diff line number Diff line change
Expand Up @@ -416,12 +416,21 @@ The per-subscriber winner changing travels as an ANNOUNCE_UPDATE; the last quali
When serving a subscription, a publisher MUST select the source by that same exclusion; if only excluded sources remain, the subscription is unroutable.
Applying one rule to both advertisement and dispatch keeps advertised paths truthful, which is what prevents subscription cycles of any length.

When resolving a path covered by several routes (across any number of streams), the subscriber SHOULD prefer the most specific covering route (see [Resolution](#resolution)), then a path that contains no 0 Hop ID over one that does, then the lowest Warm Route Cost after adding each arriving link's cost (see [Cost Parameter](#cost-parameter)), breaking ties toward the lowest Cold Route Cost, then toward the shortest path, and then toward the most recently received, so a reconnecting publisher is not outranked by the stale session it replaced.

A route's identity is its first hop: the endpoint that originated it (see [ANNOUNCE_START](#announce-start)).
Two routes covering one path with the same non-zero first hop are the same origin reached different ways, and a relay MAY move a live subscription between them, resuming at a group boundary, so a route change the identity survives (a reconnect, a cheaper path, a draining session) is invisible to the subscriber.
Across differing first hops, or where either is 0, the routes promise nothing about each other's content: a relay MUST NOT splice a live subscription across them, and when the serving session ends, in-flight subscriptions end with it (a reset) and the subscriber re-requests through the best remaining route.
Equal first hops promise the same origin, not interchangeable bytes; what a resuming relay serves next is whatever that origin publishes next at the group boundary.
When resolving a path covered by several routes (across any number of streams), the subscriber SHOULD prefer the most specific covering route (see [Resolution](#resolution)), then a path that contains no 0 Hop ID over one that does, then the lowest Warm Route Cost after adding each arriving link's cost (see [Cost Parameter](#cost-parameter)), breaking ties toward the lowest Cold Route Cost, then toward the shortest path, then toward the lowest Spread Hash, and then toward the most recently received, so a reconnecting publisher is not outranked by the stale session it replaced.

The Spread Hash is the 64-bit FNV-1a hash, with offset basis `0x420C0DECB00B` and the standard FNV-64 prime, of the requested path's UTF-8 bytes followed by each Hop ID of the route's path, oldest first, as 8 little-endian bytes.
It is keyed on the requested path rather than the route's prefix, so equal-cost advertisers of one prefix share its paths instead of the first one taking them all, while one path resolves to the same advertiser on every relay that holds the same routes.
When choosing which route to advertise for a prefix, the requested path is the prefix itself.

A subscription's identity is the origin serving it, named by the `Origin` field of the reply that carries its content: [SUBSCRIBE_OK](#subscribe-ok) for a subscription and [FETCH_OK](#fetch-ok) for a fetch.
A relay learns it from the reply rather than the route: a route promises only that paths under its prefix are servable, and an advertiser serving a prefix from several origins advertises one route for all of them.
Two sources of one track whose replies name the same non-zero Origin are the same origin reached different ways, and a relay MAY move a live subscription between them, resuming at a group boundary, so a change the identity survives (a reconnect, a draining session, a relay failing over within a pool) is invisible to the subscriber.
Across differing Origins, or where either is 0, the sources promise nothing about each other's content: a relay MUST NOT splice a live subscription across them, and instead ends it (a reset) so the subscriber re-requests through the best remaining route.
Group Streams are not ordered with the Subscribe Stream, so a relay that splices on Origin MUST hold a subscription's groups until its SUBSCRIBE_OK names their origin, and discard them if that origin is not the one it serves.
Datagrams cannot be held, so such a relay MUST drop a subscription's datagrams until then; a publisher sends SUBSCRIBE_OK before a subscription's first datagram as well as its first Group Stream.
A relay learns a replacement's Origin only once it replies, so it SHOULD keep serving from a live source when a better route appears rather than trade it for a source that may end the subscription.
A source reached over an earlier version, whose replies carry no Origin, is identified by its route's first hop: the endpoint that originated the route (see [ANNOUNCE_START](#announce-start)).
Equal Origins promise the same origin, not interchangeable bytes; what a resuming relay serves next is whatever that origin publishes next at the group boundary.

#### Resolution {#resolution}
A SUBSCRIBE, FETCH, or TRACK request names a path, and the receiver resolves it against the routes covering that path, after the per-subscriber exclusion above.
Expand Down Expand Up @@ -467,9 +476,9 @@ A subscriber opens a Fetch Stream (0x3) to request a single Group from a Track.

The subscriber sends a FETCH message containing the broadcast path, track name, priority, group sequence, and the frame range within that group.
Unlike SUBSCRIBE, FETCH works on both live and ended broadcasts; it is the only way to read an ended one.
The publisher responds with FRAME messages directly on the same bidirectional stream — there is no response header.
The publisher responds with a FETCH_OK naming the origin serving the group, followed by FRAME messages on the same bidirectional stream.
The Subscribe ID, Group Sequence, and index of the first returned frame are implicit, taken from the original FETCH request.
Because there is no response header, a publisher that cannot serve the requested frame range in full MUST reset the stream rather than return a shorter run; the subscriber has no way to learn where a truncated response actually started.
Because the response carries no position, a publisher that cannot serve the requested frame range in full MUST reset the stream rather than return a shorter run; the subscriber has no way to learn where a truncated response actually started.
As with a subscription, the subscriber MUST already have the track's [TRACK_INFO](#track-info) to parse the returned frames; because the properties are immutable, a single Track Stream lookup is reused across every FETCH of that track (group-by-group fetches do not re-fetch it).
The publisher FINs the stream after the last frame, or resets the stream on error.

Expand Down Expand Up @@ -1119,13 +1128,14 @@ Common values include `1000` (milliseconds), `1000000` (microseconds), `48000` (
A SUBSCRIBE_OK message confirms a subscription and resolves its absolute start position.
It is the first message the publisher sends on the Subscribe Stream, once the start position is known.

This is the trimmed-down counterpart of MoqTransport's SUBSCRIBE_OK: it retains the name and the role of the publisher's positive response, but carries only the resolved start position (all other per-track properties live in [TRACK_INFO](#track-info)).
This is the trimmed-down counterpart of MoqTransport's SUBSCRIBE_OK: it retains the name and the role of the publisher's positive response, but carries only the resolved start position and who serves it (all other per-track properties live in [TRACK_INFO](#track-info)).

~~~
SUBSCRIBE_OK Message {
Type (i) = 0x0
Message Length (i)
Group (i)
Origin (i)
}
~~~

Expand All @@ -1146,6 +1156,12 @@ The subscriber derives the start frame from `Group` and its own request:
The second case is easy to get wrong, so to be explicit: a subscriber that requested group 5 frame 15 and receives `Group` = 6 starts at **frame 0** of group 6, not frame 15.
The frame offset belonged to group 5 and is gone along with the rest of it; it does not carry forward to whichever group the publisher resolved to.

**Origin**:
The Hop ID of the origin serving the subscription, which relays splice failover on (see [Routing](#routing)).
An endpoint names the Origin its own source's reply named, or for a source reached over an earlier version, the first hop of that source's route, and for content it produces, its own Hop ID.
Where none of these identifies anyone (0, or an endpoint without a stable Hop ID of its own), it SHOULD generate a random Hop ID for that content and name it for as long as the content lasts, so a relay downstream can still resume within it.
A value of 0 names nobody, and a relay never splices across it.

## SUBSCRIBE_END {#subscribe-end}
A SUBSCRIBE_END message is sent by the publisher to signal that no group at or after a given sequence will be produced.

Expand Down Expand Up @@ -1212,13 +1228,27 @@ The last frame to return (inclusive), encoded as the absolute frame index + 1.
A value of 0 means through the end of the group (default).
A `Frame End` below `Frame Start` once decoded is a protocol violation; equal bounds are a legal single-frame range.

The publisher responds with FRAME messages directly on the same stream — there is no response header.
The subscriber parses them using the track's [TRACK_INFO](#track-info), which it MUST already have (see the [Track Stream](#track-stream)); the group sequence and the index of the first frame are implicit from the FETCH request.
The publisher responds with a [FETCH_OK](#fetch-ok) followed by FRAME messages on the same stream.
The subscriber parses the frames using the track's [TRACK_INFO](#track-info), which it MUST already have (see the [Track Stream](#track-stream)); the group sequence and the index of the first frame are implicit from the FETCH request.
The publisher FINs the stream after the last frame, or resets on error.
There is no FETCH_ERROR message — the publisher signals failure by resetting the stream.
A publisher holding fewer frames than requested MUST reset rather than truncate, since a short response is indistinguishable from one that started elsewhere.
A group that ends before `Frame End` is not a truncation: the publisher FINs after the last frame it has, provided the group is complete and it served everything from `Frame Start` onward.

## FETCH_OK {#fetch-ok}
FETCH_OK is the publisher's answer on a Fetch Stream, sent once the group is resolved and before its first FRAME.

~~~
FETCH_OK Message {
Message Length (i)
Origin (i)
}
~~~

**Origin**:
The Hop ID of the origin serving the group, as in [SUBSCRIBE_OK](#subscribe-ok).
A relay that splices on Origin MUST NOT deliver the frames of a fetch whose Origin is not the one it serves.

## PROBE
PROBE is used to measure the available bitrate of the connection.

Expand Down Expand Up @@ -1334,6 +1364,8 @@ The `Message Length` describes the payload size on the wire.
- Added `Stream Count` to SUBSCRIBE_END: the number of Group Streams opened for the subscription. SUBSCRIBE_END is now sent once every counted Group Stream has opened, rather than as soon as the final group is known.
- Removed SUBSCRIBE_DROP and its type 0x2; a group without a Group Stream is not counted.
- The Subscribe Stream FIN now follows once every counted Group Stream has finished or been reset.
- Added `Origin` to SUBSCRIBE_OK and the FETCH_OK message ahead of a fetch's frames: the origin serving the request. A subscription's identity is now the Origin its reply names rather than its route's first hop, so a relay splices a failover only between sources naming the same non-zero Origin, including within a pool advertised by one route, and holds a subscription's groups (dropping its datagrams) until its SUBSCRIBE_OK names their origin. SUBSCRIBE_OK now precedes a subscription's first datagram too.
- Added the Spread Hash tie-break after the shortest path: a hash of the requested path and the route's Hop IDs, so equal-cost advertisers of one prefix share its paths.
- Added announce compression: ANNOUNCE_START gains `Path Base` and `Path Keep` to copy the head of a live advertisement's suffix, and ANNOUNCE_START and ANNOUNCE_UPDATE gain `Hop Base` and `Hop Keep` to copy the tail of a live advertisement's Hop ID list.

## moq-lite-06
Expand Down
21 changes: 20 additions & 1 deletion js/net/src/broadcast.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
*/
import { type GetPromise, Once, Signal } from "@moq/signals";
import type { Consumer as GroupConsumer } from "./group.ts";
import { Route } from "./hop.ts";
import { type Hop, Route, randomHop } from "./hop.ts";
import { hooks, type TrackSequence } from "./internal.ts";
import * as track from "./track.ts";
import { registerWire, trackOf, type Broadcast as Wire } from "./wire.ts";
Expand All @@ -30,6 +30,17 @@ class BroadcastState {
// Live consumer handles sharing this state (see {@link Consumer.clone}). The broadcast
// closes once the last one closes, so a shared consumer can be handed to several callers.
consumers = 0;
// The origin serving this broadcast (see the wire's `origin`): named by an upstream
// reply, or generated on first use for content nobody named.
origin?: Hop;
}

// The origin a peer is told serves this broadcast: proxied from upstream, or a random
// one for content originating here, stable for the broadcast's life and shared by every
// session serving it.
function origin(state: BroadcastState): Hop {
state.origin ??= randomHop();
return state.origin;
}

function dequeueRequest(state: BroadcastState): track.Request | undefined {
Expand Down Expand Up @@ -238,6 +249,10 @@ export class Producer {
resolveTrackInfo: (name) => resolveTrackInfo(this.#state, name),
fetchGroup: (name, sequence, options) => fetchGroup(this.#state, name, sequence, options),
requested: () => this.#requested(),
origin: () => origin(this.#state),
name: (named) => {
this.#state.origin = named;
},
};
}

Expand Down Expand Up @@ -300,6 +315,10 @@ export class Consumer {
resolveTrackInfo: (name) => resolveTrackInfo(this.#state, name),
fetchGroup: (name, sequence, options) => fetchGroup(this.#state, name, sequence, options),
requested: () => this.#requested(),
origin: () => origin(this.#state),
name: (named) => {
this.#state.origin = named;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Reject conflicting origin names instead of overwriting

If two replies for one consumed broadcast name different origins, this assignment silently replaces the first identity after its content may already have been delivered. A JavaScript application republishing that broadcast can then label existing origin A content as origin B, allowing a downstream relay to admit or splice it under the wrong identity. Treat a second distinct name as a protocol failure, matching the Rust provenance/front checks, rather than accepting malformed peer state.

AGENTS.md reference: AGENTS.md:L17-L17

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not taking this one. A second distinct name for one consumed broadcast is not malformed: the upstream relay's front ends on an origin mismatch, and a later request for the same announced broadcast can land on a new front serving another origin. Throwing there would tear down a healthy session. The real gap is granularity (JS records the origin per broadcast, not per track), which only matters for a JS app republishing a broadcast across an upstream failover; a downstream Rust relay still refuses to splice across a mismatch. Listed under Known gaps in the description for the maintainer to decide on.

(Written by Claude Opus 5.5)

},
});
}

Expand Down
14 changes: 11 additions & 3 deletions js/net/src/integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -462,10 +462,18 @@ for (const [protocol, carriesOptIn] of [
});
}

test("integration: lite draft-05 datagram delivery", async () => {
// Draft-07 holds a subscription's content until SUBSCRIBE_START names its origin, so a
// datagram-only track must still send one for any datagram to arrive.
for (const protocol of [Lite.ALPN_05, Lite.ALPN_07_WIP]) {
test(`integration: ${protocol} datagram delivery`, async () => {
await datagramDelivery(protocol);
});
}

async function datagramDelivery(protocol: string) {
const enc = new TextEncoder();
const dec = new TextDecoder();
const pair = createMockTransportPair(Lite.ALPN_05);
const pair = createMockTransportPair(protocol);
const origin = new OriginProducer();

const [client, server] = await Promise.all([
Expand Down Expand Up @@ -502,7 +510,7 @@ test("integration: lite draft-05 datagram delivery", async () => {
remote.close();
client.close();
server.close();
});
}

test("integration: lite draft-05 datagrams not sent on a non-datagram transport", async () => {
const enc = new TextEncoder();
Expand Down
20 changes: 20 additions & 0 deletions js/net/src/internal.ts
Original file line number Diff line number Diff line change
Expand Up @@ -189,3 +189,23 @@ export const hooks: {
throw new Error("broadcast.ts not loaded");
},
};

/**
* Spreads equal routes across paths: FNV-1a 64 of `path` then each hop, oldest first, as 8
* little-endian bytes. Keyed on the requested path so an equal-cost pool advertising one
* prefix shares its paths, and every node holding the same routes picks the same member.
* Mirrors `fnv_key` in `rs/moq-net`; the seed is the draft's Spread Hash offset basis.
*/
export function spreadHash(path: string, hops: readonly bigint[]): bigint {
const prime = 0x100000001b3n;
let hash = 0x420c0decb00bn;
for (const byte of new TextEncoder().encode(path)) {
hash = BigInt.asUintN(64, (hash ^ BigInt(byte)) * prime);
}
for (const hop of hops) {
for (let shift = 0n; shift < 64n; shift += 8n) {
hash = BigInt.asUintN(64, (hash ^ ((hop >> shift) & 0xffn)) * prime);
}
}
return hash;
}
19 changes: 18 additions & 1 deletion js/net/src/lite/fetch.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
import { expect, test } from "bun:test";
import { HopSchema } from "../hop.ts";
import * as Path from "../path.ts";
import { Reader, Writer } from "../stream.ts";
import { Fetch } from "./fetch.ts";
import { Fetch, FetchOk } from "./fetch.ts";
import { Version } from "./version.ts";

function concat(chunks: Uint8Array[]): Uint8Array {
Expand Down Expand Up @@ -44,3 +45,19 @@ test("Fetch round-trips on draft-03/04/05", async () => {
expect(got.group).toBe(42);
}
});

test("FetchOk names the origin on draft-07 only", async () => {
const written: Uint8Array[] = [];
const writer = new Writer(
new WritableStream<Uint8Array>({ write: (chunk) => void written.push(new Uint8Array(chunk)) }),
);
await new FetchOk(HopSchema.parse(42n)).encode(writer, Version.DRAFT_07);
writer.close();
await writer.closed;
const buf = concat(written);
expect(buf).toEqual(new Uint8Array([1, 42]));
const got = await FetchOk.decode(new Reader(undefined, buf), Version.DRAFT_07);
expect(got.origin).toBe(HopSchema.parse(42n));

await expect(new FetchOk(HopSchema.parse(42n)).encode(writer, Version.DRAFT_06)).rejects.toThrow();
});
Loading
Loading