diff --git a/doc/bin/relay/cluster.md b/doc/bin/relay/cluster.md index ddeb3ae081..13e5c90955 100644 --- a/doc/bin/relay/cluster.md +++ b/doc/bin/relay/cluster.md @@ -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 @@ -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] diff --git a/drafts/draft-lcurley-moq-lite.md b/drafts/draft-lcurley-moq-lite.md index 532b7bafb5..26000264e9 100644 --- a/drafts/draft-lcurley-moq-lite.md +++ b/drafts/draft-lcurley-moq-lite.md @@ -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. @@ -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. @@ -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) } ~~~ @@ -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. @@ -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. @@ -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 diff --git a/js/net/src/broadcast.ts b/js/net/src/broadcast.ts index 41cbeca2ef..ab0268320c 100644 --- a/js/net/src/broadcast.ts +++ b/js/net/src/broadcast.ts @@ -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"; @@ -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 { @@ -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; + }, }; } @@ -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; + }, }); } diff --git a/js/net/src/integration.test.ts b/js/net/src/integration.test.ts index 3107cb2b14..7072553973 100644 --- a/js/net/src/integration.test.ts +++ b/js/net/src/integration.test.ts @@ -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([ @@ -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(); diff --git a/js/net/src/internal.ts b/js/net/src/internal.ts index 3f235938cd..e25b037c03 100644 --- a/js/net/src/internal.ts +++ b/js/net/src/internal.ts @@ -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; +} diff --git a/js/net/src/lite/fetch.test.ts b/js/net/src/lite/fetch.test.ts index 0eb06dca94..f99f2d522b 100644 --- a/js/net/src/lite/fetch.test.ts +++ b/js/net/src/lite/fetch.test.ts @@ -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 { @@ -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({ 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(); +}); diff --git a/js/net/src/lite/fetch.ts b/js/net/src/lite/fetch.ts index 35fe2e17de..e937b3c0b1 100644 --- a/js/net/src/lite/fetch.ts +++ b/js/net/src/lite/fetch.ts @@ -1,7 +1,8 @@ +import { type Hop, HopSchema } from "../hop.ts"; import * as Path from "../path.ts"; import type { Reader, Writer } from "../stream.ts"; import * as Message from "./message.ts"; -import { hasFrameBounds, Version } from "./version.ts"; +import { hasFrameBounds, hasOrigin, Version } from "./version.ts"; function guardFetch(version: Version) { switch (version) { @@ -102,3 +103,26 @@ export class Fetch { return Message.decode(r, (r) => Fetch.#decode(r, version)); } } + +/** + * FETCH_OK: the publisher's answer on a Fetch Stream, ahead of the frames, naming the + * origin serving the group (the Hop ID a relay stitches failover on). Draft-07+ only; + * older versions answer with the frames alone. + */ +export class FetchOk { + origin: Hop; + + constructor(origin: Hop) { + this.origin = origin; + } + + async encode(w: Writer, version: Version): Promise { + if (!hasOrigin(version)) throw new Error("FETCH_OK not supported for this version"); + return Message.encode(w, (w) => w.u62(this.origin)); + } + + static async decode(r: Reader, version: Version): Promise { + if (!hasOrigin(version)) throw new Error("FETCH_OK not supported for this version"); + return Message.decode(r, async (r) => new FetchOk(HopSchema.parse(await r.u62()))); + } +} diff --git a/js/net/src/lite/origin.test.ts b/js/net/src/lite/origin.test.ts new file mode 100644 index 0000000000..6518cb44ef --- /dev/null +++ b/js/net/src/lite/origin.test.ts @@ -0,0 +1,76 @@ +import { expect, test } from "bun:test"; +import * as broadcast from "../broadcast.ts"; +import { HopSchema, randomHop, UNKNOWN_HOP } from "../hop.ts"; +import { createMockTransportPair } from "../mock.ts"; +import * as Path from "../path.ts"; +import { Reader, Stream } from "../stream.ts"; +import { wireOf } from "../wire.ts"; +import { Group as GroupMessage } from "./group.ts"; +import { StreamId } from "./stream.ts"; +import { encodeSubscribeResponse, Subscribe, SubscribeStart } from "./subscribe.ts"; +import { Subscriber } from "./subscriber.ts"; +import { TrackInfo, Track as TrackMessage } from "./track.ts"; +import { ALPN_07_WIP, Version } from "./version.ts"; + +const VERSION = Version.DRAFT_07; + +/** Whether `promise` settles within `ms`. */ +async function settlesWithin(promise: Promise, ms: number): Promise { + let timer: ReturnType | undefined; + const pending = new Promise((resolve) => { + timer = setTimeout(() => resolve(false), ms); + }); + try { + return await Promise.race([promise.then(() => true), pending]); + } finally { + clearTimeout(timer); + } +} + +test("a broadcast originating here names one random origin for its life", () => { + const producer = new broadcast.Producer(); + const origin = wireOf(producer).origin(); + expect(origin).not.toBe(UNKNOWN_HOP); + expect(wireOf(producer).origin()).toBe(origin); + // Every handle, and so every session serving it, names the same one. + expect(wireOf(producer.consume()).origin()).toBe(origin); + expect(wireOf(new broadcast.Producer()).origin()).not.toBe(origin); +}); + +test("a draft-07 group waits for SUBSCRIBE_START, whose origin the broadcast then names", async () => { + const pair = createMockTransportPair(ALPN_07_WIP); + const subscriber = new Subscriber(pair.client, VERSION, randomHop()); + const consumer = subscriber.consume(Path.from("room")); + const reader = consumer.track("video").subscribe({}); + + const info = await Stream.accept(pair.server); + if (!info) throw new Error("the subscriber never asked for TRACK_INFO"); + expect(await info.reader.u53()).toBe(StreamId.Track); + await TrackMessage.decode(info.reader, VERSION); + await new TrackInfo({ maxAge: 60_000 }).encode(info.writer, VERSION); + info.close(); + + const sub = await Stream.accept(pair.server); + if (!sub) throw new Error("the subscriber never subscribed"); + expect(await sub.reader.u53()).toBe(StreamId.Subscribe); + await Subscribe.decode(sub.reader, VERSION); + + // The group's stream races ahead of the subscribe stream. + let controller!: ReadableStreamDefaultController; + const readable = new ReadableStream({ start: (c) => (controller = c) }); + void subscriber.runGroup(new GroupMessage({ subscribe: 0n, sequence: 0 }), new Reader(readable)); + controller.enqueue(new Uint8Array([0, 1, 120])); + controller.close(); + + const next = reader.recvGroup(); + expect(await settlesWithin(next, 50)).toBe(false); + + const origin = HopSchema.parse(42n); + await encodeSubscribeResponse(sub.writer, { start: new SubscribeStart(0, origin) }, VERSION); + const group = await next; + expect(group?.sequence).toBe(0); + expect(await group?.readString()).toBe("x"); + + // Republishing the broadcast proxies the origin upstream named. + expect(wireOf(consumer).origin()).toBe(origin); +}); diff --git a/js/net/src/lite/publisher.ts b/js/net/src/lite/publisher.ts index f35fce72be..3efad512be 100644 --- a/js/net/src/lite/publisher.ts +++ b/js/net/src/lite/publisher.ts @@ -13,7 +13,7 @@ import { type Advertised, wireOf } from "../wire.ts"; import { AnnounceInit, AnnounceOk, type AnnounceRequest, encodeAnnounceBroadcast } from "./announce.ts"; import { Datagram as DatagramMessage } from "./datagram.ts"; import * as DatagramStream from "./datagram_stream.ts"; -import type { Fetch } from "./fetch.ts"; +import { type Fetch, FetchOk } from "./fetch.ts"; import { Group as GroupMessage } from "./group.ts"; import { Priority, sendOrder } from "./priority.ts"; import { Probe } from "./probe.ts"; @@ -31,6 +31,7 @@ import { hasAnnounceId, hasAnnounceOk, hasDatagrams, + hasOrigin, hasProbeRtt, hasStreamCount, resolvesStart, @@ -282,6 +283,40 @@ class SubscriptionControls { } } +/** + * A lite-05+ subscription's SUBSCRIBE_START and SUBSCRIBE_END, written in order. The group + * and datagram loops both serve the subscription, so START goes out ahead of whichever serves + * first: it resolves the start and names the origin, which a draft-07 subscriber needs before + * it delivers either. + */ +class SubscribeResponses { + #controls: SubscriptionControls; + #start: (sequence: number) => Promise; + #writes: Promise = Promise.resolve(true); + #started?: Promise; + + constructor(controls: SubscriptionControls, start: (sequence: number) => Promise) { + this.#controls = controls; + this.#start = start; + } + + /** Sends SUBSCRIBE_START at `sequence`, once; false when peer departure superseded it. */ + start(sequence: number): Promise { + this.#started ??= this.#write(() => this.#start(sequence)); + return this.#started; + } + + /** Sends SUBSCRIBE_END after any START still being written; false as for {@link start}. */ + end(write: () => Promise): Promise { + return this.#write(write); + } + + #write(write: () => Promise): Promise { + this.#writes = this.#writes.then((ok) => ok && this.#controls.response(write())); + return this.#writes; + } +} + // A microtask is too short: decoding one framed update crosses several awaits, each of which // can requeue behind the serving continuation. A task boundary lets the decoder finish whatever // the transport already delivered before the next group pop. Updates are rare, so groups do not @@ -621,12 +656,6 @@ export class Publisher { console.debug(`publish ok: broadcast=${msg.broadcast} track=${track.name}`); - // Serve datagrams concurrently with groups whenever the transport carries them - // (the writer exists iff so). No group fallback: otherwise they simply aren't sent. - if (this.#datagramWriter) { - datagrams = this.#runDatagrams(msg.id, track, timescale); - } - controls = new SubscriptionControls({ reader: stream.reader, writer: stream.writer, @@ -643,16 +672,40 @@ export class Publisher { }); }, }); + + const bounds: FrameBounds = { + startGroup: msg.startGroup, + startFrame: msg.startFrame, + endGroup: msg.endGroup, + endFrame: msg.endFrame, + }; + const responses = new SubscribeResponses(controls, async (sequence) => { + // SUBSCRIBE_START promises nothing below this sequence will be delivered. + // Arrival-order serving could later surface a straggler below the first + // group, so pin the floor to what was announced. + hooks.replaceGroups(track, { + start: { included: sequence }, + end: bounds.endGroup === undefined ? undefined : { included: bounds.endGroup }, + }); + // Read once content flowed: an upstream's SUBSCRIBE_START named the origin + // before any of its content did. + const start = new SubscribeStart(sequence, wireOf(front).origin()); + await encodeSubscribeResponse(stream.writer, { start }, this.version); + }); + + // Serve datagrams concurrently with groups whenever the transport carries them + // (the writer exists iff so, and only on lite-05+). No group fallback: otherwise + // they simply aren't sent. + if (this.#datagramWriter) { + datagrams = this.#runDatagrams(msg.id, track, timescale, responses); + } + await this.#runTrack(track, stream.writer, controls, { sub: msg.id, broadcast: msg.broadcast, timescale, - bounds: { - startGroup: msg.startGroup, - startFrame: msg.startFrame, - endGroup: msg.endGroup, - endFrame: msg.endFrame, - }, + bounds, + responses, }); console.debug(`publish done: broadcast=${msg.broadcast} track=${track.name}`); @@ -705,6 +758,9 @@ export class Publisher { // come off the same front, so the metadata and the frames are one generation. const info = await this.#resolveTrackInfo(front, msg.track); group = await wireOf(front).fetchGroup(msg.track, msg.group, { priority: msg.priority }); + if (hasOrigin(this.version)) { + await new FetchOk(wireOf(front).origin()).encode(stream.writer, this.version); + } await this.#runFetchGroup(group, stream.writer, { timescale: Timescale(info.timescale), start: msg.startFrame, @@ -736,13 +792,18 @@ export class Publisher { track: track.Subscriber, stream: Writer, controls: SubscriptionControls, - serving: { sub: bigint; broadcast: Path.Valid; timescale: Timescale; bounds: FrameBounds }, + serving: { + sub: bigint; + broadcast: Path.Valid; + timescale: Timescale; + bounds: FrameBounds; + responses: SubscribeResponses; + }, ) { - const { sub, broadcast, timescale, bounds } = serving; + const { sub, broadcast, timescale, bounds, responses } = serving; // Lite-05+ resolves the range on the subscribe stream: SUBSCRIBE_START once the // first group is known, SUBSCRIBE_END when the track finishes. const emitRange = supportsTrackStream(this.version); - let startSent = false; let endSent = false; // Lite-07+ counts the group streams in SUBSCRIBE_END, so it goes out only once every @@ -765,14 +826,12 @@ export class Publisher { const sendEnd = async (): Promise => { endSent = true; if (!emitRange) return true; - return controls.response( - (async () => { - // A group that gives up before its stream opens is never counted. - if (countStreams) while (opening.size > 0) await Promise.all(opening); - const end = new SubscribeEnd(boundary(), streams); - await encodeSubscribeResponse(stream, { end }, this.version); - })(), - ); + return responses.end(async () => { + // A group that gives up before its stream opens is never counted. + if (countStreams) while (opening.size > 0) await Promise.all(opening); + const end = new SubscribeEnd(boundary(), streams); + await encodeSubscribeResponse(stream, { end }, this.version); + }); }; // One ranking for the whole subscription, shared by every group it serves. @@ -866,26 +925,7 @@ export class Publisher { const group = recv.group; const range = frameRange(bounds, group.sequence); - if (emitRange && !startSent) { - startSent = true; - // SUBSCRIBE_START promises nothing below this sequence will be delivered. - // Arrival-order serving could later surface a straggler below the first - // group, so pin the floor to what was announced. - hooks.replaceGroups(track, { - start: { included: group.sequence }, - end: bounds.endGroup === undefined ? undefined : { included: bounds.endGroup }, - }); - if ( - !(await controls.response( - encodeSubscribeResponse( - stream, - { start: new SubscribeStart(group.sequence) }, - this.version, - ), - )) - ) - return; - } + if (emitRange && !(await responses.start(group.sequence))) return; const options: RunGroup = { sub, @@ -979,7 +1019,7 @@ export class Publisher { * * @internal */ - async #runDatagrams(sub: bigint, track: track.Subscriber, timescale: Timescale) { + async #runDatagrams(sub: bigint, track: track.Subscriber, timescale: Timescale, responses: SubscribeResponses) { const writer = this.#datagramWriter; if (!writer) return; // Only reached with a writer (see the #datagramWriter gate). const maxSize = DatagramStream.maxDatagramSize(this.#quic); @@ -988,6 +1028,7 @@ export class Publisher { for (;;) { const datagram = await track.recvDatagram(); if (!datagram) return; // Track finished; #runTrack tears the subscription down. + if (!(await responses.start(datagram.sequence))) return; // Convert the timestamp to the track's advertised timescale, matching #serveGroup. const ts = Math.round(datagram.timestamp.as(timescale)); diff --git a/js/net/src/lite/subscribe.test.ts b/js/net/src/lite/subscribe.test.ts index afd4c7c0a6..a185678a37 100644 --- a/js/net/src/lite/subscribe.test.ts +++ b/js/net/src/lite/subscribe.test.ts @@ -1,4 +1,5 @@ import { expect, test } from "bun:test"; +import { HopSchema, UNKNOWN_HOP } from "../hop.ts"; import * as Path from "../path.ts"; import { Reader, Writer } from "../stream.ts"; import { @@ -151,6 +152,19 @@ test("SubscribeStart round-trips on draft-05", async () => { expect(got.start.group).toBe(42); }); +test("SubscribeStart names the origin on draft-07", async () => { + const start = new SubscribeStart(7, HopSchema.parse(42n)); + // Type, length, group, origin; draft-06 has no room for the origin. + expect(await encode(Version.DRAFT_07, { start })).toEqual(new Uint8Array([0, 2, 7, 42])); + expect(await encode(Version.DRAFT_06, { start })).toEqual(new Uint8Array([0, 1, 7])); + const got = await responseRoundtrip(Version.DRAFT_07, { start }); + if (!("start" in got)) throw new Error("expected start"); + expect([got.start.group, got.start.origin]).toEqual([7, HopSchema.parse(42n)]); + const old = await responseRoundtrip(Version.DRAFT_06, { start }); + if (!("start" in old)) throw new Error("expected start"); + expect(old.start.origin).toBe(UNKNOWN_HOP); +}); + test("SubscribeEnd round-trips on draft-05", async () => { // Type, length, group: no stream count before draft-07. expect(await encode(Version.DRAFT_05, { end: new SubscribeEnd(7, 3) })).toEqual(new Uint8Array([1, 1, 7])); diff --git a/js/net/src/lite/subscribe.ts b/js/net/src/lite/subscribe.ts index 3d1a1f2711..81e3beabea 100644 --- a/js/net/src/lite/subscribe.ts +++ b/js/net/src/lite/subscribe.ts @@ -1,7 +1,8 @@ +import { type Hop, HopSchema, UNKNOWN_HOP } from "../hop.ts"; import * as Path from "../path.ts"; import type { Reader, Writer } from "../stream.ts"; import * as Message from "./message.ts"; -import { hasFrameBounds, hasGroupOrder, hasStreamCount, resolvesStart, Version } from "./version.ts"; +import { hasFrameBounds, hasGroupOrder, hasOrigin, hasStreamCount, resolvesStart, Version } from "./version.ts"; /** * Encode the `Group Start` field shared by SUBSCRIBE and SUBSCRIBE_UPDATE. @@ -432,19 +433,30 @@ export class SubscribeOk { */ export class SubscribeStart { group: number; + /** + * The origin serving the subscription: the Hop ID a relay stitches failover on. + * {@link UNKNOWN_HOP} names nobody. Draft-07+; older versions decode it as unknown. + */ + origin: Hop; - constructor(group: number) { + constructor(group: number, origin: Hop = UNKNOWN_HOP) { this.group = group; + this.origin = origin; } - async encode(w: Writer): Promise { + async encode(w: Writer, version: Version): Promise { return Message.encode(w, async (w) => { await w.u53(this.group); + if (hasOrigin(version)) await w.u62(this.origin); }); } - static async decode(r: Reader): Promise { - return Message.decode(r, async (r) => new SubscribeStart(await r.u53())); + static async decode(r: Reader, version: Version): Promise { + return Message.decode(r, async (r) => { + const group = await r.u53(); + const origin = hasOrigin(version) ? HopSchema.parse(await r.u62()) : UNKNOWN_HOP; + return new SubscribeStart(group, origin); + }); } } @@ -557,7 +569,7 @@ export async function encodeSubscribeResponse(w: Writer, resp: SubscribeResponse // Draft-05+: SUBSCRIBE_OK is gone; START/END/DROP carry the resolved range. if ("start" in resp) { await w.u53(0x0); - await resp.start.encode(w); + await resp.start.encode(w, version); } else if ("end" in resp) { await w.u53(0x1); await resp.end.encode(w, version); @@ -592,7 +604,7 @@ export async function decodeSubscribeResponse(r: Reader, version: Version): Prom const typ = await r.u53(); switch (typ) { case 0x0: - return { start: await SubscribeStart.decode(r) }; + return { start: await SubscribeStart.decode(r, version) }; case 0x1: return { end: await SubscribeEnd.decode(r, version) }; case 0x2: diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index 021fe1dae4..56eb5f4a9e 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -23,7 +23,7 @@ import { } from "./announce.ts"; import { Datagram as DatagramMessage } from "./datagram.ts"; import * as DatagramStream from "./datagram_stream.ts"; -import { Fetch as FetchMessage } from "./fetch.ts"; +import { Fetch as FetchMessage, FetchOk } from "./fetch.ts"; import type { Group as GroupMessage } from "./group.ts"; import { sendOrder } from "./priority.ts"; import { Probe } from "./probe.ts"; @@ -40,7 +40,15 @@ import { SubscribeUpdate, } from "./subscribe.ts"; import { TrackInfo, Track as TrackMessage } from "./track.ts"; -import { hasAnnounceId, hasAnnounceOk, hasDatagrams, hasProbeRtt, restartSupported, Version } from "./version.ts"; +import { + hasAnnounceId, + hasAnnounceOk, + hasDatagrams, + hasOrigin, + hasProbeRtt, + restartSupported, + Version, +} from "./version.ts"; // Bound on how long stream-open plus the first response (SUBSCRIBE_OK on older // drafts, or TRACK_INFO on lite-05+) may take. Browsers cap concurrent QUIC streams @@ -83,6 +91,12 @@ interface SubscribeEntry { // (SUBSCRIBE_END), once it declares them. start?: number; end?: number; + // The broadcast the subscription belongs to, which records the origin SUBSCRIBE_START + // names so a session republishing it names the same one. + broadcast: broadcast.Consumer; + // Whether SUBSCRIBE_START arrived. On draft-07, which names the serving origin there, + // group streams wait for it: until then nobody knows whose content they carry. + started: Signal; } /** @@ -499,14 +513,14 @@ export class Subscriber { for (;;) { const request = await wireOf(consumer).requested(); if (!request) break; - void this.#runSubscribe(path, request); + void this.#runSubscribe(consumer, path, request); } })(); return consumer; } - async #runSubscribe(broadcast: Path.Valid, request: track.Request) { + async #runSubscribe(consumer: broadcast.Consumer, broadcast: Path.Valid, request: track.Request) { const id = this.#subscribeNext++; const subscription = request.subscription; const initialBounds = groupBounds(subscription.groups); @@ -535,7 +549,7 @@ export class Subscriber { // Open the stream under a timeout. The stream handle flows back via `state` // so the timeout path can abort it if it finishes opening after the deadline. const state: { stream?: Stream } = {}; - const setup = this.#openSubscribe(state, msg, request, id, timescale); + const setup = this.#openSubscribe(state, msg, request, id, timescale, consumer); let opened: { stream: Stream; entry: SubscribeEntry }; try { @@ -627,6 +641,7 @@ export class Subscriber { request: track.Request, id: bigint, timescale: Signal, + consumer: broadcast.Consumer, ): Promise<{ stream: Stream; entry: SubscribeEntry }> { let producer: track.Producer; let drainOk = false; @@ -644,7 +659,13 @@ export class Subscriber { } // Register before opening SUBSCRIBE so a racing GROUP stream finds the entry. - const entry: SubscribeEntry = { track: producer, timescale, tail: new Tail() }; + const entry: SubscribeEntry = { + track: producer, + timescale, + tail: new Tail(), + broadcast: consumer, + started: new Signal(!hasOrigin(this.version)), + }; this.#subscribes.set(id, entry); state.stream = await Stream.open(this.#quic); @@ -702,6 +723,7 @@ export class Subscriber { // Open a FETCH stream for one group and stream its bare frames into a group, for the // ConsumeBroadcast backing track.Consumer.fetchGroup() (lite-05+). fetchGroup( + front: broadcast.Consumer, broadcast: Path.Valid, track: string, sequence: number, @@ -721,12 +743,13 @@ export class Subscriber { if (this.#fetches.get(key) === group) this.#fetches.delete(key); }); - return this.#runFetch(broadcast, track, sequence, options, group); + return this.#runFetch(front, broadcast, track, sequence, options, group); } // Open the FETCH stream and pump the response into the shared group. Setup errors close the // group (so coalesced mirrors observe them and the entry evicts) and reject this caller. async #runFetch( + front: broadcast.Consumer, broadcast: Path.Valid, track: string, sequence: number, @@ -748,6 +771,12 @@ export class Subscriber { stream.writer, this.version, ); + // Draft-07 answers with FETCH_OK naming the serving origin before any frame; + // record it so a session republishing the broadcast names the same one. + if (hasOrigin(this.version)) { + const ok = await FetchOk.decode(stream.reader, this.version); + wireOf(front).name(ok.origin); + } } catch (err: unknown) { stream.abort(error(err)); throw err; @@ -822,6 +851,10 @@ export class Subscriber { if ("start" in resp) { entry.start = resp.start.group; + if (hasOrigin(this.version)) { + wireOf(entry.broadcast).name(resp.start.origin); + entry.started.set(true); + } } else if ("end" in resp) { if (entry.end !== undefined) throw new ProtocolViolation("duplicate SUBSCRIBE_END"); entry.end = resp.end.group; @@ -951,11 +984,22 @@ export class Subscriber { return; } - const { track, timescale, tail } = entry; + const { track, timescale, tail, started } = entry; const producer = new netGroup.Producer(group.sequence); const read = tail.open(group.sequence); try { + // Hold the group until SUBSCRIBE_START names whose content it is: the group's + // stream can arrive before the subscribe stream's. + while (!started.peek()) { + if (track.closed.peek() !== undefined) { + producer.close(); + stream.stop(new StreamError(StreamCode.Cancel, { message: "cancel" })); + return; + } + await Signal.race(started, track.closed); + } + track.writeGroup(producer); // Block until the timescale is known; the group's stream can arrive before @@ -1062,6 +1106,10 @@ export class Subscriber { const entry = this.#subscribes.get(dg.subscribe); if (!entry) return; // Unknown or already-closed subscription. + // A datagram cannot wait for SUBSCRIBE_START to name its origin like a group stream + // does, so on draft-07 it is dropped until then. + if (!entry.started.peek()) return; + // Datagrams are lite-05+, which always negotiates a timescale; if it hasn't resolved // yet (the datagram raced ahead of TRACK_INFO), drop rather than guess. const scale = entry.timescale.peek(); @@ -1163,7 +1211,7 @@ class ConsumeBroadcast extends broadcast.Consumer { super(state); overrideBroadcastWire(this, { resolveTrackInfo: (name) => subscriber.resolveTrackInfo(path, name), - fetchGroup: (name, sequence, options) => subscriber.fetchGroup(path, name, sequence, options), + fetchGroup: (name, sequence, options) => subscriber.fetchGroup(this, path, name, sequence, options), }); this.#subscriber = subscriber; this.#path = path; diff --git a/js/net/src/lite/version.ts b/js/net/src/lite/version.ts index 2d83114b6f..96ef8efca5 100644 --- a/js/net/src/lite/version.ts +++ b/js/net/src/lite/version.ts @@ -249,6 +249,22 @@ export function hasStreamCount(version: Version): boolean { } } +/** Whether SUBSCRIBE_OK and FETCH_OK name the origin serving the request, which relays stitch failover on. Added in lite-07, with FETCH_OK itself. */ +export function hasOrigin(version: Version): boolean { + // Explicitly list older versions so future versions keep the lite-07+ behavior. + switch (version) { + case Version.DRAFT_01: + case Version.DRAFT_02: + case Version.DRAFT_03: + case Version.DRAFT_04: + case Version.DRAFT_05: + case Version.DRAFT_06: + return false; + default: + return true; + } +} + /** Whether ANNOUNCE_START and ANNOUNCE_UPDATE may copy a path head or hop-chain tail from a live announcement. Added in lite-07. */ export function hasAnnounceCompression(version: Version): boolean { // Explicitly list older versions so future versions keep the lite-07+ behavior. diff --git a/js/net/src/origin.test.ts b/js/net/src/origin.test.ts index 968eb5ea76..d50aa26072 100644 --- a/js/net/src/origin.test.ts +++ b/js/net/src/origin.test.ts @@ -3,6 +3,7 @@ import { getter } from "@moq/signals"; import { type Consumer as BroadcastConsumer, Producer as BroadcastProducer } from "./broadcast.ts"; import { StreamCode, StreamError } from "./error.ts"; import { HopSchema, Route } from "./hop.ts"; +import { spreadHash } from "./internal.ts"; import type { Consumer, Table } from "./origin.ts"; import { Producer } from "./origin.ts"; import * as Path from "./path.ts"; @@ -1360,3 +1361,45 @@ test("announced filters by arbitrary patterns and reports captures", async () => broadcast.close(); origin.close(); }); + +test("the spread hash matches rs/moq-net byte for byte", () => { + expect(spreadHash("pool/job-0", [10n])).toBe(0xefb5e20a66101c32n); + expect(spreadHash("pool/job-0", [11n])).toBe(0x0eb0a91370ff6653n); +}); + +test("an equal-cost pool spreads its paths the same way on every node", async () => { + const workers = [10n, 11n, 12n, 13n].map((id) => HopSchema.parse(id)); + const paths = Array.from({ length: 64 }, (_, i) => Path.from(`pool/job-${i}`)); + + // The worker each path resolves to on an origin whose pool arrived in `order`. + async function winners(order: typeof workers): Promise { + const origin = new Producer(); + const served = new Map(); + const handles = order.map((hop) => { + const handle = wireOf(origin).receive(Path.from("pool"), { hops: [hop], cost: 3n }); + void (async () => { + for await (const request of handle.requested()) { + served.set(request.path, hop); + request.accept(new BroadcastProducer()); + } + })(); + return handle; + }); + const requests = paths.map((path) => origin.request(path)); + await settle(); + for (const request of requests) request.close(); + for (const handle of handles) handle.close(); + origin.close(); + return paths.map((path) => served.get(path) ?? -1n); + } + + const forward = await winners(workers); + const reverse = await winners([...workers].reverse()); + expect(reverse).toEqual(forward); + + // Not an assertion about any two paths, which a correct hash may put on one worker: + // only that the set does not pile onto a few. + for (const worker of workers) { + expect(forward.filter((hop) => hop === worker).length).toBeGreaterThanOrEqual(paths.length / 16); + } +}); diff --git a/js/net/src/origin.ts b/js/net/src/origin.ts index dbd41e7d06..203b63a8fc 100644 --- a/js/net/src/origin.ts +++ b/js/net/src/origin.ts @@ -15,7 +15,7 @@ import * as announce from "./announced.ts"; import * as broadcast from "./broadcast.ts"; import { StreamCode, StreamError } from "./error.ts"; import { isAnonymous, Route, routesEqual } from "./hop.ts"; -import { hiddenBelow, hooks, scopeCaptures, scopeHead, scopeOverlaps } from "./internal.ts"; +import { hiddenBelow, hooks, scopeCaptures, scopeHead, scopeOverlaps, spreadHash } from "./internal.ts"; import * as Path from "./path.ts"; import { type Advertised, registerWire, wireOf } from "./wire.ts"; @@ -76,8 +76,16 @@ function compareRoutes(a: Route, b: Route): number { return 0; } -/** The preferred of `entries` (newest first) not skipped: the best route, then fewest hops, then newest. */ -function preferredEntry(entries: readonly RouteEntry[], skip?: (entry: RouteEntry) => boolean): RouteEntry | undefined { +/** + * The preferred of `entries` (newest first) for resolving `path`, not skipped: the best route, + * then fewest hops, then the lowest {@link spreadHash}, then newest. `path` is the requested + * path for a request, or the prefix itself for an advertisement. + */ +function preferredEntry( + path: Path.Valid, + entries: readonly RouteEntry[], + skip?: (entry: RouteEntry) => boolean, +): RouteEntry | undefined { let best: RouteEntry | undefined; for (const entry of entries) { if (skip?.(entry)) continue; @@ -87,7 +95,13 @@ function preferredEntry(entries: readonly RouteEntry[], skip?: (entry: RouteEntr } const a = entry.route.peek(); const b = best.route.peek(); - const order = compareRoutes(a, b) || a.hops.length - b.hops.length; + let order = compareRoutes(a, b) || a.hops.length - b.hops.length; + // Hashed only on a tie, so the common single-route prefix never pays for it. + if (order === 0) { + const ha = spreadHash(path, a.hops); + const hb = spreadHash(path, b.hops); + order = ha < hb ? -1 : ha > hb ? 1 : 0; + } if (order < 0) best = entry; } return best; @@ -241,7 +255,7 @@ class OriginState { const local = new Map(); const available = new Map(); for (const [path, routes] of this.routes.peek() ?? []) { - const entry = preferredEntry(routes); + const entry = preferredEntry(path, routes); if (!entry) continue; const value = { identity: entry.identity, route: entry.route.peek() }; remote.set(path, value); @@ -249,7 +263,7 @@ class OriginState { } for (const [path, front] of this.local.peek() ?? []) { const routes = this.routes.peek()?.get(path); - if (!this.localWins(path, routes && preferredEntry(routes))) continue; + if (!this.localWins(path, routes && preferredEntry(path, routes))) continue; const value = { identity: front, route: this.advertisedLocal.peek()?.get(path) ?? Route.default }; local.set(path, value); available.set(path, value.route); @@ -343,14 +357,14 @@ class OriginState { } const next = new Map(); for (const [prefix, entries] of routes ?? []) { - const mine = preferredEntry(entries, received); + const mine = preferredEntry(prefix, entries, received); if (mine) next.set(prefix, { identity: mine.identity, route: mine.route.peek() }); } // A local broadcast and an originated dynamic at one path compete on cost, as they do for requests. for (const [path, route] of advertised ?? []) { const front = local?.get(path); const entries = routes?.get(path); - if (front && this.localWins(path, entries && preferredEntry(entries, received))) { + if (front && this.localWins(path, entries && preferredEntry(path, entries, received))) { next.set(path, { identity: front, route }); } } @@ -375,7 +389,7 @@ class OriginState { let best: RouteEntry | undefined; for (const [prefix, entries] of this.routes.peek() ?? []) { if (!Path.hasPrefix(prefix, path)) continue; - const entry = preferredEntry(entries, skip); + const entry = preferredEntry(path, entries, skip); if (!entry) continue; if (bestPrefix === undefined || prefix.length > bestPrefix.length) { bestPrefix = prefix; diff --git a/js/net/src/wire.ts b/js/net/src/wire.ts index 3348b172bc..7a38ef5253 100644 --- a/js/net/src/wire.ts +++ b/js/net/src/wire.ts @@ -10,7 +10,7 @@ import type { Dispose, Getter } from "@moq/signals"; import type * as broadcast from "./broadcast.ts"; import type { Consumer as GroupConsumer } from "./group.ts"; -import type { Route } from "./hop.ts"; +import type { Hop, Route } from "./hop.ts"; import type * as origin from "./origin.ts"; import type * as Path from "./path.ts"; import type * as track from "./track.ts"; @@ -21,6 +21,13 @@ export interface Broadcast { resolveTrackInfo(name: string): Promise; fetchGroup(name: string, sequence: number, options?: track.FetchGroupOptions): Promise; requested(): Promise; + /** + * The origin serving this broadcast, named in SUBSCRIBE_OK and FETCH_OK: the one an + * upstream reply named, or else a random one generated once for the broadcast. + */ + origin(): Hop; + /** Record the origin an upstream reply named for this broadcast. */ + name(origin: Hop): void; } /** The protocol-facing operations behind an origin producer. */ diff --git a/quest/m1/README.md b/quest/m1/README.md index 88edaefce0..9ed40f63aa 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -49,6 +49,8 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. - [JavaScript FETCH](/quest/m1/js-fetch.md) - generic on-demand group serving and IETF FETCH for browser publishers - [Archive](/quest/m1/archive/README.md) - record selected tracks to any object_store and replay them over FETCH or derived HLS, on the catalog and store the release ships - [Wildcard](/quest/m1/wildcard/README.md) - a relay resolves subscriptions against advertised prefixes, a service claims the prefix it could serve and refuses the rest instead of enumerating broadcasts, and the browser player treats a covering claim as availability +- [Cluster origin reply](/quest/m1/cluster-origin.md) - a moq-transport downstream of a spreading relay never splices two pool members, and the cluster draft says how it tells them apart +- [Verified route upgrade](/quest/m1/front-upgrade.md) - a relay moves a live subscription to a cheaper route once that route's reply names the same origin - [Tooling](/quest/m1/tooling/README.md) - justfiles become a one-line menu over `sh/`, one impact map scopes CI, and every workflow step runs a recipe - [Path patterns](/quest/m1/path-patterns.md) - one matcher for every predicate over broadcast paths: tokens, origins, interest - [Setup token](/quest/m1/setup-token.md) - a moq-transport SETUP `AUTHORIZATION TOKEN` reaches the accepted handshake and the relay's auth request, so a verifier can run on it diff --git a/quest/m1/cluster-origin.md b/quest/m1/cluster-origin.md new file mode 100644 index 0000000000..9ef7c80b50 --- /dev/null +++ b/quest/m1/cluster-origin.md @@ -0,0 +1,38 @@ +# [M] Cluster origin reply + +## Goal + +A relay that spreads a pool's paths never lets a moq-transport downstream +splice one pool member's content onto another's, and +`draft-lcurley-moq-cluster` says how a downstream learns which member serves a +path. + +## Plan + +Spreading keys route selection's hash on the requested path, so equal-cost +advertisers of one prefix share its paths, and a relay advertises one route for +the whole pool. On moq-lite-07 the SUBSCRIBE_OK and FETCH_OK replies name the +serving origin, and a relay splices a failover only between replies naming the +same one. The cluster extension (moq-transport 17+) has no such field, so a +downstream relay still pins failover to the advertisement's first Hop ID, which +labels the pool rather than the member serving the path. The draft's +"Several Publishers of One Namespace" section still requires that the first +Hop ID downstream name the publisher whose Objects flow, which a spreading +relay breaks. + +Two ways out, to settle before implementing: + +- Name the serving origin in a moq-transport reply as a cluster-extension + parameter, mirroring lite-07, and stitch on it. Recommended: one identity + model on every wire. +- Keep the label truthful instead: toward a downstream that cannot learn the + origin, do not spread (or advertise per path). + +Either way the test is the one lite-07 has: two workers behind a pool relay +behind a moq-transport downstream relay, and killing the serving worker ends +the downstream subscription rather than splicing the survivor's objects. + +## Related + +- [Wildcard](/quest/m1/wildcard/README.md) - spreading and the lite-07 reply + origin landed there diff --git a/quest/m1/front-upgrade.md b/quest/m1/front-upgrade.md new file mode 100644 index 0000000000..84b0806ddd --- /dev/null +++ b/quest/m1/front-upgrade.md @@ -0,0 +1,37 @@ +# [M] Verified route upgrade + +## Goal + +A relay whose front learned its origin from a reply moves a live subscription +to a cheaper route once that route's source proves it serves the same origin, +instead of staying on the costlier route until it fails. + +## Plan + +On moq-lite-07 a front's identity is the origin its source's SUBSCRIBE_OK or +FETCH_OK names, not the route's first hop, because one route can lead to many +pool members. The front only learns a new route's origin once that source +replies, so today it stays on its live source when a better route appears +(`Pin::Stay` in `rs/moq-net/src/model/front.rs`) and only moves on failover. +That keeps content from being spliced across origins but leaves traffic on a +route that may no longer be the cheapest. + +The move to make: when a cheaper route appears, ask it before letting go of the +live source, and splice over at a group boundary only once its reply names the +front's origin. A reply naming another origin leaves the front where it is, +with nothing delivered from the candidate. Things to watch: + +- The candidate's content is held until admitted (the copy's provenance), so + asking early must not deliver anything or start demand the live source + already covers longer than needed. +- Churn: a route that keeps flapping should not keep opening and dropping + upstream subscriptions. +- Measure the cost of a verification round trip against staying put, and + benchmark it across routes and subscribers. + +## Related + +- [Wildcard](/quest/m1/wildcard/README.md) - the reply-named identity landed + there +- [Cluster origin reply](/quest/m1/cluster-origin.md) - the same identity on + moq-transport diff --git a/quest/m1/wildcard/README.md b/quest/m1/wildcard/README.md index 6f8ba90684..86ccf71609 100644 --- a/quest/m1/wildcard/README.md +++ b/quest/m1/wildcard/README.md @@ -1,4 +1,4 @@ -# Wildcard advertisements +# [S] Wildcard advertisements ## Goal @@ -40,7 +40,6 @@ claims resolved against pattern interest. The three workloads above still hold: the transcoder claims its service prefix ([Where derived output lives](#where-derived-output-lives)), and the archive claims the root and refuses what it does not have. -Spread and Demand are additive and land on main. ### What already exists, and what does not @@ -62,10 +61,11 @@ interest locally. Request resolution exists too. `Consumer::request_broadcast` mints a front per path (`rs/moq-net/src/model/front.rs`) that selects through `best_route`: a local broadcast first, then the longest covering prefix, filtered by the -requester's excluded hop and ordered by `route_order`. A refusal from that tier -is final, a front resumes only through routes sharing its first hop, and FETCH -resolves the same way. What remains is spreading one prefix's pool across -requested paths ([Spread](/quest/m1/wildcard/spread.md)). The pattern matcher itself exists: +requester's excluded hop and ordered by `route_order`, whose hash is keyed on +the requested path so one prefix's pool shares its paths. A refusal from that +tier is final, a front resumes only onto a source whose SUBSCRIBE_OK names the +same origin (the route's first hop on wires older than lite-07), and FETCH +resolves the same way. The pattern matcher itself exists: `moq_net::{Pattern, Patterns, Segment}` and `Path.Pattern` / `Path.Patterns` in `js/net/src/path.ts` own the shared matching, containment, specificity, and rebasing tokens and filters reuse. @@ -75,9 +75,10 @@ Announcement `Epoch` was specified into lite-06 by [#2611](https://github.com/moq-dev/moq/pull/2611), never implemented, and removed from the draft by #3225, which retired `draft-lcurley-moq-broadcast` with it. [moq#3312](https://github.com/moq-dev/moq/pull/3312) restored per-path identity -from the route's first hop, reversing #3225's no-splice rule, and this -questline builds its collision handling on that rather than on a generation -field. +from the route's first hop, reversing #3225's no-splice rule, and lite-07 moved +it to the origin a SUBSCRIBE_OK or FETCH_OK names, since a pool's one route labels +many origins. This questline builds its collision handling on that rather than +on a generation field. ### Decisions @@ -149,7 +150,8 @@ field. can hash one path to different workers before either concrete announcement propagates, and both land at the SAME literal path. Whichever route wins selection serves it, and a consumer moves between them only when the winner's - identity is preserved, per the first-hop resume rule ([moq#3312](https://github.com/moq-dev/moq/pull/3312)); two distinct workers are two identities, + identity is preserved, per the resume rule (the origin a SUBSCRIBE_OK + names on lite-07, the route's first hop before it); two distinct workers are two identities, so the loser's subscribers end and resubscribe rather than being spliced onto another worker's frames mid-group. Claim routing invents neither a lease nor a generation. This is weaker than the retired `Epoch` design, which could @@ -181,11 +183,6 @@ announcement shadows it. A claim names no generation, so a client that must distinguish recording generations reads the catalog's archive entry ([archive](/quest/m1/archive/README.md)) rather than announce state. -## Quests - -- [Spread](/quest/m1/wildcard/spread.md) - equal-cost advertisers of one - prefix share its paths instead of the first one taking them all - ## Related - [path-patterns](/quest/m1/path-patterns.md) - owns the pattern dialect diff --git a/quest/m1/wildcard/spread.md b/quest/m1/wildcard/spread.md deleted file mode 100644 index 963b312858..0000000000 --- a/quest/m1/wildcard/spread.md +++ /dev/null @@ -1,46 +0,0 @@ -# [M] Spread - -## Goal - -Equal-cost advertisers of one prefix share its paths: a relay spreads distinct -requested paths across the pool, and one path always resolves to the same -advertiser. A transcode pool claiming the root today sends every job to one -worker until its cost changes. - -## Plan - -`best_route` (`rs/moq-net/src/model/origin.rs`) orders the winning tier with -`route_order`, whose hash tie-break is keyed on the advertised prefix. Every -path under a pool's shared prefix hashes the same, so one worker takes them -all. Keying that hash on the requested path spreads them; cost still orders -first, so a distant worker stays overflow rather than an equal peer. The lite -draft's Routing tie-breaks name no hash, so spell this one there too. - -Decided: stitching identity comes from the reply, not the announced route. -A relay advertises one best route per prefix to each peer, and today the peer -pins a front to that route's first hop (moq#3312). Once a relay serves a path -from a different pool member than the one it advertised, that label is wrong, -and a later failover through another route with the same first hop splices -different content. The mismatch already exists narrowly (a front stays pinned -after its prefix's best route changes), and spreading makes it the common -case. - -- The subscribe and fetch replies name the origin that actually serves the - request, and a relay stitches a failover only between replies naming the - same origin. Differing origins end the subscription and the subscriber - re-requests. Where the field sits (SUBSCRIBE_OK, TRACK_INFO, or the fetch - reply) is the implementer's call; it lands in lite-07 (`moq-lite-07-wip`) - and the lite draft, and older versions keep today's first-hop pinning. -- Rejected: re-originating the prefix at the pool's relay (identity names the - relay, but a pool membership change re-hashes under the same identity), and - documenting that pool members must be interchangeable (independent encoders - are not). -- Whether the hop list is needed at all once identity moves to the reply is - the m2 plan quest added in moq#4158; this - quest does not wait on it. - -Tests: one path always selects the same advertiser; a fixed set of many paths -spreads across advertisers rather than piling onto one (do not assert two -particular paths differ, which a correct hash may violate); and whichever -option lands, a downstream failover never splices one pool member's content -onto another's. diff --git a/rs/moq-gst/src/sink/pad.rs b/rs/moq-gst/src/sink/pad.rs index 23ecdb1a33..3a6afa5b05 100644 --- a/rs/moq-gst/src/sink/pad.rs +++ b/rs/moq-gst/src/sink/pad.rs @@ -29,12 +29,12 @@ enum PadState { /// Where a pad's buffers land: a codec importer, a subtitle track, or opaque data. /// -/// Both payloads are large (a codec importer, a container producer), so each is boxed to keep the -/// enum small. +/// Every payload is large (a codec importer, a container producer, a track producer), so +/// each is boxed to keep the enum small. enum Sink { Media(Box), Text(Box), - Opaque(moq_net::track::Producer), + Opaque(Box), } /// An audio or video pad, published through a codec importer. @@ -303,7 +303,7 @@ impl Pad { .with_context(|| format!("cannot reserve track {name}"))?; // Followed at the live edge, so it keeps the default retention the media helper raises. let info = moq_net::track::Info::default().with_timescale(moq_net::Timescale::MICRO); - self.track = Some(Sink::Opaque(request.accept(info))); + self.track = Some(Sink::Opaque(Box::new(request.accept(info)))); self.caps = Some(caps.clone()); return Ok(name); } diff --git a/rs/moq-net/src/lite/fetch.rs b/rs/moq-net/src/lite/fetch.rs index c098c5990a..c0c3b6d44e 100644 --- a/rs/moq-net/src/lite/fetch.rs +++ b/rs/moq-net/src/lite/fetch.rs @@ -1,7 +1,7 @@ use std::borrow::Cow; use crate::{ - Path, + Hop, Path, coding::{Decode, DecodeError, Encode, EncodeError}, }; @@ -81,6 +81,33 @@ impl Message for Fetch<'_> { } } +/// The publisher's answer on a Fetch Stream, ahead of the FRAME messages: the +/// origin serving the group, which a relay stitches failover on. +/// +/// Lite07+ only; older versions answer with the frames alone. +#[derive(Clone, Debug)] +pub struct FetchOk { + /// [`Hop::UNKNOWN`] names nobody. + pub origin: Hop, +} + +impl Message for FetchOk { + fn decode_msg(r: &mut R, version: Version) -> Result { + if !version.has_origin() { + return Err(DecodeError::Version); + } + let origin = Hop::from_wire(u64::decode(r, version)?)?; + Ok(Self { origin }) + } + + fn encode_msg(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { + if !version.has_origin() { + return Err(EncodeError::Version); + } + self.origin.id().encode(w, version) + } +} + #[cfg(test)] mod test { use super::*; @@ -103,6 +130,21 @@ mod test { Fetch::decode_msg(&mut slice, version).unwrap() } + #[test] + fn fetch_ok_names_the_origin_on_lite07() { + let msg = FetchOk { + origin: Hop::new(42).unwrap(), + }; + let mut buf = Vec::new(); + msg.encode(&mut buf, Version::Lite07).unwrap(); + assert_eq!(buf, [1, 42]); + let got = FetchOk::decode(&mut buf.as_slice(), Version::Lite07).unwrap(); + assert_eq!(got.origin, msg.origin); + + // Older versions answer with the frames alone. + assert!(msg.encode(&mut Vec::new(), Version::Lite06).is_err()); + } + #[test] fn fetch_roundtrips() { for version in [Version::Lite03, Version::Lite04, Version::Lite05] { diff --git a/rs/moq-net/src/lite/publisher.rs b/rs/moq-net/src/lite/publisher.rs index 668ca43a1b..ba7cec00a4 100644 --- a/rs/moq-net/src/lite/publisher.rs +++ b/rs/moq-net/src/lite/publisher.rs @@ -1,5 +1,5 @@ use crate::runtime::Timers as _; -use crate::{SessionError, announce, frame, group, origin, track}; +use crate::{SessionError, announce, broadcast, frame, group, origin, track}; use std::{ collections::HashMap, ops::Bound, @@ -1008,6 +1008,7 @@ enum SubscribeState { /// Waiting for the model subscription to be confirmed. Confirm { msg: lite::Subscribe<'static>, + broadcast: broadcast::Consumer, subscribing: track::Subscribing, }, /// Streaming groups and datagrams. Boxed: by far the largest state, and the enum @@ -1112,11 +1113,15 @@ impl SubscribeServe { // duplicate demand). let track_consumer = broadcast.track(&msg.track)?; let subscribing = track_consumer.subscribe(subscription).into_inner(); - self.state = SubscribeState::Confirm { msg, subscribing }; + self.state = SubscribeState::Confirm { + msg, + broadcast, + subscribing, + }; } SubscribeState::Confirm { subscribing, .. } => { let track = ready!(subscribing.poll_ok(waiter))?; - let SubscribeState::Confirm { msg, .. } = + let SubscribeState::Confirm { msg, broadcast, .. } = std::mem::replace(&mut self.state, SubscribeState::Decode) else { unreachable!() @@ -1164,6 +1169,10 @@ impl SubscribeServe { version: self.shared.version, timescale, opens: Default::default(), + served: Served { + broadcast: Some(broadcast), + here: self.shared.self_origin, + }, }; let run = TrackRun::new(sub, track, Bounds::from(&msg), track_priority_tx); @@ -1235,6 +1244,7 @@ enum FetchState { /// Waiting for the fetched group. Fetch { msg: lite::Fetch<'static>, + broadcast: broadcast::Consumer, fetching: track::Fetching, }, /// Streaming the group's frames in order. The delta-timestamp baseline @@ -1339,12 +1349,34 @@ impl FetchServe { }, ) .into_inner(); - self.state = FetchState::Fetch { msg, fetching }; + self.state = FetchState::Fetch { + msg, + broadcast, + fetching, + }; } - FetchState::Fetch { msg, fetching } => { + FetchState::Fetch { + msg, + broadcast, + fetching, + } => { let mut group = ready!(kio::Pollable::poll(fetching, waiter))?; - // The response carries no header, so a short run is indistinguishable + // Lite-07 names the origin serving the group ahead of its frames, read + // now that the group resolved (its content flowed, so the front's + // origin is settled). + if self.shared.version.has_origin() { + let served = Served { + broadcast: Some(broadcast.clone()), + here: self.shared.self_origin, + }; + let stream = self.stream.as_mut().expect("stream present"); + stream.writer.buffer(&lite::FetchOk { + origin: served.origin(), + })?; + } + + // The response carries no position, so a short run is indistinguishable // from one that started elsewhere: only serve a range we can cover // exactly. `fetch_group` already positions the consumer, so this is a // belt-and-braces check on a promise the wire can't restate. @@ -2133,6 +2165,37 @@ struct Subscription { timescale: Option, /// The group streams this subscription opened, shared by every group it serves. opens: Arc, + /// Who SUBSCRIBE_OK names as serving it (lite-07+). + served: Served, +} + +/// Names the origin serving a request, for SUBSCRIBE_OK and FETCH_OK: the one the +/// broadcast's front serves, read when the reply goes out (a copy's content only +/// reaches the front once its origin is admitted), or ours for a broadcast that is +/// not route-fed. +#[derive(Clone)] +struct Served { + broadcast: Option, + here: Hop, +} + +#[cfg(test)] +impl Default for Served { + fn default() -> Self { + Self { + broadcast: None, + here: Hop::UNKNOWN, + } + } +} + +impl Served { + fn origin(&self) -> Hop { + self.broadcast + .as_ref() + .and_then(|broadcast| broadcast.origin()) + .unwrap_or(self.here) + } } /// Counts a subscription's group streams for lite-07's SUBSCRIBE_END. @@ -2326,22 +2389,7 @@ impl TrackRun { tracing::debug!(subscribe = self.ctx.id, track = %self.ctx.track_name, sequence, "skipping group with a missing head"); continue; } - if self.emit_range && !self.start_sent { - self.start_sent = true; - // Only the group: the subscriber derives the start frame from - // its own request (see `lite::SubscribeStart`). - stream - .writer - .buffer(&lite::SubscribeResponse::Start(lite::SubscribeStart { - group: sequence, - }))?; - // SUBSCRIBE_OK is an implicit drop of everything below the - // resolved start (the subscriber records it as a permanent - // miss), so a lower group arriving late must not be served - // after all. A widening SUBSCRIBE_UPDATE re-lowers the floor, - // renegotiating the resolved start along with the demand. - self.track.start_at(sequence); - } + self.start(stream, sequence)?; let frame_start = group.index(); tracing::debug!(subscribe = self.ctx.id, track = %self.ctx.track_name, sequence, "serving group"); @@ -2358,7 +2406,10 @@ impl TrackRun { self.children .push(GroupServe::new(self.ctx.clone(), sequence, frame_start, handle, group)); } - Recv::Datagram(datagram) => self.ctx.serve_datagram(datagram), + Recv::Datagram(datagram) => { + self.start(stream, datagram.sequence)?; + self.ctx.serve_datagram(datagram); + } Recv::Boundary(group) => { // The track declared its exclusive final sequence. Forward it now, // even if trailing groups (below `group`) are still in flight, then @@ -2391,6 +2442,30 @@ impl TrackRun { return Poll::Pending; } } + + /// Send SUBSCRIBE_START ahead of the first group served, by stream or datagram: it + /// resolves the start and names the origin, which a lite-07 subscriber needs before + /// it delivers either. + fn start(&mut self, stream: &mut Stream, sequence: u64) -> Result<(), Error> { + if !self.emit_range || self.start_sent { + return Ok(()); + } + self.start_sent = true; + // Only the group: the subscriber derives the start frame from its own request + // (see `lite::SubscribeStart`). + stream + .writer + .buffer(&lite::SubscribeResponse::Start(lite::SubscribeStart { + group: sequence, + origin: self.ctx.served.origin(), + }))?; + // SUBSCRIBE_OK is an implicit drop of everything below the resolved start (the + // subscriber records it as a permanent miss), so a lower group arriving late + // must not be served after all. A widening SUBSCRIBE_UPDATE re-lowers the floor, + // renegotiating the resolved start along with the demand. + self.track.start_at(sequence); + Ok(()) + } } /// Serves one group on its own unidirectional stream: the header, then every @@ -2810,6 +2885,7 @@ mod serve_group_test { version: Version::Lite06, timescale: Some(crate::Timescale::default()), opens: Default::default(), + served: Served::default(), }; let track = track::Producer::new(Arc::new(broadcast::Info::default()), "test", None); @@ -2851,6 +2927,7 @@ mod serve_group_test { version: Version::Lite06, timescale: Some(crate::Timescale::default()), opens: Default::default(), + served: Served::default(), }; let track = track::Producer::new(Arc::new(broadcast::Info::default()), "test", None); @@ -2896,6 +2973,7 @@ mod serve_group_test { version: Version::Lite06, timescale: Some(crate::Timescale::default()), opens: Default::default(), + served: Served::default(), }; let track = track::Producer::new(Arc::new(broadcast::Info::default()), "test", None); @@ -2959,6 +3037,7 @@ mod serve_group_test { version: Version::Lite06, timescale: Some(crate::Timescale::default()), opens: Default::default(), + served: Served::default(), }; let track = track::Producer::new(Arc::new(broadcast::Info::default()), "test", None); @@ -3029,6 +3108,7 @@ mod serve_group_test { version: Version::Lite06, timescale: Some(crate::Timescale::default()), opens: Default::default(), + served: Served::default(), }; let track = track::Producer::new(Arc::new(broadcast::Info::default()), "test", None); @@ -3075,6 +3155,11 @@ mod serve_group_test { version: Version::Lite07, timescale: Some(crate::Timescale::default()), opens: Default::default(), + // Content originating here: SUBSCRIBE_START names our own hop. + served: Served { + broadcast: None, + here: Hop::new(9).unwrap(), + }, }; let bounds = Bounds { start_group: Some(0), @@ -3108,8 +3193,8 @@ mod serve_group_test { track.finish().unwrap(); assert!(matches!(run.await.unwrap(), TrackEnd::Finished)); - // SUBSCRIBE_START at 0, then SUBSCRIBE_END at 3 with 2 streams. - assert_eq!(*log.writes.lock().unwrap(), [0, 1, 0, 1, 2, 3, 2]); + // SUBSCRIBE_START at 0 from origin 9, then SUBSCRIBE_END at 3 with 2 streams. + assert_eq!(*log.writes.lock().unwrap(), [0, 2, 0, 9, 1, 2, 3, 2]); } /// A track that ends without a group still ends the subscription, with no stream owed. @@ -3145,7 +3230,7 @@ mod serve_group_test { write_group(&mut track, 1, 1000); track.finish().unwrap(); assert!(futures::poll!(run.as_mut()).is_pending(), "group 1 is still opening"); - assert_eq!(*log.writes.lock().unwrap(), [0, 1, 0], "only SUBSCRIBE_START so far"); + assert_eq!(*log.writes.lock().unwrap(), [0, 2, 0, 9], "only SUBSCRIBE_START so far"); let Ok(mut open) = gate.write() else { panic!("transport gate closed"); @@ -3153,7 +3238,7 @@ mod serve_group_test { *open = true; drop(open); run.await.unwrap(); - assert_eq!(*log.writes.lock().unwrap(), [0, 1, 0, 1, 2, 2, 1]); + assert_eq!(*log.writes.lock().unwrap(), [0, 2, 0, 9, 1, 2, 2, 1]); } } diff --git a/rs/moq-net/src/lite/subscribe.rs b/rs/moq-net/src/lite/subscribe.rs index 4656e94521..a3513f8019 100644 --- a/rs/moq-net/src/lite/subscribe.rs +++ b/rs/moq-net/src/lite/subscribe.rs @@ -1,7 +1,7 @@ use std::borrow::Cow; use crate::{ - Path, + Hop, Path, coding::{Decode, DecodeError, Encode, EncodeError, Sizer}, }; @@ -295,6 +295,9 @@ impl Message for SubscribeOk { #[derive(Clone, Debug)] pub struct SubscribeStart { pub group: u64, + /// The origin serving the subscription: the Hop ID a relay stitches failover on. + /// [`Hop::UNKNOWN`] names nobody. Lite07+; older versions decode it as unknown. + pub origin: Hop, } impl Message for SubscribeStart { @@ -302,16 +305,23 @@ impl Message for SubscribeStart { if !version.has_track_stream() { return Err(DecodeError::Version); } - Ok(Self { - group: u64::decode(r, version)?, - }) + let group = u64::decode(r, version)?; + let origin = match version.has_origin() { + true => Hop::from_wire(u64::decode(r, version)?)?, + false => Hop::UNKNOWN, + }; + Ok(Self { group, origin }) } fn encode_msg(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { if !version.has_track_stream() { return Err(EncodeError::Version); } - self.group.encode(w, version) + self.group.encode(w, version)?; + if version.has_origin() { + self.origin.id().encode(w, version)?; + } + Ok(()) } } @@ -578,7 +588,10 @@ mod test { #[test] fn subscribe_start_roundtrips_on_lite05() { - let resp = SubscribeResponse::Start(SubscribeStart { group: 42 }); + let resp = SubscribeResponse::Start(SubscribeStart { + group: 42, + origin: Hop::UNKNOWN, + }); let mut buf = Vec::new(); resp.encode(&mut buf, Version::Lite05).unwrap(); let mut slice = buf.as_slice(); @@ -602,6 +615,31 @@ mod test { } } + #[test] + fn subscribe_start_carries_the_origin_on_lite07() { + let resp = SubscribeResponse::Start(SubscribeStart { + group: 7, + origin: Hop::new(42).unwrap(), + }); + let mut buf = Vec::new(); + resp.encode(&mut buf, Version::Lite07).unwrap(); + // Type, length, group, origin. + assert_eq!(buf, [0, 2, 7, 42]); + match SubscribeResponse::decode(&mut buf.as_slice(), Version::Lite07).unwrap() { + SubscribeResponse::Start(start) => assert_eq!((start.group, start.origin), (7, Hop::new(42).unwrap())), + other => panic!("expected Start, got {other:?}"), + } + + // Older versions have no room for it. + let mut buf = Vec::new(); + resp.encode(&mut buf, Version::Lite06).unwrap(); + assert_eq!(buf, [0, 1, 7]); + match SubscribeResponse::decode(&mut buf.as_slice(), Version::Lite06).unwrap() { + SubscribeResponse::Start(start) => assert_eq!(start.origin, Hop::UNKNOWN), + other => panic!("expected Start, got {other:?}"), + } + } + #[test] fn subscribe_end_carries_the_stream_count_on_lite07() { let resp = SubscribeResponse::End(SubscribeEnd { group: 7, streams: 3 }); diff --git a/rs/moq-net/src/lite/subscriber.rs b/rs/moq-net/src/lite/subscriber.rs index 20987dbbc7..edd22f8f5b 100644 --- a/rs/moq-net/src/lite/subscriber.rs +++ b/rs/moq-net/src/lite/subscriber.rs @@ -90,6 +90,10 @@ struct TrackEntry { timescale: Option, /// The groups received so far, so the subscription's end can wait for the ones owed. tail: kio::Producer, + /// Whether this subscription's SUBSCRIBE_OK arrived. On a wire that names the + /// serving origin there, group streams wait for it: until then nobody knows + /// whose content they carry. + started: kio::Shared, } impl Subscriber { @@ -434,6 +438,12 @@ impl Subscriber { return Ok(()); }; + // A datagram cannot wait for its subscription's SUBSCRIBE_OK like a group stream + // does, so until the origin it names is admitted, the datagram is dropped. + if self.version.has_origin() && !entry.producer.provenance().is_admitted() { + return Ok(()); + } + // Datagrams are lite-05+, which always negotiates a timescale; default defensively. let scale = entry.timescale.unwrap_or_default(); let timestamp = @@ -695,6 +705,11 @@ struct GroupRecv { enum GroupRecvState { /// Reading the GROUP header. Header, + /// Holding the group until its subscription's SUBSCRIBE_OK names the serving origin + /// and the front admits it, on a wire that names one. + Hold { + header: lite::Group, + }, /// Filling the group, bailing if the track or group dies first. Serve { /// Guarded: dropping this machine mid-group is a cancellation, not a clean end. @@ -719,7 +734,31 @@ impl GroupRecv { match &mut self.state { GroupRecvState::Header => { let mut cx = std::task::Context::from_waker(waiter.waker()); - let hdr = ready!(self.reader.poll_decode::(&mut cx))?; + let header = ready!(self.reader.poll_decode::(&mut cx))?; + self.state = GroupRecvState::Hold { header }; + } + GroupRecvState::Hold { header } => { + if self.subscriber.version.has_origin() { + let entry = self + .subscriber + .subscribes + .lock() + .get(&header.subscribe) + .cloned() + .ok_or(Error::Cancel)?; + if let Poll::Ready(err) = entry.producer.poll_closed(waiter) { + return Poll::Ready(Err(err)); + } + ready!(entry.started.poll(waiter, |started| match **started { + true => Poll::Ready(()), + false => Poll::Pending, + })); + ready!(entry.producer.provenance().poll_admitted(waiter))?; + } + let GroupRecvState::Hold { header: hdr } = std::mem::replace(&mut self.state, GroupRecvState::Done) + else { + unreachable!() + }; let (group, track, timescale) = { let mut subs = self.subscriber.subscribes.lock(); @@ -1369,6 +1408,7 @@ mod tests { producer, timescale: Some(Timescale::default()), tail: Default::default(), + started: kio::Shared::new(true), }, ); @@ -1405,6 +1445,65 @@ mod tests { ); } + /// On lite-07 a datagram carries no origin of its own, so until its subscription's + /// origin is admitted it is dropped rather than risk splicing another origin's content. + #[test] + fn datagram_waits_for_the_admitted_origin() { + let version = Version::Lite07; + let origin = origin::Config::new(crate::Hop::new(1).unwrap()).produce(); + let subscriber = Subscriber::new(SubscriberConfig { + runtime: crate::time::Clock::tokio(), + session: SinkSession::default(), + origin, + recv_bandwidth: None, + version, + peer_setup: Default::default(), + peer_hop: None, + cost: None, + going_away: Default::default(), + }); + + let broadcast = crate::broadcast::Info::new().produce(); + let producer = broadcast.create_track("datagrams", None).unwrap(); + let provenance = producer.provenance(); + let mut received = producer.subscribe(None); + subscriber.subscribes.lock().insert( + 7, + TrackEntry { + producer, + timescale: Some(Timescale::default()), + tail: Default::default(), + started: kio::Shared::new(false), + }, + ); + + let payload = |sequence| { + lite::Datagram { + subscribe: 7, + sequence, + timestamp: sequence, + payload: bytes::Bytes::from_static(b"x"), + } + .encode_bytes(version) + .unwrap() + }; + + // Named but not admitted: the front serves another origin. + let named = crate::Hop::new(42).unwrap(); + provenance.name(named).unwrap(); + provenance.admit(crate::Hop::new(43).unwrap()); + subscriber.route_datagram(payload(1)).unwrap(); + assert!( + received.recv_datagram().now_or_never().is_none(), + "delivered before admission" + ); + + provenance.admit(named); + subscriber.route_datagram(payload(2)).unwrap(); + let datagram = received.recv_datagram().now_or_never().unwrap().unwrap().unwrap(); + assert_eq!(datagram.sequence, 2); + } + /// `establish` puts exactly one SUBSCRIBE on the wire, and the id is registered /// before any of it reaches the transport. /// @@ -2543,6 +2642,8 @@ struct SubStream { requested: Option, /// The groups received for this subscription, shared with its [`TrackEntry`]. tail: kio::Producer, + /// Whether SUBSCRIBE_OK arrived, shared with its [`TrackEntry`]. + started: kio::Shared, /// The first group the publisher serves (SUBSCRIBE_START), once declared. served: Option, /// The track's exclusive end (SUBSCRIBE_END), once declared. @@ -2912,12 +3013,14 @@ impl TrackServe { tracing::info!(id, broadcast = %self.subscriber.log_path(&self.path), track = %self.name, "subscribe started"); let tail = kio::Producer::new(Tail::default()); + let started = kio::Shared::new(!self.subscriber.version.has_origin()); self.subscriber.subscribes.lock().insert( id, TrackEntry { producer: producer.clone(), timescale, tail: tail.clone(), + started: started.clone(), }, ); @@ -2929,6 +3032,7 @@ impl TrackServe { id, subscription, tail, + started, state: EstablishState::Open, } } @@ -3029,6 +3133,7 @@ struct Establish { id: u64, subscription: Subscription, tail: kio::Producer, + started: kio::Shared, state: EstablishState, } @@ -3114,6 +3219,7 @@ impl Establish { priority: self.subscription.priority, requested: self.subscription.start, tail: self.tail.clone(), + started: self.started.clone(), served: None, end: None, } @@ -3284,10 +3390,11 @@ impl TrackInfoFetch { // window matches what the upstream advertises (relays re-serve with // the same bound). `broadcast` is left at its default here; // `track::Request::accept` stamps the track's real broadcast. - let model = track::Info::default() + let mut model = track::Info::default() .with_timescale(info.timescale) .with_max_age(info.max_age) .with_priority(info.priority); + model.names_origin = serve.subscriber.version.has_origin(); return Poll::Ready(Ok(model)); } } @@ -3425,8 +3532,12 @@ impl ServeLoop { match self.dynamic.poll_requested_group(waiter) { Poll::Ready(Ok(req)) => { if self.supports_fetch { - self.fetches - .push(FetchServeRun::new(serve.clone(), req, self.timescale)); + self.fetches.push(FetchServeRun::new( + serve.clone(), + req, + self.timescale, + self.serving.provenance(), + )); } else { req.reject(Error::Version); } @@ -3504,6 +3615,18 @@ impl ServeLoop { // signal, so a spliced reader waiting on a skipped group // fails over instead of stalling on a live route. lite::SubscribeResponse::Start(start) => { + // The reply names whose content this subscription + // carries. A copy has one origin, so another one is a + // different track: hand it back. The held group + // streams go ahead, and still wait on the front + // admitting that origin. + if serve.subscriber.version.has_origin() { + if let Err(err) = self.serving.provenance().name(start.origin) { + tracing::debug!(track = %serve.name, origin = start.origin.id(), "subscription changed origin"); + return Poll::Ready(ServeEnd::GiveBack(err)); + } + *active.started.lock() = true; + } // A START describes the demand the SUBSCRIBE carried. // It applies only while the current start still matches // that demand (updates get no fresh START, so an update @@ -3586,9 +3709,14 @@ struct FetchServeRun { session: S, timescale: Option, group: u64, + /// The copy's origin, which FETCH_OK names on lite-07. + provenance: track::Provenance, state: FetchRunState, } +// A state machine's enum is its storage: one transient instance per stream, so the +// big variant is the working state, not padding held in bulk. +#[allow(clippy::large_enum_variant)] enum FetchRunState { Open { request: Option, @@ -3604,6 +3732,12 @@ enum FetchRunState { stream: Stream, frame_start: u64, }, + /// FETCH_OK named the origin: hold the frames until the front admits it. + Admit { + request: Option, + stream: Stream, + frame_start: u64, + }, Ingest { stream: Stream, producer: group::Producer, @@ -3613,7 +3747,12 @@ enum FetchRunState { } impl FetchServeRun { - fn new(serve: TrackServe, request: group::Request, timescale: Option) -> Self { + fn new( + serve: TrackServe, + request: group::Request, + timescale: Option, + provenance: track::Provenance, + ) -> Self { let session = serve.subscriber.session.clone(); let group = request.sequence(); Self { @@ -3621,6 +3760,7 @@ impl FetchServeRun { session, timescale, group, + provenance, state: FetchRunState::Open { request: Some(request) }, } } @@ -3721,11 +3861,16 @@ impl kio::Task for FetchServeRun { }; } FetchRunState::Answer { stream, .. } => { - // Lite has no FETCH_OK: a publisher without the group resets the - // stream instead. Accepting before the first byte (or a FIN, for an - // empty group) would resolve every joined `fetch_group` to a group - // that only fails on its first read, so wait for the answer. - let answered = ready!(stream.reader.poll_has_more(&mut cx)); + // A publisher without the group resets the stream. Accepting before + // its answer (FETCH_OK on lite-07, the first byte or a FIN before) + // would resolve every joined `fetch_group` to a group that only fails + // on its first read, so wait for it. FETCH_OK names whose content + // follows, which a copy only takes from one origin. + let answered = match self.serve.subscriber.version.has_origin() { + true => ready!(stream.reader.poll_decode::(&mut cx)) + .and_then(|ok| self.provenance.name(ok.origin)), + false => ready!(stream.reader.poll_has_more(&mut cx)).map(|_| ()), + }; let FetchRunState::Answer { request, stream, @@ -3734,10 +3879,35 @@ impl kio::Task for FetchServeRun { else { unreachable!() }; - let request = request.expect("request pending"); if let Err(err) = answered { tracing::debug!(track = %self.serve.name, group = self.group, %err, "fetch refused"); stream.writer.abort(&err); + request.expect("request pending").reject(err); + return Poll::Ready(()); + } + self.state = FetchRunState::Admit { + request, + stream, + frame_start, + }; + } + FetchRunState::Admit { .. } => { + let admitted = match self.serve.subscriber.version.has_origin() { + true => ready!(self.provenance.poll_admitted(waiter)), + false => Ok(()), + }; + let FetchRunState::Admit { + request, + stream, + frame_start, + } = std::mem::replace(&mut self.state, FetchRunState::Done) + else { + unreachable!() + }; + let request = request.expect("request pending"); + if let Err(err) = admitted { + tracing::debug!(track = %self.serve.name, group = self.group, %err, "fetch from another origin"); + stream.writer.abort(&err); request.reject(err); return Poll::Ready(()); } diff --git a/rs/moq-net/src/lite/version.rs b/rs/moq-net/src/lite/version.rs index 3c81f2a994..f560cbfe71 100644 --- a/rs/moq-net/src/lite/version.rs +++ b/rs/moq-net/src/lite/version.rs @@ -206,6 +206,18 @@ impl Version { } } + /// Whether SUBSCRIBE_OK and FETCH_OK name the origin serving the request, which is + /// what a relay stitches a failover on. Added in lite-07 (with FETCH_OK itself); + /// older versions leave a relay the route's first hop. + #[allow(clippy::match_like_matches_macro)] + pub(crate) fn has_origin(self) -> bool { + // Match form so future versions default forward (CLAUDE.md convention). + match self { + Self::Lite01 | Self::Lite02 | Self::Lite03 | Self::Lite04 | Self::Lite05 | Self::Lite06 => false, + _ => true, + } + } + /// Whether announcements carry the route cost: the marginal cost of pulling /// the broadcast via this route, accumulated per link. Added in lite-06. /// Older versions carry nothing, so a received route stays at zero and ranks diff --git a/rs/moq-net/src/model/broadcast.rs b/rs/moq-net/src/model/broadcast.rs index 7cdf0618eb..2e8211cd78 100644 --- a/rs/moq-net/src/model/broadcast.rs +++ b/rs/moq-net/src/model/broadcast.rs @@ -117,6 +117,10 @@ struct SplicedState { // Names awaiting assignment to a route, in request order. pending: VecDeque>, + + // The origin the front serves, as its replies name it. Set by the front as its + // identity settles; `None` until then. + origin: Option, } impl BroadcastState { @@ -413,6 +417,17 @@ impl Producer { Poll::Ready((name, producer)) } + /// Record the origin the front serves; see [`Consumer::origin`]. Route-fed + /// broadcasts only. + pub(crate) fn set_origin(&self, origin: crate::Hop) { + if self.state.read().spliced.as_ref().and_then(|spliced| spliced.origin) == Some(origin) { + return; + } + let mut state = self.state.lock(); + let spliced = state.spliced.as_mut().expect("origin of a route-fed broadcast"); + spliced.origin = Some(origin); + } + /// Let go of every spliced track, aborting with `err` the ones never handed /// out by [`Self::poll_spliced_assigned`]. Called when the broadcast ends: /// whoever took the others decides how they end. @@ -943,6 +958,12 @@ impl Consumer { self.is_closed() || self.state.read().closing } + /// The origin serving this broadcast, as a reply for it names it: `None` for a + /// broadcast that is not route-fed, whose content originates on this origin. + pub(crate) fn origin(&self) -> Option { + self.state.read().spliced.as_ref().and_then(|spliced| spliced.origin) + } + /// Whether the broadcast ended via a deliberate [`Producer::finish`], as opposed /// to aborting or losing its producer. `false` while the broadcast is still live. pub fn is_finished(&self) -> bool { diff --git a/rs/moq-net/src/model/front.rs b/rs/moq-net/src/model/front.rs index 0cc387086c..de1d8451b8 100644 --- a/rs/moq-net/src/model/front.rs +++ b/rs/moq-net/src/model/front.rs @@ -21,7 +21,8 @@ use std::{ use crate::{Error, Hop, runtime::Instant, track}; /// A route the table selected for the front: the entry id, the endpoint that -/// originated it, and whether it is a broadcast published on this origin. +/// originated it (`None` for an empty chain: announced on this origin), and +/// whether it is a broadcast published on this origin. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub(super) struct Candidate { pub route: u64, @@ -74,6 +75,9 @@ pub(super) enum Event { result: Result<(), Error>, delivered: bool, }, + /// A spliced copy's reply named the origin serving it. Its content is held + /// until the front [admits](Front::admit) that origin. + Origin { track: Arc, source: u64, origin: Hop }, /// A reader arrived on the track. Used { track: Arc }, /// The last reader left the track. @@ -131,13 +135,20 @@ pub(super) enum Identity { /// closing the front ends instead, so a newcomer gets a fresh broadcast /// rather than being spliced into one that is over. Local, - /// The serving route's first hop was absent or [`Hop::UNKNOWN`], which - /// identifies nobody: the front cannot resume through any other route, so - /// its source ending, or `route` leaving the table, ends it. - Anonymous { route: u64 }, - /// The first hop of the serving route. Routes sharing it are the same - /// origin reached another way and safe to resume through. + /// The serving route names no other origin: its chain is empty (a handler on + /// this origin, `here`) or its first hop is [`Hop::UNKNOWN`], which + /// identifies nobody. The front cannot resume through any other route, so its + /// source ending, or `route` leaving the table, ends it. + Anonymous { route: u64, here: bool }, + /// The first hop of the serving route, for sources whose replies name no + /// origin (older wires). Routes sharing it are the same origin reached + /// another way and safe to resume through. Publisher(Hop), + /// The origin the first source's replies named. Any route may lead back to it, + /// since an advertiser serving a prefix from a pool advertises one route for many + /// origins, so a replacement source is verified by its own replies instead of by + /// its route. + Origin(Hop), } /// Which routes qualify for a front's (re)selection. @@ -151,15 +162,19 @@ pub(super) enum Pin { Publisher(Hop), /// Only this route: the front never fails over. Route(u64), + /// Any served route, but this one while it stands: a replacement's content + /// is only known once it replies, so a live source is not traded for it. + Stay(u64), } impl Identity { - fn pin(self) -> Pin { + /// The origin the front's content comes from, as its own replies name it: + /// `None` for this origin, [`Hop::UNKNOWN`] for nobody identifiable. + fn origin(self) -> Option { match self { - Self::Undetermined => Pin::Any, - Self::Local => Pin::Local, - Self::Anonymous { route } => Pin::Route(route), - Self::Publisher(hop) => Pin::Publisher(hop), + Self::Local | Self::Anonymous { here: true, .. } => None, + Self::Undetermined | Self::Anonymous { here: false, .. } => Some(Hop::UNKNOWN), + Self::Publisher(hop) | Self::Origin(hop) => Some(hop), } } } @@ -195,6 +210,11 @@ struct Track { #[derive(Clone, Debug)] pub(super) struct Front { identity: Identity, + /// The origin the replies established, and the only one whose content is + /// admitted. `None` until a copy whose wire names origins replies. + named: Option, + /// The first source attached: its replies may name what the route only labelled. + first: Option, /// The attached source and the route that produced it. serving: Option<(u64, u64)>, /// Whether the serving source has begun closing, as of the last event that @@ -223,6 +243,8 @@ impl Front { pub(super) fn new(linger: Duration) -> Self { Self { identity: Identity::Undetermined, + named: None, + first: None, serving: None, serving_closing: false, upstream: None, @@ -238,7 +260,27 @@ impl Front { /// Which routes qualify for the next selection. pub(super) fn pin(&self) -> Pin { - self.identity.pin() + match self.identity { + Identity::Undetermined => Pin::Any, + Identity::Local => Pin::Local, + Identity::Anonymous { route, .. } => Pin::Route(route), + Identity::Publisher(hop) => Pin::Publisher(hop), + Identity::Origin(_) => match self.serving { + Some((_, route)) => Pin::Stay(route), + None => Pin::Any, + }, + } + } + + /// The origin this front's replies name: `None` for this origin, see + /// [`Identity::origin`]. + pub(super) fn origin(&self) -> Option { + self.identity.origin() + } + + /// The origin whose copies may deliver content, once replies established one. + pub(super) fn admit(&self) -> Option { + self.named } /// The attached source, if any. @@ -289,6 +331,7 @@ impl Front { result, delivered, } => self.track_ended(track, source, closing, result, delivered, &mut actions), + Event::Origin { track, source, origin } => self.origin_named(track, source, origin, &mut actions), Event::Used { track } => self.used(track, &mut actions), Event::Unused { track, now } => self.unused(track, now, &mut actions), Event::Deadline { now } => self.deadline(now, &mut actions), @@ -378,6 +421,7 @@ impl Front { } } self.serving = Some((source, route)); + self.first.get_or_insert(source); self.serving_closing = false; if !self.resolved { self.resolved = true; @@ -404,7 +448,10 @@ impl Front { self.identity = match (candidate.local, candidate.first) { (true, _) => Identity::Local, (false, Some(hop)) if hop != Hop::UNKNOWN => Identity::Publisher(hop), - (false, _) => Identity::Anonymous { route: candidate.route }, + (false, first) => Identity::Anonymous { + route: candidate.route, + here: first.is_none(), + }, }; } @@ -429,7 +476,7 @@ impl Front { // A local publisher ending ends its broadcast; a newcomer at the // path gets a fresh one. An anonymous source can never be resumed. Identity::Local | Identity::Anonymous { .. } | Identity::Undetermined => self.end(Error::Dropped, actions), - Identity::Publisher(_) => actions.push(Action::Reselect), + Identity::Publisher(_) | Identity::Origin(_) => actions.push(Action::Reselect), } } @@ -441,10 +488,17 @@ impl Front { result: Result, actions: &mut Vec, ) { - let Some(track) = self.tracks.get_mut(&name) else { + if self.tracks.get(&name).map(|track| &track.state) != Some(&TrackState::Querying { source }) { return; - }; - if track.state != (TrackState::Querying { source }) { + } + // Once replies named the front's origin, a copy whose replies cannot is + // content nobody can vouch for: let it go, so the end cannot splice it + // either, and end the front; a re-request gets a fresh one. + if let Ok(info) = &result + && self.named.is_some() + && !info.names_origin + { + self.reject(source, actions); return; } let verdict = match result { @@ -468,6 +522,7 @@ impl Front { }; match verdict { Ok(()) => { + let track = self.tracks.get_mut(&name).expect("querying a known track"); track.state = TrackState::Spliced { source }; actions.push(Action::Splice { track: name.clone(), @@ -478,6 +533,71 @@ impl Front { } } + /// A spliced copy's reply named its origin: admit it, or end the front if it is + /// another origin's content. The copy held everything so far, so nothing of it + /// was delivered. + fn origin_named(&mut self, name: Arc, source: u64, origin: Hop, actions: &mut Vec) { + if self.serving.map(|(serving, _)| serving) != Some(source) || !self.tracks.contains_key(&name) { + return; + } + if !self.vouch(source, origin) { + self.reject(source, actions); + } + } + + /// End the front over a source carrying another origin's content. Its copies + /// delivered nothing, so a track spliced from one is aborted rather than left + /// to end with it, and the source is let go. + fn reject(&mut self, source: u64, actions: &mut Vec) { + let rejected: Vec> = self + .tracks + .iter() + .filter(|(_, track)| { + matches!(track.state, TrackState::Querying { source: s } | TrackState::Spliced { source: s } if s == source) + }) + .map(|(name, _)| name.clone()) + .collect(); + for name in rejected { + self.tracks.remove(&name); + actions.push(Action::Abort { + track: name, + err: Error::Dropped, + }); + } + actions.push(Action::Detach { source }); + self.end(Error::Dropped, actions); + } + + /// Whether `origin`, named by a reply from `source`, is the front's origin, + /// establishing it from the first source's first reply. + fn vouch(&mut self, source: u64, origin: Hop) -> bool { + if let Some(named) = self.named { + return named == origin; + } + match self.identity { + // The first source's reply names what its route could only label, since + // an advertiser serving a pool advertises one route for all of it. + Identity::Publisher(_) | Identity::Anonymous { .. } if self.first == Some(source) => { + self.identity = match origin { + Hop::UNKNOWN => Identity::Anonymous { + route: self + .serving + .map(|(_, route)| route) + .expect("a reply comes from a serving source"), + here: false, + }, + origin => Identity::Origin(origin), + }; + } + // Content already flowed under the route's label, from a wire that names no + // origin: only that origin matches. + Identity::Publisher(hop) if hop == origin => {} + _ => return false, + } + self.named = Some(origin); + true + } + fn track_ended( &mut self, name: Arc, @@ -652,6 +772,23 @@ mod tests { track::Info::default() } + /// Track metadata from a wire whose replies name the origin. + fn vouching() -> track::Info { + track::Info { + names_origin: true, + ..info() + } + } + + /// `source`'s copy of `video` named `origin`. + fn named(source: u64, origin: u64) -> Event { + Event::Origin { + track: name("video"), + source, + origin: Hop::from_wire(origin).unwrap(), + } + } + fn name(s: &str) -> Arc { Arc::from(s) } @@ -664,6 +801,11 @@ mod tests { /// A front serving `source` through `candidate`, with `video` spliced and read. fn serving(candidate: Candidate, source: u64) -> Front { + serving_with(candidate, source, None) + } + + /// [`serving`], with the copy of `video` replying that `origin` serves it. + fn serving_with(candidate: Candidate, source: u64, origin: Option) -> Front { let mut front = Front::new(LINGER); assert_actions( front.step(Event::Selected { @@ -693,13 +835,19 @@ mod tests { track: name("video"), source, closing: false, - result: Ok(info()), + result: Ok(match origin { + Some(_) => vouching(), + None => info(), + }), }), &[Action::Splice { track: name("video"), source, }], ); + if let Some(origin) = origin { + assert_actions(front.step(named(source, origin)), &[]); + } front } @@ -710,6 +858,177 @@ mod tests { assert_eq!(front.pin(), Pin::Publisher(hop(10))); } + /// An advertiser serving a pool advertises one route for all of it, so the + /// reply, not the route, names what the front serves. + #[test] + fn the_first_reply_names_the_origin() { + let front = serving_with(remote(1, 10), 100, Some(20)); + assert_eq!(front.identity, Identity::Origin(hop(20))); + assert_eq!(front.origin(), Some(hop(20))); + assert_eq!(front.admit(), Some(hop(20))); + // Any route may lead back to that origin, but the serving one stays. + assert_eq!(front.pin(), Pin::Stay(1)); + + // A reply naming nobody pins the front to its route. + let front = serving_with(remote(1, 10), 100, Some(0)); + assert_eq!(front.identity, Identity::Anonymous { route: 1, here: false }); + assert_eq!(front.origin(), Some(Hop::UNKNOWN)); + + // Nothing is admitted until a reply names something. + assert_eq!(serving(remote(1, 10), 100).admit(), None); + } + + /// A replacement source up to its TRACK_INFO: re-requested, resolved, and its + /// copy of `video` spliced (holding its content until its origin is admitted). + fn fail_over(front: &mut Front, candidate: Candidate, source: u64, info: track::Info) -> Vec { + front.step(Event::SourceClosed { + source: front.serving().unwrap(), + }); + front.step(Event::Selected { + best: Some(candidate), + serving_closing: false, + }); + front.step(Event::Resolved { + route: candidate.route, + result: Ok(source), + }); + front.step(Event::TrackInfo { + track: name("video"), + source, + closing: false, + result: Ok(info), + }) + } + + /// The serving source dies and the best remaining route has another first + /// hop, but its reply names the same origin: the front resumes there. + #[test] + fn a_replacement_naming_the_same_origin_resumes() { + let mut front = serving_with(remote(1, 10), 100, Some(20)); + assert_actions( + front.step(Event::SourceClosed { source: 100 }), + &[Action::Detach { source: 100 }, Action::Reselect], + ); + assert_eq!(front.pin(), Pin::Any); + assert_actions( + front.step(Event::Selected { + best: Some(remote(2, 11)), + serving_closing: false, + }), + &[Action::Request { route: 2 }], + ); + assert_actions( + front.step(Event::Resolved { + route: 2, + result: Ok(200), + }), + &[Action::Query { + track: name("video"), + source: 200, + }], + ); + assert_actions( + front.step(Event::TrackInfo { + track: name("video"), + source: 200, + closing: false, + result: Ok(vouching()), + }), + &[Action::Splice { + track: name("video"), + source: 200, + }], + ); + assert_actions(front.step(named(200, 20)), &[]); + assert_eq!(front.pin(), Pin::Stay(2)); + } + + /// A replacement whose reply names another origin is different content: its + /// copy, which held everything, is let go and the front ends. + #[test] + fn a_replacement_naming_another_origin_ends_the_front() { + for reply in [21, 0] { + let mut front = serving_with(remote(1, 10), 100, Some(20)); + fail_over(&mut front, remote(2, 10), 200, vouching()); + assert_actions( + front.step(named(200, reply)), + &[ + Action::Abort { + track: name("video"), + err: Error::Dropped, + }, + Action::Detach { source: 200 }, + Action::End { err: Error::Dropped }, + ], + ); + } + } + + /// A replacement whose wire names no origin cannot vouch for the front's: it + /// is let go before it is spliced. + #[test] + fn a_replacement_that_cannot_vouch_ends_the_front() { + let mut front = serving_with(remote(1, 10), 100, Some(20)); + assert_actions( + fail_over(&mut front, remote(2, 10), 200, info()), + &[ + Action::Abort { + track: name("video"), + err: Error::Dropped, + }, + Action::Detach { source: 200 }, + Action::End { err: Error::Dropped }, + ], + ); + } + + /// Once content flowed under a route's label from a wire naming no origin, a + /// later reply can only confirm the label, never replace it. + #[test] + fn a_reply_after_unnamed_content_must_match_the_route_label() { + let mut front = serving(remote(1, 10), 100); + fail_over(&mut front, remote(2, 10), 200, vouching()); + assert_actions( + front.step(named(200, 20)), + &[ + Action::Abort { + track: name("video"), + err: Error::Dropped, + }, + Action::Detach { source: 200 }, + Action::End { err: Error::Dropped }, + ], + ); + + let mut front = serving(remote(1, 10), 100); + fail_over(&mut front, remote(2, 10), 200, vouching()); + assert_actions(front.step(named(200, 10)), &[]); + assert_eq!(front.admit(), Some(hop(10))); + } + + /// A reply from a source the front already let go changes nothing. + #[test] + fn a_stale_reply_is_ignored() { + let mut front = serving_with(remote(1, 10), 100, Some(20)); + fail_over(&mut front, remote(2, 10), 200, vouching()); + assert_actions(front.step(named(100, 21)), &[]); + } + + /// A handler on this origin has an empty chain: its content is named by this + /// origin, though the front still never leaves its route. + #[test] + fn an_empty_chain_originates_here() { + let candidate = Candidate { + route: 1, + first: None, + local: false, + }; + let front = serving(candidate, 100); + assert_eq!(front.identity, Identity::Anonymous { route: 1, here: true }); + assert_eq!(front.origin(), None); + assert_eq!(front.pin(), Pin::Route(1)); + } + #[test] fn nothing_routable_ends_an_unresolved_front() { let mut front = Front::new(LINGER); @@ -1227,6 +1546,22 @@ mod tests { closing: false, result: Err(Error::NotFound), }, + Event::TrackInfo { + track: name("v"), + source: 200, + closing: false, + result: Ok(vouching()), + }, + Event::Origin { + track: name("v"), + source: 100, + origin: hop(20), + }, + Event::Origin { + track: name("v"), + source: 200, + origin: hop(21), + }, Event::TrackEnded { track: name("v"), source: 100, @@ -1275,6 +1610,16 @@ mod tests { .filter(|action| matches!(action, Action::Splice { .. })) .count(); assert!(splices <= 1); + // A copy's content is only admitted when the front serves the origin + // its reply named. + if let Event::Origin { source, origin, .. } = event + && !next.ended() + && front.serving() == Some(*source) + && front.tracks.contains_key("v") + { + assert_eq!(next.admit(), Some(*origin)); + assert_eq!(next.origin(), Some(*origin)); + } // A Detach names a source the front no longer serves from. for action in &actions { if let Action::Detach { source } = action { diff --git a/rs/moq-net/src/model/origin.rs b/rs/moq-net/src/model/origin.rs index 53b58a6d03..c5f06f7761 100644 --- a/rs/moq-net/src/model/origin.rs +++ b/rs/moq-net/src/model/origin.rs @@ -582,20 +582,25 @@ fn fnv_key(name: &str, origins: impl IntoIterator) -> u64 { hash } -/// Ordering key for a route entry covering one prefix. Lower wins: an identified +/// Ordering key for a route entry resolving `path`. Lower wins: an identified /// chain (no 0) outranks an anonymous one regardless of cost, then the cheapest /// cost, then a broadcast published on this origin (it serves what is here, not a /// claim that has to ask), then the shortest hop chain, then a deterministic hash -/// of the prefix and chain so every node converges on the same winner, and finally +/// of `path` and the chain so every node converges on the same winner, and finally /// the newest announcement, so a reconnect under an otherwise identical route wins /// the moment it lands instead of after the transport retires the old session. -fn route_order(prefix: &Path, entry: &RouteEntry) -> (bool, Cost, bool, usize, u64, Reverse) { +/// +/// `path` is what is being resolved: the requested path for a request, the prefix +/// itself for an advertisement. Keying the hash on the requested path is what +/// spreads equal-cost advertisers of one prefix: each path picks its own winner +/// from the pool, rather than every path under the prefix hashing alike. +fn route_order(path: &Path, entry: &RouteEntry) -> (bool, Cost, bool, usize, u64, Reverse) { ( entry.is_anonymous(), entry.cost, !entry.local, entry.hops.len(), - fnv_key(prefix.as_str(), entry.hops.iter().copied()), + fnv_key(path.as_str(), entry.hops.iter().copied()), Reverse(entry.id), ) } @@ -767,7 +772,7 @@ impl RouteEntry { /// Whether `pin` admits this entry for a front's selection. fn qualifies(&self, pin: Pin) -> bool { match pin { - Pin::Any => true, + Pin::Any | Pin::Stay(_) => true, Pin::Local => self.local, Pin::Publisher(first) => self.hops.iter().next() == Some(&first), Pin::Route(id) => self.id == id, @@ -1978,6 +1983,8 @@ struct FrontTask { request: kio::Producer, /// Published for requesters once the first source fixes it; see [`RemoteFront::pin`]. pin: kio::Lock, + /// This origin's hop, which names content originating here. + hop: Hop, timers: Clock, } @@ -2004,6 +2011,18 @@ struct TrackIo { head: Option, /// Whether the track had a reader as of the last demand edge. used: bool, + /// Whether the spliced copy's named origin was handed to the machine. + reported: bool, +} + +impl TrackIo { + /// Every copy the driver holds for the track, with its source. + fn copies(&self) -> impl Iterator { + let query = self.query.as_ref().map(|(source, copy, _)| (*source, copy)); + let staged = self.staged.as_ref().map(|(source, copy)| (*source, copy)); + let copy = self.copy.as_ref().map(|(source, copy)| (*source, copy)); + query.into_iter().chain(staged).chain(copy) + } } /// Drives one front: feeds the world's events to a [`Front`] and performs the @@ -2018,6 +2037,7 @@ async fn run_front(task: FrontTask) { watch, request, pin, + hop, timers, } = task; @@ -2027,6 +2047,7 @@ async fn run_front(task: FrontTask) { Resolved(u64, Result), SourceClosed(u64), Info(Arc, u64, Result), + Origin(Arc, u64, Hop), Ended(Arc, u64, Result<(), Error>), Demand(Arc), Deadline, @@ -2034,6 +2055,16 @@ async fn run_front(task: FrontTask) { } let mut front = Front::new(TRACK_IDLE_LINGER); + // The origin the front's replies name. Content nobody identifies (an anonymous + // source, or an origin without a hop) gets a random one for this front's life, + // rather than 0: a downstream relay can still resume within it, and nothing + // else ever names it. + let anonymous = Hop::random(); + let named = |front: &Front| match front.origin() { + None if hop != Hop::UNKNOWN => hop, + Some(origin) if origin != Hop::UNKNOWN => origin, + _ => anonymous, + }; let mut sources: HashMap = HashMap::new(); let mut next_source = 0u64; // The in-flight upstream request: the route and its pending channel. @@ -2071,7 +2102,21 @@ async fn run_front(task: FrontTask) { loop { while let Some(event) = events.pop_front() { - for action in front.step(event) { + let actions = front.step(event); + // Published before any copy is admitted below: once a copy's content + // flows, a reply for it names this origin. + *pin.lock() = front.pin(); + broadcast.set_origin(named(&front)); + if let Some(origin) = front.admit() { + for io in tracks.values() { + for (_, copy) in io.copies() { + if let Some(provenance) = copy.provenance() { + provenance.admit(origin); + } + } + } + } + for action in actions { match action { Action::Reselect => events.push_back(select(&mut front, &sources, &mut seen)), Action::Request { route } => { @@ -2106,6 +2151,7 @@ async fn run_front(task: FrontTask) { }; front.identify(candidate); *pin.lock() = front.pin(); + broadcast.set_origin(named(&front)); if let Some(source) = source { let id = next_source; next_source += 1; @@ -2175,8 +2221,14 @@ async fn run_front(task: FrontTask) { Action::Detach { source } => { sources.remove(&source); // Its copies go with it; the segments they delivered stay - // spliced until a replacement resumes past them. + // spliced until a replacement resumes past them. Whatever they + // still hold for admission never will be. for io in tracks.values_mut() { + for (_, copy) in io.copies().filter(|(s, _)| *s == source) { + if let Some(provenance) = copy.provenance() { + provenance.refuse(); + } + } if io.copy.as_ref().is_some_and(|(s, _)| *s == source) { io.copy = None; } @@ -2233,6 +2285,7 @@ async fn run_front(task: FrontTask) { // edge the copy is asked to advance. io.edge = io.resume.resume_position(); io.copy = Some((source, copy)); + io.reported = false; } Action::Park { track: name } => { let Some(io) = tracks.get_mut(&name) else { continue }; @@ -2271,6 +2324,12 @@ async fn run_front(task: FrontTask) { Action::Abort { track: name, err } => { if let Some(mut io) = tracks.remove(&name) { tracing::debug!(name = %name, %err, "aborting track"); + // Nothing a copy still holds for admission will be read. + for (_, copy) in io.copies() { + if let Some(provenance) = copy.provenance() { + provenance.refuse(); + } + } let _ = io.resume.abort(err); } } @@ -2293,6 +2352,16 @@ async fn run_front(task: FrontTask) { // ends as that copy does. let waiting = io.staged.take().map(|(_, copy)| copy); let waiting = waiting.or_else(|| io.query.take().map(|(_, copy, _)| copy)); + // Whatever a copy still holds is only delivered if it is the + // front's origin; with none established, nobody vouches for it. + for copy in waiting.iter().chain(io.copy.as_ref().map(|(_, copy)| copy)) { + if let Some(provenance) = copy.provenance() { + match front.admit() { + Some(origin) => provenance.admit(origin), + None => provenance.refuse(), + } + } + } if let Some(copy) = waiting && io.resume.is_used() { @@ -2348,6 +2417,12 @@ async fn run_front(task: FrontTask) { { return Poll::Ready(Step::Ended(name.clone(), *source, result)); } + if let Some((source, copy)) = &io.copy + && !io.reported && let Some(provenance) = copy.provenance() + && let Poll::Ready(origin) = provenance.poll_named(waiter) + { + return Poll::Ready(Step::Origin(name.clone(), *source, origin)); + } // Watch the demand edge in whichever direction is unmet. let edge = match io.used { true => io.resume.poll_unused(waiter), @@ -2377,6 +2452,7 @@ async fn run_front(task: FrontTask) { warm: None, head: None, used: false, + reported: false, }, ); Event::TrackAssigned { track: name } @@ -2404,6 +2480,15 @@ async fn run_front(task: FrontTask) { } } Step::SourceClosed(source) => Event::SourceClosed { source }, + Step::Origin(name, source, origin) => { + let Some(io) = tracks.get_mut(&name) else { continue }; + io.reported = true; + Event::Origin { + track: name, + source, + origin, + } + } Step::Info(name, source, result) => { let closing = sources.get(&source).is_some_and(|s| s.is_closing()); let Some(io) = tracks.get_mut(&name) else { continue }; @@ -2941,12 +3026,26 @@ impl OriginState { /// prefix, the cheapest served one is picked by [`route_order`]. /// /// Only announced routes are candidates: an unannounced broadcast serves - /// nobody, and does not shadow anything either. `pin` is the front's + /// nobody, and does not shadow anything either. The hash tie-break is keyed + /// on `path`, so equal-cost advertisers of one prefix share its paths. `pin` is the front's /// identity: only routes it admits are candidates, since a route from anyone /// else is different content rather than an alternate path (see [`Front`]). /// A broadcast published on this origin competes on cost like any other /// route and wins a tie. fn best_route(&self, path: &Path, horizon: Horizon, pin: Pin) -> Option<&RouteEntry> { + // A route a front stays on wins while it can still serve, whatever else + // appeared since: see [`Pin::Stay`]. + if let Pin::Stay(route) = pin + && let Some(entry) = self.routes.covering(path).find(|entry| { + entry.id == route + && entry.advertised + && entry.scope.matches(path.as_str()) + && horizon.admits(entry) + && entry.serves(path) + }) { + return Some(entry); + } + // Covering prefixes of one path form a chain, so the deepest node with a // candidate holds the unique longest prefix; walking down, the last such // node decides. @@ -2964,7 +3063,7 @@ impl OriginState { if candidates.peek().is_some() { best = candidates .filter(|entry| entry.serves(path)) - .min_by_key(|entry| route_order(&entry.prefix, entry)); + .min_by_key(|entry| route_order(path, entry)); } } best @@ -3708,6 +3807,7 @@ impl Consumer { watch, request, pin, + hop: self.hop, timers: self.timers.clone(), })); kio::Pending::new(Requesting::queued(consumer).with_path(requested).with_stats(scope)) @@ -4309,6 +4409,56 @@ mod tests { announced.assert_next_ended("room"); } + /// Equal-cost advertisers of one prefix share its paths: a set of requested + /// paths spreads across the pool, and one path always resolves to the same + /// advertiser, whatever order the routes arrived in. + /// Pinned so `spreadHash` in `js/net` picks the same pool member for a path. + #[test] + fn spread_hash_matches_js() { + assert_eq!(fnv_key("pool/job-0", [origin(10)]), 0xefb5e20a66101c32); + assert_eq!(fnv_key("pool/job-0", [origin(11)]), 0x0eb0a91370ff6653); + } + + #[tokio::test] + async fn equal_cost_pool_spreads_paths() { + const WORKERS: [u64; 4] = [10, 11, 12, 13]; + const PATHS: usize = 64; + + // The first hop of the route each path resolves to on a node whose pool + // arrived in `order`. + fn winners(order: impl Iterator) -> Vec { + let producer = origin(1).produce(); + let _pool: Vec = order + .map(|id| { + producer + .dynamic("pool", Route::default().with_hops(hops(&[id])).with_cost(3)) + .unwrap() + }) + .collect(); + let table = producer.shared.read(); + (0..PATHS) + .map(|i| { + let path = Path::new(&format!("pool/job-{i}")).to_owned(); + let entry = table + .best_route(&path.as_path(), Horizon::default(), Pin::Any) + .expect("the pool serves every path"); + entry.hops.iter().next().copied().unwrap() + }) + .collect() + } + + let forward = winners(WORKERS.into_iter()); + let reverse = winners(WORKERS.into_iter().rev()); + assert_eq!(forward, reverse, "a path must resolve the same way on every node"); + + // Not an assertion about any two paths, which a correct hash may put on + // one worker: only that the set does not pile onto a few. + for worker in WORKERS { + let share = forward.iter().filter(|hop| **hop == origin(worker)).count(); + assert!(share >= PATHS / 16, "worker {worker} took {share} of {PATHS} paths"); + } + } + #[tokio::test] async fn identical_reannounce_is_invisible() { let producer = origin(1).produce(); @@ -5992,6 +6142,176 @@ mod tests { pending.await.expect("re-request resolves through the rival"); } + /// A copy of "video" whose reply named `member`, as a Lite07 session's is. + fn vouching(source: &broadcast::Producer, member: u64) -> track::Producer { + let info = track::Info { + names_origin: true, + ..Default::default() + }; + let track = source.create_track("video", info).unwrap(); + track.provenance().name(origin(member)).unwrap(); + track + } + + /// Wait for the front to admit or refuse a vouching copy, which a Lite07 session + /// waits on before delivering anything. + async fn admission(track: &track::Producer) -> Result<(), Error> { + let provenance = track.provenance(); + let mut admitted = None; + settle(|| match provenance.poll_admitted(&kio::Waiter::noop()) { + Poll::Ready(result) => { + admitted = Some(result); + true + } + Poll::Pending => false, + }) + .await; + admitted.unwrap() + } + + /// A pool front: "room/alice" served through a route labelled `first` (at + /// cost 5) by a copy whose reply named `member`, one "before" group delivered. + async fn pool_rig(first: u64, member: u64) -> (ResumeRig, Dynamic, broadcast::Producer) { + let producer = origin(1).produce(); + let consumer = producer.consume(); + let server = producer + .dynamic("room", Route::default().with_hops(hops(&[first])).with_cost(5)) + .unwrap(); + + let pending = consumer.request_broadcast("room/alice"); + let request = queued(&server).await; + let source = broadcast::Info::new().produce(); + let track = vouching(&source, member); + request.accept(&source); + + let resolved = pending.await.expect("resolves"); + let mut subscription = resolved + .track("video") + .unwrap() + .subscribe(None) + .await + .expect("subscribe"); + admission(&track) + .await + .expect("the first reply names the front's origin"); + let mut group = track.append_group().unwrap(); + group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap(); + group.finish().unwrap(); + let mut group = subscription + .recv_group() + .await + .expect("recv group") + .expect("track ended early"); + let frame = group.read_frame().await.expect("read frame").expect("frame"); + assert_eq!(&frame.payload[..], b"before"); + + let rig = ResumeRig { + producer, + resolved, + subscription, + incumbent_track: track, + }; + (rig, server, source) + } + + /// A downstream relay reaches a pool through advertisements whose first hop + /// labels the pool, not the member serving the path. The reply names the + /// member, so a failover through a route with another label still resumes + /// when the replacement names the same member, and the relay names that + /// member in its own replies. + #[tokio::test] + async fn failover_resumes_on_the_origin_the_reply_names() { + let (mut rig, incumbent, source) = pool_rig(10, 20).await; + assert_eq!(rig.resolved.origin(), Some(origin(20))); + + let standby_server = rig.standby(&[11]); + drop(incumbent); + drop(source); + + let request = queued(&standby_server).await; + let replacement = broadcast::Info::new().produce(); + let track = vouching(&replacement, 20); + request.accept(&replacement); + admission(&track).await.expect("the same origin is admitted"); + let mut group = track.append_group().unwrap(); + group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap(); + group.finish().unwrap(); + let mut group = track.append_group().unwrap(); + group.write_frame(crate::Timestamp::ZERO, b"resumed".as_ref()).unwrap(); + group.finish().unwrap(); + + let mut group = rig + .subscription + .recv_group() + .await + .expect("subscription survives the failover") + .expect("track ended early"); + let frame = group.read_frame().await.expect("read frame").expect("frame"); + assert_eq!(&frame.payload[..], b"resumed"); + assert_eq!(rig.resolved.origin(), Some(origin(20))); + } + + /// Two pool members behind the same advertised label are different content: + /// a failover never splices one member's frames onto another's. + #[tokio::test] + async fn failover_never_splices_another_pool_member() { + let (mut rig, incumbent, source) = pool_rig(10, 20).await; + let standby_server = rig.standby(&[10]); + drop(incumbent); + drop(source); + rig.incumbent_track.abort(Error::Dropped).unwrap(); + + let request = queued(&standby_server).await; + let rival = broadcast::Info::new().produce(); + let track = vouching(&rival, 21); + request.accept(&rival); + assert!( + admission(&track).await.is_err(), + "another member's content is never admitted" + ); + + let err = rig.subscription.recv_group().await.err().expect("subscription ends"); + assert!(matches!(err, Error::Dropped), "unexpected end: {err}"); + } + + /// A front serving a named origin does not trade its live source for a + /// better route: which member the new route leads to is unknown until it + /// replies. Newcomers join it rather than starting a second copy. + #[tokio::test] + async fn a_named_origin_stays_on_its_live_route() { + let (rig, _incumbent, _source) = pool_rig(10, 20).await; + let cheaper = rig + .producer + .dynamic("room", Route::default().with_hops(hops(&[11])).with_cost(1)) + .unwrap(); + + let consumer = rig.producer.consume(); + let joined = consumer.request_broadcast("room/alice").await.expect("joins"); + assert!(joined.is_clone(&rig.resolved), "a newcomer joins the live front"); + for _ in 0..100 { + tokio::task::yield_now().await; + } + assert!( + cheaper.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(), + "the front must not re-request through the new route" + ); + } + + /// A front names the origin of what it serves in its replies: this origin's hop + /// for content originating here, and a random hop of its own for content nobody + /// identifies, never 0. + #[tokio::test] + async fn a_front_names_an_origin_for_its_content() { + let (rig, _server, _source) = ResumeRig::new(&[0]).await; + let anonymous = rig.resolved.origin().expect("a front names its origin"); + assert_ne!(anonymous, Hop::UNKNOWN); + assert_ne!(anonymous, origin(1)); + + // Content originating here is named by this origin. + let (rig, _server, _source) = ResumeRig::new(&[]).await; + assert_eq!(rig.resolved.origin(), Some(origin(1))); + } + #[tokio::test] async fn anonymous_routes_never_resume() { // An empty hop chain identifies nobody, so two of them must not pass for diff --git a/rs/moq-net/src/model/track.rs b/rs/moq-net/src/model/track.rs index 8ac10f45b9..c78d27b217 100644 --- a/rs/moq-net/src/model/track.rs +++ b/rs/moq-net/src/model/track.rs @@ -106,6 +106,10 @@ pub struct Info { /// The publisher's priority for this track, used only to break ties between /// subscriptions of equal subscriber priority. Reported in TRACK_INFO (Lite05+). pub priority: u8, + /// Whether a copy's replies name the origin serving it (SUBSCRIBE_OK and FETCH_OK + /// on Lite07+), recorded in its [`Provenance`]. The origin's failover only + /// splices a copy that can vouch for the front's origin this way. + pub(crate) names_origin: bool, } impl Default for Info { @@ -114,6 +118,7 @@ impl Default for Info { timescale: Timescale::default(), max_age: DEFAULT_MAX_AGE, priority: 0, + names_origin: false, } } } @@ -141,6 +146,89 @@ impl Info { } } +/// Which origin a copy's content comes from: the one its replies named, and the one +/// its consumer admits. +/// +/// A relay's front splices copies from several sources into one logical track, and +/// only copies of one origin may join. A session whose replies name the serving +/// origin records each reply here and holds the copy's content until the front +/// admits that origin, so a copy of another origin never delivers anything. +#[derive(Clone, Default)] +pub(crate) struct Provenance(kio::Shared); + +#[derive(Default)] +struct ProvenanceState { + named: Option, + admit: Option, + refused: bool, +} + +impl ProvenanceState { + fn admitted(&self) -> bool { + self.named.is_some() && self.named == self.admit + } +} + +impl Provenance { + /// Record the origin a reply named. A copy has one origin for its whole life, so + /// a reply naming another means the upstream is serving a different track. + pub(crate) fn name(&self, origin: crate::Hop) -> Result<()> { + let mut state = self.0.lock(); + match state.named { + Some(named) if named != origin => Err(Error::Dropped), + Some(_) => Ok(()), + None => { + state.named = Some(origin); + Ok(()) + } + } + } + + /// Wait for a reply to name the origin. + pub(crate) fn poll_named(&self, waiter: &kio::Waiter) -> Poll { + self.0 + .poll(waiter, |state| match state.named { + Some(_) => Poll::Ready(()), + None => Poll::Pending, + }) + .map(|state| state.named.expect("predicate guaranteed a name")) + } + + /// The origin whose content may be delivered. Replaces an earlier one: a front + /// learns its origin from the first reply, after it had only the route's label. + pub(crate) fn admit(&self, origin: crate::Hop) { + if self.0.read().admit != Some(origin) { + self.0.lock().admit = Some(origin); + } + } + + /// Nothing more of this copy will be admitted: its consumer let it go. + pub(crate) fn refuse(&self) { + if !self.0.read().refused { + self.0.lock().refused = true; + } + } + + /// Whether the named origin is admitted now, for content that cannot wait for it. + pub(crate) fn is_admitted(&self) -> bool { + self.0.read().admitted() + } + + /// Wait until the named origin is admitted, or fail once the copy is refused. + /// Admission wins: a refusal only stops what was still waiting. + pub(crate) fn poll_admitted(&self, waiter: &kio::Waiter) -> Poll> { + self.0 + .poll(waiter, |state| match state.admitted() || state.refused { + true => Poll::Ready(()), + false => Poll::Pending, + }) + .map(|state| match state.admitted() { + true => Ok(()), + false => Err(Error::Dropped), + }) + } +} + #[derive(Default)] pub(crate) struct TrackState { // The publisher's properties, once known; always Some for Subscriber/Producer. @@ -249,6 +337,10 @@ pub(crate) struct TrackState { // The reverse fetch queue (see [`FetchState`]), same reasoning: cache-miss // `fetch_group` calls enqueue here and a `Dynamic` drains. fetch: kio::Shared, + + // Which origin this copy's content comes from, same reasoning: the session + // records replies, the consumer admits. + provenance: Provenance, } /// A cached group plus its bookkeeping in the track's `lookup` map. @@ -1708,6 +1800,11 @@ impl Producer { }) } + /// Which origin this copy's content comes from; see [`Provenance`]. + pub(crate) fn provenance(&self) -> Provenance { + self.state.read().provenance.clone() + } + /// Create a [`Dynamic`] handle that serves on-demand fetches of uncached /// (old) groups. Most producers never need this; a relay creates one to fetch /// past groups from upstream. @@ -2422,6 +2519,15 @@ impl Consumer { } } + /// Which origin this copy's content comes from; see [`Provenance`]. `None` for a + /// spliced logical track, which is made of copies. + pub(crate) fn provenance(&self) -> Option { + match &self.inner { + ConsumerKind::Plain(state) => Some(state.read().provenance.clone()), + ConsumerKind::Spliced(_) => None, + } + } + /// The newest group, when it is already cached: resolved synchronously, without /// counting as a fetch or a delivery. The IETF publisher snapshots its frame count to /// resolve Largest Object; a group that is not immediately available reads as no edge. diff --git a/rs/moq-net/tests/datagram.rs b/rs/moq-net/tests/datagram.rs index 4cc78b1297..3cd596cf25 100644 --- a/rs/moq-net/tests/datagram.rs +++ b/rs/moq-net/tests/datagram.rs @@ -35,6 +35,10 @@ struct Fixture { /// A publisher and a subscriber joined over the mock transport, sharing one /// datagram-carrying track. async fn connect_datagram_track() -> Fixture { + connect_datagram_track_on("moq-lite-05").await +} + +async fn connect_datagram_track_on(version: &str) -> Fixture { let publisher = produce_origin(1); let consumer_origin = produce_origin(2); @@ -42,7 +46,7 @@ async fn connect_datagram_track() -> Fixture { let producer = broadcast.create_track("datagrams", None).unwrap(); broadcast.announce(Default::default()).unwrap(); - let mut options = MockConnectOptions::new("moq-lite-05".parse::().unwrap()); + let mut options = MockConnectOptions::new(version.parse::().unwrap()); options.server_publish = Some(publisher.consume()); options.client_subscribe = Some(consumer_origin.clone()); let pair = connect_mock(options).await; @@ -110,6 +114,34 @@ async fn datagrams_reach_the_subscriber_in_order() { .expect("timed out"); } +/// Lite-07 holds a subscription's content until SUBSCRIBE_OK names its origin, so a +/// track that only ever sends datagrams must still send SUBSCRIBE_OK for any to arrive. +/// Datagrams racing ahead of it are dropped, so this keeps sending until one lands. +#[tokio::test] +async fn datagram_only_track_names_its_origin_on_lite07() { + tokio::time::timeout(TEST_TIMEOUT, async { + let mut fixture = connect_datagram_track_on("moq-lite-07-wip").await; + + for sequence in 0.. { + fixture + .producer + .insert_datagram( + sequence, + Timestamp::from_millis(sequence).unwrap(), + bytes::Bytes::from_static(PAYLOAD), + ) + .unwrap(); + let arrived = tokio::time::timeout(Duration::from_millis(10), fixture.subscriber.recv_datagram()).await; + if let Ok(datagram) = arrived { + assert_eq!(&datagram.unwrap().unwrap().payload[..], PAYLOAD); + break; + } + } + }) + .await + .expect("no datagram arrived: SUBSCRIBE_OK never named the origin"); +} + /// MoQ Transport has no datagram mapping: groups still flow, inserted datagrams do not. #[tokio::test] async fn ietf_does_not_deliver_datagrams() { diff --git a/rs/moq-net/tests/pool_failover.rs b/rs/moq-net/tests/pool_failover.rs new file mode 100644 index 0000000000..d70083ba65 --- /dev/null +++ b/rs/moq-net/tests/pool_failover.rs @@ -0,0 +1,117 @@ +//! A relay spreading one prefix across a pool of workers, seen from downstream. +//! +//! Workers claim the same prefix, so the pool relay advertises one route for all +//! of them, labelled by whichever member ranks first for the prefix. Which member +//! serves a path is the relay's per-path choice, so the label cannot say what a +//! downstream relay is receiving; the SUBSCRIBE_OK and FETCH_OK replies do. When +//! the serving worker dies, the pool relay re-serves the path from another worker, +//! and the downstream relay must end its subscription rather than splice that +//! worker's frames onto the first one's. + +mod support; + +use std::time::Duration; + +use moq_net::{Hop, Timestamp, Version, origin}; +use support::harness::{MockConnectOptions, connect_mock}; + +const TIMEOUT: Duration = Duration::from_secs(10); + +fn produce_origin(hop: u64) -> origin::Producer { + let (producer, driver) = origin::Producer::new(origin::Config::new(Hop::new(hop).unwrap())); + tokio::spawn(support::harness::run(driver)); + producer +} + +/// A worker claiming the "pool" prefix: every request is answered with a fresh +/// broadcast whose "video" track keeps producing groups carrying the worker's name. +fn worker(hop: u64) -> origin::Producer { + let producer = produce_origin(hop); + let dynamic = producer.dynamic("pool", origin::Route::default()).unwrap(); + let name = format!("w{hop}").into_bytes(); + tokio::spawn(async move { + while let Ok(request) = dynamic.requested_broadcast().await { + let broadcast = moq_net::broadcast::Info::new().produce(); + let track = broadcast.create_track("video", None).unwrap(); + request.accept(&broadcast); + let name = name.clone(); + tokio::spawn(async move { + let _broadcast = broadcast; + loop { + let Ok(mut group) = track.append_group() else { return }; + group.write_frame(Timestamp::ZERO, name.clone()).unwrap(); + group.finish().unwrap(); + tokio::time::sleep(Duration::from_millis(5)).await; + } + }); + } + }); + producer +} + +#[tokio::test] +async fn downstream_failover_never_splices_another_pool_member() { + tokio::time::timeout(TIMEOUT, async { + let version: Version = "moq-lite-07-wip".parse().unwrap(); + let pool = produce_origin(10); + let downstream = produce_origin(30); + + // Each worker publishes its claim to the pool relay. + let mut members = Vec::new(); + for hop in [20, 21] { + let producer = worker(hop); + let mut options = MockConnectOptions::new(version); + options.client_publish = Some(producer.consume()); + options.server_subscribe = Some(pool.clone()); + members.push((hop, producer, connect_mock(options).await)); + } + + // The downstream relay reaches the pool through one session. + let mut options = MockConnectOptions::new(version); + options.server_publish = Some(pool.consume()); + options.client_subscribe = Some(downstream.clone()); + let _link = connect_mock(options).await; + + let consumer = downstream.consume(); + consumer.routed("pool/job").await.unwrap(); + let remote = consumer.request_broadcast("pool/job").await.unwrap(); + let mut subscription = remote.track("video").unwrap().subscribe(None).await.unwrap(); + + let mut group = subscription.recv_group().await.unwrap().expect("a first group"); + let first = group.read_frame().await.unwrap().expect("a frame").payload.to_vec(); + + // Kill whichever worker the pool relay picked for this path. + let serving = members + .iter() + .position(|(hop, ..)| first == format!("w{hop}").into_bytes()) + .expect("a pool member served the path"); + let (_, _, pair) = members.remove(serving); + pair.client.abort(moq_net::Error::Cancel); + + // Every group this subscription delivers is the first worker's; it ends + // instead of carrying on with the survivor's. + while let Ok(Some(mut group)) = subscription.recv_group().await { + while let Some(frame) = group.read_frame().await.unwrap_or(None) { + assert_eq!(frame.payload.to_vec(), first, "spliced another pool member's frames"); + } + } + + // A fresh request is served by the survivor. + let (hop, ..) = &members[0]; + let survivor = format!("w{hop}").into_bytes(); + let remote = consumer.request_broadcast("pool/job").await.unwrap(); + let track = remote.track("video").unwrap(); + let mut subscription = track.subscribe(None).await.unwrap(); + let mut group = subscription.recv_group().await.unwrap().expect("a group"); + let frame = group.read_frame().await.unwrap().expect("a frame"); + assert_eq!(frame.payload.to_vec(), survivor); + + // A group from before the subscription is fetched through both relays, each + // FETCH_OK naming the survivor. + let mut fetched = track.fetch_group(0, None).await.unwrap(); + let frame = fetched.read_frame().await.unwrap().expect("a frame"); + assert_eq!(frame.payload.to_vec(), survivor); + }) + .await + .expect("timed out"); +}