From 751c07f21465ec05d0f2ca11cc2cde528daab19a Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Wed, 30 Sep 2026 09:33:29 -0700 Subject: [PATCH 1/4] quest(m0): claim ietf-legal-input Co-Authored-By: Claude Opus 5.5 From 49f19ad1b409a79225aba1b1068aa194a78589b0 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Wed, 30 Sep 2026 10:22:43 -0700 Subject: [PATCH 2/4] fix(ietf): decode every legal request and refuse per request Draft-20+ FETCH, legal request parameters (AUTHORIZATION TOKEN, Range Filters, NEW_GROUP_REQUEST, FILL_TIMEOUT), INCLUDE_PROPERTIES as a uint8, FORWARD=0, and TRACK_STATUS parameters now decode; what we don't serve is refused NOT_SUPPORTED instead of closing the session. Mirrored in js/net. Co-Authored-By: Claude Opus 5.5 --- doc/concept/standard.md | 14 +- js/net/src/ietf/connection.ts | 6 + js/net/src/ietf/fetch.ts | 97 ++++--- js/net/src/ietf/filter.ts | 5 + js/net/src/ietf/ietf.test.ts | 103 +++++++ js/net/src/ietf/parameters.ts | 63 +++- js/net/src/ietf/publisher.test.ts | 47 ++- js/net/src/ietf/publisher.ts | 83 ++++-- js/net/src/ietf/subscribe.ts | 39 ++- js/net/src/ietf/track.ts | 28 +- quest/m0/README.md | 1 - quest/m0/ietf-fin-not-cancel.md | 4 - quest/m0/ietf-legal-input.md | 49 ---- quest/m1/auth/request-token.md | 15 +- quest/m1/moxygen/fetch.md | 7 +- rs/moq-net/src/fuzz.rs | 27 +- rs/moq-net/src/ietf/fetch.rs | 240 ++++++++++++++-- rs/moq-net/src/ietf/parameters.rs | 141 ++++++++- rs/moq-net/src/ietf/publish.rs | 5 +- rs/moq-net/src/ietf/publish_namespace.rs | 38 ++- rs/moq-net/src/ietf/publisher.rs | 129 ++++++++- rs/moq-net/src/ietf/subscribe.rs | 316 +++++++++++++++++---- rs/moq-net/src/ietf/subscribe_namespace.rs | 14 +- rs/moq-net/src/ietf/subscriber.rs | 3 + rs/moq-net/src/ietf/track.rs | 73 +++-- rs/moq-net/src/ietf/version.rs | 18 +- 26 files changed, 1238 insertions(+), 327 deletions(-) delete mode 100644 quest/m0/ietf-legal-input.md diff --git a/doc/concept/standard.md b/doc/concept/standard.md index 54801adec8..29ac388222 100644 --- a/doc/concept/standard.md +++ b/doc/concept/standard.md @@ -40,7 +40,8 @@ On drafts 14–19, the Rust publisher serves relative joining `FETCH` requests with offset zero for `NextObject` subscriptions. The fetch delivers the saved current-group prefix, and the subscription delivers later objects. Standalone, absolute joining, and nonzero-offset fetches are refused. Draft-20 uses -subscription fills instead. JavaScript publishing does not yet serve `FETCH`; +subscription fills instead, so its `FETCH` is refused. JavaScript publishing +refuses every `FETCH`; Rust and JavaScript subscribers request unfiltered delivery on older drafts because they do not issue joining fetches. Other publishers may replay a cached backlog for that filter; selecting the next group instead would leave static @@ -53,7 +54,16 @@ and bytes to the application unverified; a relay forwards them to its [auth server](/bin/relay/auth#the-contract). An alias reference (`DELETE`, `USE_ALIAS`) closes the session with `PROTOCOL_VIOLATION`, a structure that does not decode with `KEY_VALUE_FORMATTING_ERROR`, and a second token is -refused. +refused. An `AUTHORIZATION TOKEN` parameter on a request is read and ignored: +the session's credential is what authorizes it. + +A legal request that is not served is refused on its own with `NOT_SUPPORTED`, +leaving the session open: a `SUBSCRIBE` with `FORWARD=0`, a `SUBSCRIBE` or +`FETCH` carrying Range Filters (no `MAX_FILTER_RANGES` is advertised), +`TRACK_STATUS`, and the `FETCH` forms above. `NEW_GROUP_REQUEST` is ignored, as +the draft allows a publisher without dynamic groups to do. A parameter the +negotiated draft does not define still closes the session with +`PROTOCOL_VIOLATION`, as the draft requires. Several project drafts extend the IETF wire without breaking it, since `SETUP` ignores unknown parameters: [cluster](/draft/moq-cluster) routing hop lists, diff --git a/js/net/src/ietf/connection.ts b/js/net/src/ietf/connection.ts index c2a11987c7..d4eb993892 100644 --- a/js/net/src/ietf/connection.ts +++ b/js/net/src/ietf/connection.ts @@ -10,6 +10,7 @@ import { type Reader, Readers, type Stream } from "../stream.ts"; import { registerWire } from "../wire.ts"; import { ControlStreamAdapter, NativeSession, type Session } from "./adapter.ts"; import * as Cluster from "./cluster.ts"; +import { Fetch } from "./fetch.ts"; import { GoAway } from "./goaway.ts"; import { Group } from "./object.ts"; import { Publish } from "./publish.ts"; @@ -247,6 +248,11 @@ export class Connection implements Established { await this.#publisher.runTrackStatusRequest(msg, stream); break; } + case Fetch.id: { + const msg = await Fetch.decode(stream.reader, this.#session.version); + await this.#publisher.runFetch(msg, stream); + break; + } // Subscriber handles incoming notifications case PublishNamespace.id: { diff --git a/js/net/src/ietf/fetch.ts b/js/net/src/ietf/fetch.ts index d4a9f8c92a..1263555689 100644 --- a/js/net/src/ietf/fetch.ts +++ b/js/net/src/ietf/fetch.ts @@ -1,7 +1,9 @@ -import type * as Path from "../path.ts"; import type { Reader, Writer } from "../stream.ts"; +import { hasRangeFilters, isDraft20 } from "./filter.ts"; import * as Message from "./message.ts"; -import type { IetfVersion } from "./version.ts"; +import * as Namespace from "./namespace.ts"; +import { Parameters } from "./parameters.ts"; +import { type IetfVersion, Version } from "./version.ts"; /** * The header that begins a fetch stream (draft-20 section 11.4.4), naming the request it @@ -29,49 +31,17 @@ export class FetchHeader { } } +/** + * A FETCH request. We serve none, so only what a refusal needs is kept; the rest is + * decoded so a legal request is refused rather than breaking the stream. + */ export class Fetch { static id = 0x16; requestId: bigint; - trackNamespace: Path.Valid; - trackName: string; - subscriberPriority: number; - groupOrder: number; - startGroup: bigint; - startObject: bigint; - endGroup: bigint; - endObject: bigint; - constructor({ - requestId, - trackNamespace, - trackName, - subscriberPriority, - groupOrder, - startGroup, - startObject, - endGroup, - endObject, - }: { - requestId: bigint; - trackNamespace: Path.Valid; - trackName: string; - subscriberPriority: number; - groupOrder: number; - startGroup: bigint; - startObject: bigint; - endGroup: bigint; - endObject: bigint; - }) { + constructor({ requestId }: { requestId: bigint }) { this.requestId = requestId; - this.trackNamespace = trackNamespace; - this.trackName = trackName; - this.subscriberPriority = subscriberPriority; - this.groupOrder = groupOrder; - this.startGroup = startGroup; - this.startObject = startObject; - this.endGroup = endGroup; - this.endObject = endObject; } async #encode(_w: Writer): Promise { @@ -82,12 +52,49 @@ export class Fetch { return Message.encode(w, this.#encode.bind(this)); } - static async decode(r: Reader, _version: IetfVersion): Promise { - return Message.decode(r, Fetch.#decode); - } - - static async #decode(_r: Reader): Promise { - throw new Error("FETCH messages are not supported"); + static async decode(r: Reader, version: IetfVersion): Promise { + return Message.decode(r, (mr) => Fetch.#decode(mr, version)); + } + + static async #decode(r: Reader, version: IetfVersion): Promise { + const requestId = await r.u62(); + if (version === Version.DRAFT_17) { + await r.u62(); // required_request_id_delta + } + + if (version === Version.DRAFT_14) { + await r.u8(); // subscriber_priority + await r.u8(); // group_order + } + + if (isDraft20(version)) { + // Draft-20 names the track up front and moves the range into LOCATION_FILTER. + await Namespace.decode(r); + await r.string(); + } else { + const fetchType = await r.u53(); + switch (fetchType) { + case 0x1: // Standalone: namespace, name, start and end Locations + await Namespace.decode(r); + await r.string(); + for (let i = 0; i < 4; i++) await r.u62(); + break; + case 0x2: // Relative Joining: subscription, group offset + case 0x3: // Absolute Joining: subscription, group + await r.u62(); + await r.u62(); + break; + default: + throw new Error(`unknown fetch type: ${fetchType}`); + } + } + + const params = await Parameters.decode(r, version); + if (params.rangeFilters && !hasRangeFilters(version)) { + throw new Error("Range Filters need draft-19"); + } + + return new Fetch({ requestId }); } } diff --git a/js/net/src/ietf/filter.ts b/js/net/src/ietf/filter.ts index 9112990372..d3157b1a7f 100644 --- a/js/net/src/ietf/filter.ts +++ b/js/net/src/ietf/filter.ts @@ -48,6 +48,11 @@ export function isDraft20(version: IetfVersion): boolean { ); } +/** Whether the Range Filters (draft-19) exist on this draft. */ +export function hasRangeFilters(version: IetfVersion): boolean { + return isDraft20(version) || version === Version.DRAFT_19; +} + /** * Draft-17 replaced QUIC's two-bit-length varint with a leading-1-bits one. The two agree * below 64 and diverge above it, so getting this wrong is invisible until group or object diff --git a/js/net/src/ietf/ietf.test.ts b/js/net/src/ietf/ietf.test.ts index cfd7d29525..12f74c5220 100644 --- a/js/net/src/ietf/ietf.test.ts +++ b/js/net/src/ietf/ietf.test.ts @@ -3,6 +3,7 @@ import * as Path from "../path.ts"; import { type Cursor, Reader, Writer } from "../stream.ts"; import { Timescale, Timestamp } from "../time.ts"; import * as Varint from "../varint.ts"; +import { Fetch } from "./fetch.ts"; import * as GoAway from "./goaway.ts"; import * as Namespace from "./namespace.ts"; import { FetchFrame, type FetchPosition, Frame, Group, type GroupFlags } from "./object.ts"; @@ -1797,3 +1798,105 @@ test("FetchFrame: an unstamped first object keeps its ids", async () => { 0x01, ]); }); + +/** Frame a control message body with its 16-bit Length, as every draft does. */ +function framed(body: number[]): Uint8Array { + return new Uint8Array([body.length >> 8, body.length & 0xff, ...body]); +} + +/** Request ID 1, Track Namespace ("live"), Track Name ("video"): the head of a SUBSCRIBE or draft-20 FETCH. */ +const TRACK_HEAD = [0x01, 0x01, 0x04, ...new TextEncoder().encode("live"), 0x05, ...new TextEncoder().encode("video")]; + +// Every parameter draft-20 lets a SUBSCRIBE carry that we don't act on still decodes, with +// what we have to refuse recorded rather than failing the stream. +test("Subscribe v20: decodes every legal parameter", async () => { + const body = framed([ + ...TRACK_HEAD, + 0x06, // Number of Parameters + ...[0x03, 0x03, 0x03, 0x00, 0xaa], // AUTHORIZATION TOKEN + ...[0x00, 0x03, 0x03, 0x00, 0xbb], // a second one, which the draft allows + ...[0x0d, 0x00], // FORWARD (0x10) = 0 + ...[0x15, 0x03, 0x00, 0x00, 0x01], // SUBGROUP_FILTER (0x25) + ...[0x0d, 0x07], // NEW_GROUP_REQUEST (0x32) + ...[0x03, 0x00], // INCLUDE_PROPERTIES (0x35) = 0, a uint8 + ]); + + const msg = await decodeVersioned(body, Subscribe.Subscribe.decode, Version.DRAFT_20); + expect(msg.forward).toBe(false); + expect(msg.rangeFilters).toBe(true); + expect(msg.propertiesWanted).toBe(false); +}); + +// INCLUDE_PROPERTIES is a uint8: one raw byte with no Length. +test("Subscribe v20: INCLUDE_PROPERTIES is a uint8", async () => { + const msg = new Subscribe.Subscribe({ + requestId: 1n, + trackNamespace: Path.from("live"), + trackName: "video", + subscriberPriority: 128, + propertiesWanted: false, + }); + const encoded = await encodeVersioned(msg, Version.DRAFT_20); + expect(Array.from(encoded.slice(-2))).toEqual([0x13, 0x00]); + expect(Array.from(encoded.slice(-3, -2))).not.toEqual([0x01]); +}); + +// FORWARD=0 decodes on every draft, so the request can be refused rather than failing. +test("Subscribe: FORWARD=0 decodes", async () => { + for (const version of [Version.DRAFT_14, Version.DRAFT_15, Version.DRAFT_17, Version.DRAFT_20] as const) { + const msg = new Subscribe.Subscribe({ + requestId: 1n, + trackNamespace: Path.from("live"), + trackName: "video", + subscriberPriority: 128, + forward: false, + }); + const decoded = await decodeVersioned(await encodeVersioned(msg, version), Subscribe.Subscribe.decode, version); + expect(decoded.forward).toBe(false); + } +}); + +// A draft-20 FETCH names the track up front and carries its range in LOCATION_FILTER. +test("Fetch v20: decodes every legal parameter", async () => { + const body = framed([ + ...TRACK_HEAD, + 0x07, // Number of Parameters + ...[0x03, 0x03, 0x03, 0x00, 0xaa], // AUTHORIZATION TOKEN + ...[0x07, 0x64], // FILL_TIMEOUT (0x0A) + ...[0x16, 0x40], // SUBSCRIBER_PRIORITY (0x20) + ...[0x01, 0x03, 0x04, 0x00, 0x02], // LOCATION_FILTER (0x21) + ...[0x01, 0x01], // GROUP_ORDER (0x22) + ...[0x04, 0x02, 0x00, 0x05], // OBJECTID_FILTER (0x26) + ...[0x0f, 0x00], // INCLUDE_PROPERTIES (0x35) + ]); + + const msg = await decodeVersioned(body, Fetch.decode, Version.DRAFT_20); + expect(msg.requestId).toBe(1n); +}); + +// The tagged forms through draft-19, with a token. +test("Fetch v15: decodes a joining FETCH", async () => { + const body = framed([ + 0x01, // Request ID + ...[0x02, 0x03, 0x00], // Relative Joining: subscription 3, offset 0 + 0x01, // Number of Parameters + ...[0x03, 0x03, 0x03, 0x00, 0xaa], // AUTHORIZATION TOKEN + ]); + + const msg = await decodeVersioned(body, Fetch.decode, Version.DRAFT_15); + expect(msg.requestId).toBe(1n); +}); + +// TRACK_STATUS is identical to SUBSCRIBE on every draft, fields and parameters included. +test("TrackStatusRequest: carries SUBSCRIBE fields", async () => { + const cases: [IetfVersion, number[]][] = [ + // Priority, group order, forward, AbsoluteStart {5, 1}, no parameters. + [Version.DRAFT_14, [0x80, 0x02, 0x00, 0x03, 0x05, 0x01, 0x00]], + // AUTHORIZATION TOKEN and INCLUDE_PROPERTIES. + [Version.DRAFT_20, [0x02, 0x03, 0x03, 0x03, 0x00, 0xaa, 0x32, 0x00]], + ]; + for (const [version, rest] of cases) { + const msg = await decodeVersioned(framed([...TRACK_HEAD, ...rest]), Track.TrackStatusRequest.decode, version); + expect(msg.trackName).toBe("video"); + } +}); diff --git a/js/net/src/ietf/parameters.ts b/js/net/src/ietf/parameters.ts index f2a6a3120b..6b7ccf841b 100644 --- a/js/net/src/ietf/parameters.ts +++ b/js/net/src/ietf/parameters.ts @@ -194,6 +194,10 @@ export class SetupOptions { // Varint parameter IDs (even) const MSG_PARAM_DELIVERY_TIMEOUT = 0x02n; +/// FILL_TIMEOUT, on a FETCH. Ignored: a refusal doesn't wait. +const MSG_PARAM_FILL_TIMEOUT = 0x0an; +/// NEW_GROUP_REQUEST. Ignored, as the draft lets a publisher without dynamic groups do. +const MSG_PARAM_NEW_GROUP_REQUEST = 0x32n; /// SUBGROUP_DELIVERY_TIMEOUT, alongside the per-object one above. const MSG_PARAM_SUBGROUP_DELIVERY_TIMEOUT = 0x06n; const MSG_PARAM_MAX_CACHE_DURATION = 0x04n; @@ -202,22 +206,30 @@ const MSG_PARAM_PUBLISHER_PRIORITY = 0x0en; const MSG_PARAM_FORWARD = 0x10n; const MSG_PARAM_SUBSCRIBER_PRIORITY = 0x20n; const MSG_PARAM_GROUP_ORDER = 0x22n; +/// INCLUDE_PROPERTIES, draft-20's opt-out from Track Properties. +const MSG_PARAM_INCLUDE_PROPERTIES = 0x35n; /// ROUTE_COST, from the MoQ Cluster extension. See `cluster.ts`. const MSG_PARAM_ROUTE_COST = 0x40b58n; /// HIDDEN, from the MoQ Hidden extension. See `hidden.ts`. const MSG_PARAM_HIDDEN = 0x40b5en; // Bytes parameter IDs (odd) +/// AUTHORIZATION TOKEN. Ignored: the session's grant is what authorizes a request. +const MSG_PARAM_AUTHORIZATION_TOKEN = 0x03n; const MSG_PARAM_LARGEST_OBJECT = 0x09n; const MSG_PARAM_SUBSCRIPTION_FILTER = 0x21n; /// FILL_PARAMETERS, draft-20's request for a backfill. const MSG_PARAM_FILL_PARAMETERS = 0x23n; -/// INCLUDE_PROPERTIES, draft-20's opt-out from Track Properties. Its own section calls the -/// value a uint8, but 0x35 is odd, so the Key-Value-Pair rule length prefixes it. -const MSG_PARAM_INCLUDE_PROPERTIES = 0x35n; /// HOP_PATH, from the MoQ Cluster extension. See `cluster.ts`. const MSG_PARAM_HOP_PATH = 0x40b57n; +/// The Range Filters (draft-19): SUBGROUP, OBJECTID, PRIORITY, OBJECT_PROPERTY and +/// TRACK_PROPERTY. Each is length prefixed whatever the parity of its id. +const MSG_PARAM_RANGE_FILTERS: readonly bigint[] = [0x25n, 0x26n, 0x27n, 0x28n, 0x29n]; + +/// The parameters whose definitions let them repeat within one message. +const MSG_PARAM_REPEATABLE: readonly bigint[] = [MSG_PARAM_AUTHORIZATION_TOKEN, ...MSG_PARAM_RANGE_FILTERS]; + type MessageParamKind = "varint" | "uint8" | "bool" | "location" | "bytes"; /** A `{Group, Object}` pair carried by a message parameter, such as LARGEST_OBJECT. */ export type MessageLocation = { groupId: bigint; objectId: bigint }; @@ -225,6 +237,8 @@ export type MessageLocation = { groupId: bigint; objectId: bigint }; function getMessageParamKind(id: bigint): MessageParamKind { switch (id) { case MSG_PARAM_DELIVERY_TIMEOUT: + case MSG_PARAM_FILL_TIMEOUT: + case MSG_PARAM_NEW_GROUP_REQUEST: case MSG_PARAM_SUBGROUP_DELIVERY_TIMEOUT: case MSG_PARAM_MAX_CACHE_DURATION: case MSG_PARAM_EXPIRES: @@ -234,17 +248,19 @@ function getMessageParamKind(id: bigint): MessageParamKind { case MSG_PARAM_PUBLISHER_PRIORITY: case MSG_PARAM_SUBSCRIBER_PRIORITY: case MSG_PARAM_GROUP_ORDER: + case MSG_PARAM_INCLUDE_PROPERTIES: return "uint8"; case MSG_PARAM_FORWARD: return "bool"; case MSG_PARAM_LARGEST_OBJECT: return "location"; + case MSG_PARAM_AUTHORIZATION_TOKEN: case MSG_PARAM_SUBSCRIPTION_FILTER: case MSG_PARAM_FILL_PARAMETERS: - case MSG_PARAM_INCLUDE_PROPERTIES: case MSG_PARAM_HOP_PATH: return "bytes"; default: + if (MSG_PARAM_RANGE_FILTERS.includes(id)) return "bytes"; throw new Error(`unknown message parameter id: ${id.toString()}`); } } @@ -277,11 +293,28 @@ export class Parameters { vars: Map; bytes: Map; #locations: Map; + /** Every instance of a parameter that may repeat, decoded only; we never send one. */ + #repeated: Map; constructor() { this.vars = new Map(); this.bytes = new Map(); this.#locations = new Map(); + this.#repeated = new Map(); + } + + #repeat(id: bigint, value: Uint8Array) { + const values = this.#repeated.get(id) ?? []; + values.push(value); + this.#repeated.set(id, values); + } + + /** + * Whether the message carried a Range Filter. We advertise no MAX_FILTER_RANGES, so a + * request with one is refused rather than served unfiltered. + */ + get rangeFilters(): boolean { + return MSG_PARAM_RANGE_FILTERS.some((id) => this.#repeated.has(id)); } // --- Numeric accessors --- @@ -395,17 +428,15 @@ export class Parameters { /** INCLUDE_PROPERTIES: whether the peer wants Track Properties on the response. */ get includeProperties(): boolean | undefined { - const data = this.bytes.get(MSG_PARAM_INCLUDE_PROPERTIES); - if (!data) return undefined; + const v = this.vars.get(MSG_PARAM_INCLUDE_PROPERTIES); + if (v === undefined) return undefined; // The draft allows exactly 0 or 1; anything else is a protocol violation. - if (data.length !== 1 || data[0] > 1) { - throw new Error(`invalid INCLUDE_PROPERTIES value: ${data.join(",")}`); - } - return data[0] === 1; + if (v > 1n) throw new Error(`invalid INCLUDE_PROPERTIES value: ${v}`); + return v === 1n; } set includeProperties(v: boolean) { - this.bytes.set(MSG_PARAM_INCLUDE_PROPERTIES, new Uint8Array([v ? 1 : 0])); + this.vars.set(MSG_PARAM_INCLUDE_PROPERTIES, v ? 1n : 0n); } /** HOP_PATH: the hop chain an advertisement traversed, as its raw parameter value. */ @@ -552,7 +583,9 @@ export class Parameters { } else { const size = await r.u53(); const bytes = await r.read(size); - if (id === MSG_PARAM_LARGEST_OBJECT) { + if (MSG_PARAM_REPEATABLE.includes(id)) { + params.#repeat(id, bytes); + } else if (id === MSG_PARAM_LARGEST_OBJECT) { if (params.#locations.has(id)) { throw new Error(`duplicate message parameter id: ${id.toString()}`); } @@ -567,6 +600,12 @@ export class Parameters { continue; } + if (MSG_PARAM_REPEATABLE.includes(id)) { + const size = await r.u53(); + params.#repeat(id, await r.read(size)); + continue; + } + if (params.vars.has(id) || params.bytes.has(id) || params.#locations.has(id)) { throw new Error(`duplicate message parameter id: ${id.toString()}`); } diff --git a/js/net/src/ietf/publisher.test.ts b/js/net/src/ietf/publisher.test.ts index e7aa4b2cd6..b7a4a0889e 100644 --- a/js/net/src/ietf/publisher.test.ts +++ b/js/net/src/ietf/publisher.test.ts @@ -13,7 +13,7 @@ import type { Producer as TrackProducer } from "../track.ts"; import { wireOf } from "../wire.ts"; import { NativeSession, type Session } from "./adapter.ts"; import type * as Cluster from "./cluster.ts"; -import { FetchHeader } from "./fetch.ts"; +import { Fetch, FetchHeader } from "./fetch.ts"; import { Frame, Group as GroupMessage } from "./object.ts"; import { PublishDone } from "./publish.ts"; import { PublishNamespace, PublishNamespaceUpdate } from "./publish_namespace.ts"; @@ -174,6 +174,51 @@ test("TRACK_STATUS gets exact NOT_SUPPORTED refusal bytes on every draft", async } }); +// Legal requests we don't serve are refused NOT_SUPPORTED one at a time. +test("FETCH and a non-forwarding SUBSCRIBE get NOT_SUPPORTED", async () => { + const refusal = async (version: IetfVersion, run: (pub: Publisher, stream: Stream) => Promise) => { + const pair = createMockTransportPair(ALPN.DRAFT_19); + const session = new NativeSession(pair.server, version, true); + const { pub, origin } = publisher(pair.server, { session }); + const written: Uint8Array[] = []; + const stream = new Stream({ + readable: new ReadableStream(), + writable: new WritableStream({ + write: (chunk) => { + written.push(new Uint8Array(chunk)); + }, + }), + version, + }); + await run(pub, stream); + await stream.writer.closed; + origin.close(); + return written.flatMap((chunk) => Array.from(chunk)); + }; + + for (const version of [Version.DRAFT_14, Version.DRAFT_16, Version.DRAFT_20] as const) { + const fetch = await refusal(version, (pub, stream) => pub.runFetch(new Fetch({ requestId: 7n }), stream)); + // FETCH_ERROR on draft-14, REQUEST_ERROR after; the code follows the Length and any Request ID. + expect(fetch[0]).toBe(version === Version.DRAFT_14 ? 0x19 : 0x05); + expect(fetch[version <= Version.DRAFT_16 ? 4 : 3]).toBe(0x3); + + const paused = await refusal(version, (pub, stream) => + pub.runSubscribe( + new Subscribe({ + requestId: 7n, + trackNamespace: Path.from("test"), + trackName: "video", + subscriberPriority: 0, + forward: false, + }), + stream, + ), + ); + expect(paused[0]).toBe(0x05); + expect(paused[version <= Version.DRAFT_16 ? 4 : 3]).toBe(0x3); + } +}); + // The header is part of the group's lifetime too. If it blocks on flow control, advancing // the live edge must reset the stream without waiting for that write to finish. test("a blocked group header is reset when the group expires", async () => { diff --git a/js/net/src/ietf/publisher.ts b/js/net/src/ietf/publisher.ts index a251b0878b..477d4faba2 100644 --- a/js/net/src/ietf/publisher.ts +++ b/js/net/src/ietf/publisher.ts @@ -14,8 +14,8 @@ import * as Varint from "../varint.ts"; import { type Advertised, type Advertisements, wireOf } from "../wire.ts"; import type { Session } from "./adapter.ts"; import * as Cluster from "./cluster.ts"; -import { requestReason, toRequestCode } from "./error.ts"; -import { FetchHeader } from "./fetch.ts"; +import { type RequestKind, requestReason, toRequestCode } from "./error.ts"; +import { type Fetch, FetchError, FetchHeader } from "./fetch.ts"; import * as Filter from "./filter.ts"; import { FetchFrame, Frame, Group as GroupMessage } from "./object.ts"; import { fromWire, toWire } from "./priority.ts"; @@ -285,20 +285,34 @@ export class Publisher { const name = msg.trackNamespace; let broadcast: broadcast.Consumer | undefined; let refusal: { errorCode: number; reasonPhrase: string } | undefined; - try { - broadcast = - this.#publish && (wireOf(this.#publish).local(name) ?? (await wireOf(this.#publish).demand(name))); - if (!broadcast) { - refusal = { - errorCode: toRequestCode("does_not_exist", "subscribe", version), - reasonPhrase: "broadcast not found", - }; + + // Legal requests we can't honor are refused one at a time. A subscription that + // forwards nothing is only useful to a subscriber that later turns forwarding on, + // and serving a Range Filter unfiltered would deliver objects it excluded. + const unsupported = !msg.forward + ? "FORWARD=0 not supported" + : msg.rangeFilters + ? "range filters not supported" + : undefined; + + if (unsupported) { + refusal = { errorCode: toRequestCode("not_supported", "subscribe", version), reasonPhrase: unsupported }; + } else { + try { + broadcast = + this.#publish && (wireOf(this.#publish).local(name) ?? (await wireOf(this.#publish).demand(name))); + if (!broadcast) { + refusal = { + errorCode: toRequestCode("does_not_exist", "subscribe", version), + reasonPhrase: "broadcast not found", + }; + } + } catch (err: unknown) { + const e = error(err); + const condition = + e instanceof StreamError && e.code === StreamCode.NotFound ? "does_not_exist" : "internal"; + refusal = { errorCode: toRequestCode(condition, "subscribe", version), reasonPhrase: reason(e) }; } - } catch (err: unknown) { - const e = error(err); - const condition = - e instanceof StreamError && e.code === StreamCode.NotFound ? "does_not_exist" : "internal"; - refusal = { errorCode: toRequestCode(condition, "subscribe", version), reasonPhrase: reason(e) }; } if (refusal) { @@ -1254,22 +1268,41 @@ export class Publisher { * @internal */ async runTrackStatusRequest(msg: TrackStatusRequest, stream: Stream) { + // TRACK_STATUS_ERROR is 0x0f on draft-14. + await this.#refuseUnsupported(stream, msg.requestId, "track_status", 0x0f, "TRACK_STATUS is not supported"); + } + + /** + * Handles an incoming FETCH on a bidi stream. We serve none. + * + * @internal + */ + async runFetch(msg: Fetch, stream: Stream) { + await this.#refuseUnsupported(stream, msg.requestId, "fetch", FetchError.id, "FETCH is not supported"); + } + + /** + * Refuse a request NOT_SUPPORTED. Draft-14 gives each request its own error message, + * `errorId14`, with the SUBSCRIBE_ERROR body; later drafts use REQUEST_ERROR. + */ + async #refuseUnsupported( + stream: Stream, + requestId: bigint, + kind: RequestKind, + errorId14: number, + reasonPhrase: string, + ) { const version = this.#session.version; - const errorCode = toRequestCode("not_supported", "track_status", version); + const errorCode = toRequestCode("not_supported", kind, version); if (version === Version.DRAFT_14) { - // TRACK_STATUS_ERROR shares the SUBSCRIBE_ERROR body on draft-14. - await stream.writer.u53(0x0f); - await new SubscribeError({ - requestId: msg.requestId, - errorCode, - reasonPhrase: "TRACK_STATUS is not supported", - }).encode(stream.writer, version); + await stream.writer.u53(errorId14); + await new SubscribeError({ requestId, errorCode, reasonPhrase }).encode(stream.writer, version); } else { await stream.writer.u53(RequestError.id); await new RequestError({ - requestId: version === Version.DRAFT_15 || version === Version.DRAFT_16 ? msg.requestId : undefined, + requestId: version === Version.DRAFT_15 || version === Version.DRAFT_16 ? requestId : undefined, errorCode, - reasonPhrase: "TRACK_STATUS is not supported", + reasonPhrase, }).encode(stream.writer, version); } stream.close(); diff --git a/js/net/src/ietf/subscribe.ts b/js/net/src/ietf/subscribe.ts index 51973c7baa..84ec371890 100644 --- a/js/net/src/ietf/subscribe.ts +++ b/js/net/src/ietf/subscribe.ts @@ -38,6 +38,12 @@ export class Subscribe { /** Whether the subscriber wants Track Properties on the response (INCLUDE_PROPERTIES). */ propertiesWanted: boolean; + /** Whether Objects are forwarded (FORWARD). We only serve forwarding subscriptions. */ + forward: boolean; + + /** Whether the request carried a Range Filter, which we refuse. Never encoded. */ + rangeFilters: boolean; + constructor({ requestId, trackNamespace, @@ -46,6 +52,8 @@ export class Subscribe { filter, fill, propertiesWanted, + forward, + rangeFilters, }: { requestId: bigint; trackNamespace: Path.Valid; @@ -54,6 +62,8 @@ export class Subscribe { filter?: Filter.Filter; fill?: Filter.Fill; propertiesWanted?: boolean; + forward?: boolean; + rangeFilters?: boolean; }) { this.requestId = requestId; this.trackNamespace = trackNamespace; @@ -62,6 +72,8 @@ export class Subscribe { this.filter = filter ?? { kind: "unfiltered" }; this.fill = fill; this.propertiesWanted = propertiesWanted ?? true; + this.forward = forward ?? true; + this.rangeFilters = rangeFilters ?? false; } async #encode(w: Writer, version: IetfVersion): Promise { @@ -75,7 +87,7 @@ export class Subscribe { if (version === Version.DRAFT_14) { await w.u8(this.subscriberPriority); await w.u8(GROUP_ORDER); - await w.bool(true); // forward = true + await w.bool(this.forward); await w.write(Filter.encode(this.filter, version)); await w.u53(0); // no parameters } else { @@ -83,7 +95,7 @@ export class Subscribe { const params = new Parameters(); params.subscriberPriority = this.subscriberPriority; params.groupOrder = GROUP_ORDER; - params.forward = true; + params.forward = this.forward; params.subscriptionFilter = Filter.encode(this.filter, version); // FILL_PARAMETERS and INCLUDE_PROPERTIES arrived in draft-20. An older peer reads @@ -131,15 +143,11 @@ export class Subscribe { } const forward = await r.bool(); - if (!forward) { - throw new Error(`unsupported forward value: ${forward}`); - } - const filter = await Filter.decodeInline(r); await Parameters.decode(r, version); // ignore parameters - return new Subscribe({ requestId, trackNamespace, trackName, subscriberPriority, filter }); + return new Subscribe({ requestId, trackNamespace, trackName, subscriberPriority, filter, forward }); } // v15+: fields are in parameters const params = await Parameters.decode(r, version); @@ -152,18 +160,17 @@ export class Subscribe { groupOrder = GROUP_ORDER; // default to descending } - const forward = params.forward ?? true; - if (!forward) { - throw new Error(`unsupported forward value: ${forward}`); - } - - // FILL_PARAMETERS and INCLUDE_PROPERTIES are draft-20 additions. An unknown message - // parameter is a protocol violation, so they stay rejected on the drafts that predate - // them rather than being quietly tolerated. + // FILL_PARAMETERS and INCLUDE_PROPERTIES are draft-20 additions, and the Range + // Filters draft-19 ones. An unknown message parameter is a protocol violation, so + // they stay rejected on the drafts that predate them rather than being quietly + // tolerated. const draft20 = Filter.isDraft20(version); if ((params.fillParameters !== undefined || params.includeProperties !== undefined) && !draft20) { throw new Error("FILL_PARAMETERS and INCLUDE_PROPERTIES need draft-20"); } + if (params.rangeFilters && !Filter.hasRangeFilters(version)) { + throw new Error("Range Filters need draft-19"); + } // An absent LOCATION_FILTER means the subscription is unfiltered. const raw = params.subscriptionFilter; @@ -180,6 +187,8 @@ export class Subscribe { fill, // Defaults to 1, so an absent parameter means the subscriber wants them. propertiesWanted: params.includeProperties ?? true, + forward: params.forward ?? true, + rangeFilters: params.rangeFilters, }); } } diff --git a/js/net/src/ietf/track.ts b/js/net/src/ietf/track.ts index 668d3ad86d..3dbeb118ac 100644 --- a/js/net/src/ietf/track.ts +++ b/js/net/src/ietf/track.ts @@ -3,6 +3,7 @@ import type { Reader, Writer } from "../stream.ts"; import * as Message from "./message.ts"; import * as Namespace from "./namespace.ts"; import { Parameters } from "./parameters.ts"; +import { Subscribe } from "./subscribe.ts"; import { type IetfVersion, Version } from "./version.ts"; // we only support Group Order descending @@ -50,29 +51,12 @@ export class TrackStatusRequest { return Message.encode(w, (mw) => this.#encode(mw, version)); } + /** + * Every draft defines TRACK_STATUS as identical to SUBSCRIBE, so it decodes as one and + * keeps only what names the track. We refuse the request, so the rest goes unread. + */ static async decode(r: Reader, version: IetfVersion): Promise { - return Message.decode(r, (mr) => TrackStatusRequest.#decode(mr, version)); - } - - static async #decode(r: Reader, version: IetfVersion): Promise { - const requestId = await r.u62(); - if (version === Version.DRAFT_17) { - await r.u62(); // required_request_id_delta - } - const trackNamespace = await Namespace.decode(r); - const trackName = await r.string(); - - if (version === Version.DRAFT_14) { - await r.u8(); // subscriber_priority - await r.u8(); // group_order - await r.bool(); // forward - await r.u53(); // filter_type - await Parameters.decode(r, version); // parameters - } else { - // v15+: just parameters - await Parameters.decode(r, version); - } - + const { requestId, trackNamespace, trackName } = await Subscribe.decode(r, version); return new TrackStatusRequest({ requestId, trackNamespace, trackName }); } } diff --git a/quest/m0/README.md b/quest/m0/README.md index fae3a99c49..a78e28969e 100644 --- a/quest/m0/README.md +++ b/quest/m0/README.md @@ -41,7 +41,6 @@ Published API or wire breaks still land on dev; each quest's Plan says so. ## Required -- [Legal IETF input](/quest/m0/ietf-legal-input.md) - draft-20+ FETCH, allowed parameters, INCLUDE_PROPERTIES and FORWARD=0 decode and are refused per request, not session-fatal - [IETF FIN semantics](/quest/m0/ietf-fin-not-cancel.md) - a request stream FIN stops updates without cancelling, and REQUEST_UPDATE on a subscribe is parsed - [Subgroup refusal](/quest/m0/ietf-subgroup-refusal.md) - a non-zero moq-transport subgroup costs that one stream, never the session - [IETF stream types](/quest/m0/ietf-uni-stream-types.md) - padding streams are discarded stream-only and an unknown uni type closes the session, per draft-21 diff --git a/quest/m0/ietf-fin-not-cancel.md b/quest/m0/ietf-fin-not-cancel.md index ccae4e3d28..0f1d880919 100644 --- a/quest/m0/ietf-fin-not-cancel.md +++ b/quest/m0/ietf-fin-not-cancel.md @@ -26,10 +26,6 @@ silently ends the subscription. Public API: none. Wire: conformance fix; no draft change. -## Required - -- [Legal IETF input](/quest/m0/ietf-legal-input.md) - lands first; both change `ietf/fetch.rs` and the request close path - ## Related - [Lite request streams](/quest/m1/request-stream-serve.md) - lite deliberately treats a FIN as ending the request; don't unify the two diff --git a/quest/m0/ietf-legal-input.md b/quest/m0/ietf-legal-input.md deleted file mode 100644 index c90b5f547b..0000000000 --- a/quest/m0/ietf-legal-input.md +++ /dev/null @@ -1,49 +0,0 @@ -# [L] Legal IETF input never fails the session - -## Goal - -Every message a conformant draft-14 to draft-22 peer may send decodes, and -anything we don't support is refused per request, not by closing the session. -moxygen and libquicr send several of these today, and Seattle interop is -2026-10-12. - -## Plan - -Each of these closes the session today with `PROTOCOL_VIOLATION`: - -- FETCH on drafts 20+ still decodes the removed Fetch Type field - (`ietf/fetch.rs`). Decode the draft-20 layout: namespace, name and params, - with the range in LOCATION_FILTER. Refuse it `NOT_SUPPORTED` until - [moxygen FETCH](/quest/m1/moxygen/fetch.md) serves it. Fix the encoder and - the pinned test in `version.rs`. -- Request parameters the draft allows on a message fail `decode_params!`: - AUTHORIZATION TOKEN (0x03) anywhere, NEW_GROUP_REQUEST (0x32) and the - object and track filters (0x25 to 0x29) on SUBSCRIBE and REQUEST_UPDATE, - FILL_TIMEOUT (0x0A) and 0x21 and 0x35 on FETCH, and any parameter on TRACK_STATUS. Accept each - one where the draft allows it. The token is decoded and ignored, since the - session grant still applies; [Request token](/quest/m1/auth/request-token.md) - gives it meaning (decided 2026-09-29). An unknown key stays fatal, per the - draft. -- INCLUDE_PROPERTIES (0x35) is read as length-prefixed; the draft defines a - uint8 (`ietf/subscribe.rs`). Encode and decode it as `u8`, and delete the - doc comment arguing for key-parity framing: that rule covers Setup and - Properties KVPs, not message parameters. -- FORWARD=0 returns `DecodeError::Unsupported`. Decode it, then refuse the - request `NOT_SUPPORTED`. - -Not in scope: pre-draft-20 subscribe filters other than Largest Object stay -ignored. `older_drafts_are_ignored` in `ietf/publisher.rs` records that -choice, since honoring them would change what existing peers receive. - -- Tests: one decode test per message and draft, built from the draft's - message figures and bytes captured from moxygen and libquicr, plus a - session test showing each refusal leaves the session open. -- Mirror the decode in `js/net`. - -Public API: none. Wire: fixes conformance; no draft change. - -## Related - -- [Request token](/quest/m1/auth/request-token.md) - also edits `decode_params!`, turning the ignored token into a per-request grant -- [IETF FIN semantics](/quest/m0/ietf-fin-not-cancel.md) - the other interop blocker -- [IETF stream types](/quest/m0/ietf-uni-stream-types.md) - same stream-scoped-before-fatal rule for uni streams diff --git a/quest/m1/auth/request-token.md b/quest/m1/auth/request-token.md index f733630897..a171f7383b 100644 --- a/quest/m1/auth/request-token.md +++ b/quest/m1/auth/request-token.md @@ -11,9 +11,8 @@ does not cover it, by the token on the request; with neither it is refused `UNAUTHORIZED`. The token's grant covers only the request it rode on and lives exactly as long as that request, and a REQUEST_UPDATE carrying a new token replaces it, which is how a peer refreshes. It scopes by path, never -by method. [Legal IETF input](/quest/m0/ietf-legal-input.md) decodes and -ignores the key first, so a token no longer fails the session; this quest -gives it meaning. +by method. Every request already decodes the key and ignores it, so a token +no longer fails the session; this quest gives it meaning. ## Plan @@ -21,9 +20,9 @@ gives it meaning. (`rs/moq-net/src/ietf/token.rs`, `js/net/src/ietf/token.ts`): `USE_VALUE` yields the token, `REGISTER` is a value since we advertise no `MAX_AUTH_TOKEN_CACHE_SIZE`, and `DELETE` or `USE_ALIAS` closes with `PROTOCOL_VIOLATION`. Both decoder families change: the strict - `decode_params!` path, which ignores the key after - [Legal IETF input](/quest/m0/ietf-legal-input.md), and the generic KVP - path the legacy drafts use, which also ignores it. + `decode_params!` path, where each request reads the repeatable key into an + ignored `Vec`, and draft-14's `Parameters::skip`, which consumes it + unread. `js/net` keeps every instance in `Parameters` and reads none. - Fallback only: a request the session grant already covers is served without verifying its token. Otherwise its token becomes an `auth::Request` on the session's `auth::Handle`, the seam an AUTH stream's @@ -75,7 +74,3 @@ to) and on `moq_auth::Client` (the per-request lease). Wire: none new; the param - [Relay tokens](/quest/m1/auth/relay-refresh.md) - supplies the lease revalidation the per-request lease reuses - -## Related - -- [Legal IETF input](/quest/m0/ietf-legal-input.md) - also edits `decode_params!`; refresh this Plan when it lands, since the strict decoder then decodes and ignores the token diff --git a/quest/m1/moxygen/fetch.md b/quest/m1/moxygen/fetch.md index e38d1e533b..0398027380 100644 --- a/quest/m1/moxygen/fetch.md +++ b/quest/m1/moxygen/fetch.md @@ -12,7 +12,8 @@ subscriber sees. Standalone FETCH and a non-zero joining FETCH are refused today with "not supported". Those forms exist only before draft-20; on draft-20+ the -range comes from LOCATION_FILTER. Walk `track::Consumer::fetch_group`, one group, then the +range comes from LOCATION_FILTER, which decodes to `FetchType::Filtered` and +is refused today. A FETCH carrying Range Filters stays refused. Walk `track::Consumer::fetch_group`, one group, then the next. Do not add an archive. A joining FETCH is the same walk for the groups it names. A form the walk @@ -21,10 +22,6 @@ cannot express is still an explicit refusal, not a hang. The moxygen FETCH cases that ask for whole groups are the check. The rest of that suite is not. -## Required - -- [Legal IETF input](/quest/m0/ietf-legal-input.md) - decodes the draft-20+ FETCH layout this serves - ## Related - [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line this belongs to diff --git a/rs/moq-net/src/fuzz.rs b/rs/moq-net/src/fuzz.rs index a7dd026195..cbf344b69f 100644 --- a/rs/moq-net/src/fuzz.rs +++ b/rs/moq-net/src/fuzz.rs @@ -693,18 +693,29 @@ pub fn seeds() -> Vec { // FETCH is the one arm the sweep above cannot reach: its body ends in a parameter // count, and a uniform fill never lands a zero there. Build it from the encoder // instead, which also means a field change breaks the build rather than the seed. + // Draft-20 dropped the joining forms, so each draft gets whichever layout it has. for (index, version) in IETF_VERSIONS.iter().enumerate() { - let fetch = ietf::Fetch { - request_id: ietf::RequestId(0), - subscriber_priority: 128, - group_order: ietf::GroupOrder::Ascending, - fetch_type: ietf::FetchType::AbsoluteJoining { + let fetch_types = [ + ietf::FetchType::AbsoluteJoining { subscriber_request_id: ietf::RequestId(0), group_id: 0, }, - }; - - let Ok(encoded) = fetch.encode_bytes(*version) else { + ietf::FetchType::Filtered { + namespace: crate::Path::new("a"), + track: "b".into(), + filter: ietf::Filter::Relative(1), + }, + ]; + let Some(encoded) = fetch_types.into_iter().find_map(|fetch_type| { + let fetch = ietf::Fetch { + request_id: ietf::RequestId(0), + subscriber_priority: 128, + group_order: ietf::GroupOrder::Ascending, + fetch_type, + range_filters: false, + }; + fetch.encode_bytes(*version).ok() + }) else { continue; }; diff --git a/rs/moq-net/src/ietf/fetch.rs b/rs/moq-net/src/ietf/fetch.rs index 5987d8efbc..3e1da017af 100644 --- a/rs/moq-net/src/ietf/fetch.rs +++ b/rs/moq-net/src/ietf/fetch.rs @@ -4,8 +4,9 @@ use crate::{ Path, coding::{Decode, DecodeError, Encode, EncodeError}, ietf::{ - GroupOrder, Location, Parameters, RequestId, + Filter, GroupOrder, Location, Opaque, Parameters, RequestId, namespace::{decode_namespace, encode_namespace}, + subscribe::has_range_filters, }, }; @@ -13,9 +14,13 @@ use super::Message; use super::Version; +/// What a FETCH asks for. +/// +/// Through draft-19 a Fetch Type tag picks one of the first three. Draft-20 dropped the +/// tag and the joining forms, leaving [`Self::Filtered`]. #[derive(Debug, Clone, PartialEq, Eq)] pub enum FetchType<'a> { - // + /// An inclusive range of a track, through draft-19. Standalone { namespace: Path<'a>, track: Cow<'a, str>, @@ -30,6 +35,13 @@ pub enum FetchType<'a> { subscriber_request_id: RequestId, group_id: u64, }, + /// A track, with the range as a LOCATION_FILTER, from draft-20 on. Unfiltered is + /// everything from `{0, 0}` up to Largest Object. + Filtered { + namespace: Path<'a>, + track: Cow<'a, str>, + filter: Filter, + }, } impl Encode for FetchType<'_> { @@ -63,6 +75,8 @@ impl Encode for FetchType<'_> { subscriber_request_id.encode(w, version)?; group_id.encode(w, version)?; } + // Draft-20 has no Fetch Type tag to write. + FetchType::Filtered { .. } => return Err(EncodeError::Version), } Ok(()) } @@ -111,6 +125,10 @@ pub struct Fetch<'a> { pub subscriber_priority: u8, pub group_order: GroupOrder, pub fetch_type: FetchType<'a>, + /// Whether the request carried a Range Filter (0x25-0x28). We advertise no + /// MAX_FILTER_RANGES, so the request is refused rather than served unfiltered. + /// Never encoded; we send no range filters. + pub range_filters: bool, } impl Message for Fetch<'_> { @@ -129,13 +147,31 @@ impl Message for Fetch<'_> { self.fetch_type.encode(w, version)?; 0u8.encode(w, version)?; // no parameters } - _ => { + Version::Draft15 | Version::Draft16 | Version::Draft17 | Version::Draft18 | Version::Draft19 => { self.fetch_type.encode(w, version)?; encode_params!(w, version, 0x20 => self.subscriber_priority, 0x22 => self.group_order, ); } + _ => { + // The joining forms and the Fetch Type tag are gone in draft-20. + let FetchType::Filtered { + namespace, + track, + filter, + } = &self.fetch_type + else { + return Err(EncodeError::Version); + }; + encode_namespace(w, namespace, version)?; + track.encode(w, version)?; + encode_params!(w, version, + 0x20 => self.subscriber_priority, + 0x21 => *filter, + 0x22 => self.group_order, + ); + } } Ok(()) } @@ -146,37 +182,88 @@ impl Message for Fetch<'_> { let _required_request_id_delta = u64::decode(buf, version)?; } - match version { + // The token is ignored: the session's grant is what authorizes the request. We + // refuse or serve a FETCH the same whatever its FILL_TIMEOUT or INCLUDE_PROPERTIES. + let (fetch_type, subscriber_priority, group_order, range_filters) = match version { Version::Draft14 => { let subscriber_priority = u8::decode(buf, version)?; let group_order = GroupOrder::decode(buf, version)?; let fetch_type = FetchType::decode(buf, version)?; - let _params = Parameters::decode(buf, version)?; - Ok(Self { - request_id, - subscriber_priority, - group_order, - fetch_type, - }) + Parameters::skip(buf, version)?; + (fetch_type, Some(subscriber_priority), Some(group_order), false) } - _ => { + Version::Draft15 | Version::Draft16 | Version::Draft17 | Version::Draft18 | Version::Draft19 => { let fetch_type = FetchType::decode(buf, version)?; decode_params!(buf, version, + 0x03 => _authorization_token: Vec, + 0x0A => fill_timeout: Option, 0x20 => subscriber_priority: Option, 0x22 => group_order: Option, + 0x25 => subgroup_filter: Vec, + 0x26 => object_id_filter: Vec, + 0x27 => priority_filter: Vec, + 0x28 => object_property_filter: Vec, ); + let range_filters = [ + subgroup_filter, + object_id_filter, + priority_filter, + object_property_filter, + ] + .iter() + .any(|filter| !filter.is_empty()); + + // An unknown message parameter is a protocol violation, so each stays rejected + // on the drafts that predate it: FILL_TIMEOUT arrived in draft-18. + let has_fill_timeout = !matches!(version, Version::Draft15 | Version::Draft16 | Version::Draft17); + if (fill_timeout.is_some() && !has_fill_timeout) || (range_filters && !has_range_filters(version)) { + return Err(DecodeError::InvalidValue); + } - let subscriber_priority = subscriber_priority.unwrap_or(128); - let group_order = group_order.unwrap_or(GroupOrder::Descending); - - Ok(Self { - request_id, - subscriber_priority, - group_order, - fetch_type, - }) + (fetch_type, subscriber_priority, group_order, range_filters) } - } + // Draft-20 names the track up front and moves the range into LOCATION_FILTER. + _ => { + let namespace = decode_namespace(buf, version)?; + let track = Cow::::decode(buf, version)?; + decode_params!(buf, version, + 0x03 => _authorization_token: Vec, + 0x0A => _fill_timeout: Option, + 0x20 => subscriber_priority: Option, + 0x21 => filter: Option, + 0x22 => group_order: Option, + 0x25 => subgroup_filter: Vec, + 0x26 => object_id_filter: Vec, + 0x27 => priority_filter: Vec, + 0x28 => object_property_filter: Vec, + 0x35 => _include_properties: Option, + ); + let range_filters = [ + subgroup_filter, + object_id_filter, + priority_filter, + object_property_filter, + ] + .iter() + .any(|filter| !filter.is_empty()); + + let fetch_type = FetchType::Filtered { + namespace, + track, + // An absent LOCATION_FILTER fetches the whole track. + filter: filter.unwrap_or(Filter::Unfiltered), + }; + (fetch_type, subscriber_priority, group_order, range_filters) + } + }; + + Ok(Self { + request_id, + subscriber_priority: subscriber_priority.unwrap_or(128), + group_order: group_order.unwrap_or(GroupOrder::Descending), + fetch_type, + range_filters, + }) } } @@ -229,7 +316,7 @@ impl Message for FetchOk { let group_order = GroupOrder::decode(buf, version)?; let end_of_track = bool::decode(buf, version)?; let end_location = Location::decode(buf, version)?; - let _params = Parameters::decode(buf, version)?; + Parameters::skip(buf, version)?; Ok(Self { request_id, group_order, @@ -556,6 +643,7 @@ mod tests { start: Location { group: 0, object: 0 }, end: Location { group: 10, object: 5 }, }, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft14); @@ -577,6 +665,7 @@ mod tests { start: Location { group: 0, object: 0 }, end: Location { group: 10, object: 5 }, }, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft15); @@ -615,6 +704,7 @@ mod tests { start: Location { group: 0, object: 0 }, end: Location { group: 10, object: 5 }, }, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft16); @@ -636,6 +726,7 @@ mod tests { start: Location { group: 0, object: 0 }, end: Location { group: 10, object: 5 }, }, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft17); @@ -708,6 +799,7 @@ mod tests { start: Location { group: 0, object: 0 }, end: Location { group: 10, object: 5 }, }, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft18); @@ -755,6 +847,108 @@ mod tests { ]; assert_eq!(encode_message(&msg, Version::Draft18), expected); } + + /// The head of a draft-20 FETCH (Figure 16) up to its Number of Parameters: Request + /// ID 1, Track Namespace ("live"), Track Name ("video"). There is no Fetch Type. + const FETCH_HEAD: &[u8] = &[ + 0x01, 0x01, 0x04, b'l', b'i', b'v', b'e', 0x05, b'v', b'i', b'd', b'e', b'o', + ]; + + /// A draft-20 FETCH names the track up front and carries its range in LOCATION_FILTER. + /// Reading a Fetch Type there instead takes the namespace's field count for one. + #[test] + fn test_fetch_v20_decodes_every_parameter() { + #[rustfmt::skip] + let body = [FETCH_HEAD, &[ + 0x07, // Number of Parameters + 0x03, 0x03, 0x03, 0x00, 0xAA, // AUTHORIZATION TOKEN + 0x07, 0x64, // FILL_TIMEOUT (0x0A) = 100, a varint + 0x16, 0x40, // SUBSCRIBER_PRIORITY (0x20) = 64 + 0x01, 0x03, 0x04, 0x00, 0x02, // LOCATION_FILTER (0x21): groups 4 through 6 + 0x01, 0x01, // GROUP_ORDER (0x22) = Ascending + 0x04, 0x02, 0x00, 0x05, // OBJECTID_FILTER (0x26): SetID 0, from 5 + 0x0F, 0x00, // INCLUDE_PROPERTIES (0x35) = 0 + ]].concat(); + + for version in [Version::Draft20, Version::Draft21, Version::Draft22] { + let fetch: Fetch = decode_message(&body, version).unwrap_or_else(|e| panic!("{version}: {e}")); + assert_eq!(fetch.request_id, RequestId(1)); + assert_eq!(fetch.subscriber_priority, 64); + assert_eq!(fetch.group_order, GroupOrder::Ascending); + assert!(fetch.range_filters, "{version}"); + assert_eq!( + fetch.fetch_type, + FetchType::Filtered { + namespace: Path::new("live"), + track: "video".into(), + filter: Filter::Absolute { + start: Location { group: 4, object: 0 }, + end: Some(crate::ietf::EndLocation { group: 6, object: None }), + }, + }, + "{version}" + ); + } + } + + /// With no LOCATION_FILTER a draft-20 FETCH is the whole track, and that is what we + /// write back: the parameter is omitted rather than sent empty. + #[test] + fn test_fetch_v20_wire() { + let fetch = Fetch { + request_id: RequestId(1), + subscriber_priority: 128, + group_order: GroupOrder::Descending, + fetch_type: FetchType::Filtered { + namespace: Path::new("live"), + track: "video".into(), + filter: Filter::Unfiltered, + }, + range_filters: false, + }; + + #[rustfmt::skip] + let expected = [FETCH_HEAD, &[ + 0x02, // Number of Parameters + 0x20, 0x80, // SUBSCRIBER_PRIORITY = 128 + 0x02, 0x02, // GROUP_ORDER (0x22) = Descending + ]].concat(); + assert_eq!(encode_message(&fetch, Version::Draft20), expected); + assert_eq!(decode_message::(&expected, Version::Draft20).unwrap(), fetch); + + // The tagged forms have no spelling from draft-20 on, and this one none before it. + let mut buf = Vec::new(); + assert!(fetch.encode_msg(&mut buf, Version::Draft19).is_err()); + } + + /// FILL_TIMEOUT arrived in draft-18 and the Range Filters in draft-19; each is still + /// an unknown parameter before then. + #[test] + fn test_fetch_parameters_follow_their_draft() { + #[rustfmt::skip] + let joining = |params: &[u8]| [&[ + 0x01, // Request ID + 0x02, 0x03, 0x00, // Relative Joining: subscription 3, offset 0 + ][..], params].concat(); + + let token = joining(&[0x01, 0x03, 0x03, 0x03, 0x00, 0xAA]); + let fill_timeout = joining(&[0x01, 0x0A, 0x20]); + let range_filter = joining(&[0x01, 0x25, 0x00]); + + for (version, body, ok) in [ + (Version::Draft15, &token, true), + (Version::Draft16, &fill_timeout, false), + (Version::Draft18, &fill_timeout, true), + (Version::Draft18, &range_filter, false), + (Version::Draft19, &range_filter, true), + ] { + let decoded = decode_message::(body, version); + assert_eq!(decoded.is_ok(), ok, "{version}: {body:x?}"); + if let Ok(fetch) = decoded { + assert_eq!(fetch.range_filters, body == &range_filter, "{version}"); + } + } + } } /// The Object serialization on a fetch stream (draft-20 section 11.4.4), which a fill's diff --git a/rs/moq-net/src/ietf/parameters.rs b/rs/moq-net/src/ietf/parameters.rs index 7b3b426f88..2e7330a801 100644 --- a/rs/moq-net/src/ietf/parameters.rs +++ b/rs/moq-net/src/ietf/parameters.rs @@ -211,6 +211,39 @@ impl Encode for Parameters { } impl Parameters { + /// Consume a draft-14 message parameter block without interpreting it. + /// + /// Draft-14 section 9.2 has a receiver ignore unrecognized parameters and allow their + /// duplicates, and lets AUTHORIZATION TOKEN repeat. We act on none of these, so unlike + /// [`Parameters::decode`], which SETUP uses, a repeat is not refused. + pub fn skip(r: &mut R, version: Version) -> Result<(), DecodeError> { + let count = u64::decode(r, version)?; + if count > MAX_PARAMS { + return Err(DecodeError::TooMany); + } + + for _ in 0..count { + // Parity frames a Key-Value-Pair: even is one varint, odd is length prefixed. + match u64::decode(r, version)? % 2 { + 0 => { + u64::decode(r, version)?; + } + _ => { + let len = usize::try_from(u64::decode(r, version)?).map_err(|_| DecodeError::BoundsExceeded)?; + if len > MAX_KVP_VALUE_LEN { + return Err(DecodeError::BoundsExceeded); + } + if r.remaining() < len { + return Err(DecodeError::Short); + } + r.advance(len); + } + } + } + + Ok(()) + } + pub fn get_varint(&self, kind: ParameterVarInt) -> Option { self.vars.get(&kind).copied() } @@ -245,6 +278,14 @@ pub trait Param: Sized { fn param_present(&self) -> bool { true } + + /// Fold a repeat of this parameter into the value already decoded. + /// + /// A parameter may appear once unless its definition says otherwise, so the default + /// refuses the repeat. + fn param_repeat(self, _next: Self) -> Result { + Err(DecodeError::Duplicate) + } } impl Param for u8 { @@ -353,6 +394,13 @@ impl Param for Option { self.is_some() } + fn param_repeat(self, next: Self) -> Result { + match (self, next) { + (Some(prev), Some(next)) => Ok(Some(prev.param_repeat(next)?)), + _ => Err(DecodeError::Duplicate), + } + } + fn param_encode(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { match self { Some(v) => v.param_encode(w, version), @@ -365,6 +413,45 @@ impl Param for Option { } } +/// A parameter whose definition lets it repeat, such as AUTHORIZATION TOKEN (0x03) or a +/// Range Filter (0x25-0x29). Every instance is kept, in wire order. +impl Param for Vec { + fn param_present(&self) -> bool { + !self.is_empty() + } + + fn param_encode(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { + // `encode_params!` writes the key once, so only a single instance fits behind it. + match self.as_slice() { + [value] => value.param_encode(w, version), + _ => Err(EncodeError::Unsupported), + } + } + + fn param_decode(r: &mut R, version: Version) -> Result { + Ok(vec![T::param_decode(r, version)?]) + } + + fn param_repeat(mut self, next: Self) -> Result { + self.extend(next); + Ok(self) + } +} + +/// A length-prefixed parameter value, consumed without being interpreted. +#[derive(Clone, Debug, Default, PartialEq, Eq)] +pub struct Opaque(pub Vec); + +impl Param for Opaque { + fn param_encode(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { + self.0.encode(w, version) + } + + fn param_decode(r: &mut R, version: Version) -> Result { + Ok(Self(Vec::::decode(r, version)?)) + } +} + /// Encode message parameters with compile-time sorted keys. /// /// Keys must be listed in ascending order (enforced at compile time). @@ -429,7 +516,8 @@ macro_rules! encode_params { /// for parameters where `T::default()` is an acceptable fallback. /// /// Unknown parameters cause `DecodeError::InvalidValue`. -/// Duplicate parameters cause `DecodeError::Duplicate`. +/// Duplicate parameters cause `DecodeError::Duplicate`, unless the type allows a repeat +/// (see [`Param::param_repeat`]), such as `Vec`. /// /// ```ignore /// decode_params!(r, version, @@ -487,10 +575,11 @@ macro_rules! decode_params { // the macro captures it as an expression, which is not a legal pattern. $( if _key == $key { - if $name.is_some() { - return Err($crate::coding::DecodeError::Duplicate); - } - $name = Some(<$ty as $crate::ietf::Param>::param_decode($r, _version)?); + let _value = <$ty as $crate::ietf::Param>::param_decode($r, _version)?; + $name = Some(match $name.take() { + None => _value, + Some(_prev) => $crate::ietf::Param::param_repeat(_prev, _value)?, + }); continue; } )* @@ -975,4 +1064,46 @@ mod tests { ); } } + + /// A parameter the draft lets repeat, such as AUTHORIZATION TOKEN, keeps every + /// instance instead of failing the message as a duplicate. + #[test] + fn test_param_repeat_allowed() { + for version in [Version::Draft15, Version::Draft16, Version::Draft17, Version::Draft20] { + let mut buf = BytesMut::new(); + 2usize.encode(&mut buf, version).unwrap(); + 0x03u64.encode(&mut buf, version).unwrap(); + Opaque(vec![0xAA]).param_encode(&mut buf, version).unwrap(); + // The second key: absolute before draft-16, a zero delta after. + let second: u64 = if version == Version::Draft15 { 0x03 } else { 0 }; + second.encode(&mut buf, version).unwrap(); + Opaque(vec![0xBB]).param_encode(&mut buf, version).unwrap(); + + let mut bytes = buf.freeze(); + let tokens = (|| -> Result, DecodeError> { + decode_params!(&mut bytes, version, 0x03 => tokens: Vec); + Ok(tokens) + })() + .unwrap_or_else(|e| panic!("{version}: {e}")); + assert_eq!(tokens, vec![Opaque(vec![0xAA]), Opaque(vec![0xBB])], "{version}"); + assert!(!bytes.has_remaining(), "{version}"); + } + } + + /// Draft-14 lets AUTHORIZATION TOKEN repeat and has unknown parameters, duplicates + /// included, ignored. The block is consumed whole either way. + #[test] + fn test_skip_allows_draft14_repeats() { + #[rustfmt::skip] + let block = [ + 0x04, // Number of Parameters + 0x03, 0x01, 0xAA, // AUTHORIZATION TOKEN + 0x03, 0x01, 0xBB, // and again + 0x3E, 0x05, // an unknown varint parameter + 0x3E, 0x06, // and again + ]; + let mut buf = &block[..]; + Parameters::skip(&mut buf, Version::Draft14).unwrap(); + assert!(buf.is_empty()); + } } diff --git a/rs/moq-net/src/ietf/publish.rs b/rs/moq-net/src/ietf/publish.rs index cf93727fd1..706fa5a16f 100644 --- a/rs/moq-net/src/ietf/publish.rs +++ b/rs/moq-net/src/ietf/publish.rs @@ -309,7 +309,7 @@ impl Message for Publish<'_> { }; let forward = bool::decode(r, version)?; // parameters - let _params = Parameters::decode(r, version)?; + Parameters::skip(r, version)?; Ok(Self { request_id, @@ -334,6 +334,7 @@ impl Message for Publish<'_> { // letting the request reach its NOT_SUPPORTED response. decode_params!(r, version, 0x02 => object_delivery_timeout: Option, + 0x03 => _authorization_token: Vec, 0x06 => subgroup_delivery_timeout: Option, 0x08 => _expires: Option, 0x09 => largest_location: Option, @@ -437,7 +438,7 @@ impl Message for PublishOk { let filter = Filter::decode(r, version)?; // no parameters - let _params = Parameters::decode(r, version)?; + Parameters::skip(r, version)?; Ok(Self { request_id, diff --git a/rs/moq-net/src/ietf/publish_namespace.rs b/rs/moq-net/src/ietf/publish_namespace.rs index 863826dfeb..ae090688bd 100644 --- a/rs/moq-net/src/ietf/publish_namespace.rs +++ b/rs/moq-net/src/ietf/publish_namespace.rs @@ -36,7 +36,24 @@ impl PublishNamespace<'_> { let _required_request_id_delta = u64::decode(r, version)?; } let track_namespace = decode_namespace(r, version)?; - let cluster = decode_cluster_params(r, version, negotiated)?; + + // The token is ignored: the session's grant is what authorizes the request. + decode_params!(r, version, + 0x03 => _authorization_token: Vec, + cluster::HOP_PATH => hops: Option, + cluster::ROUTE_COST => cost: Option, + ); + + let cluster = match negotiated { + true => Some(cluster::Advert { + hops: hops.ok_or(DecodeError::InvalidValue)?, + cost: cost.unwrap_or(0), + }), + // An endpoint must not append these on a session that did not negotiate the + // extension, so either is a violation. + false if hops.is_some() || cost.is_some() => return Err(DecodeError::InvalidValue), + false => None, + }; Ok(Self { request_id, @@ -127,7 +144,9 @@ impl Message for PublishNamespaceUpdate { } _ => RequestId::decode(r, version)?, }; + // The token is ignored: the session's grant is what authorizes the request. decode_params!(r, version, + 0x03 => _authorization_token: Vec, cluster::HOP_PATH => hops: Option, cluster::ROUTE_COST => cost: Option, ); @@ -157,28 +176,21 @@ pub(super) fn encode_cluster_params( Ok(()) } -/// Read the Parameters field of an advertisement. See [`encode_cluster_params`]. +/// Read the Parameters field of a NAMESPACE on a session that negotiated the extension. +/// See [`encode_cluster_params`]. pub(super) fn decode_cluster_params( r: &mut R, version: Version, - negotiated: bool, -) -> Result, DecodeError> { - if !negotiated { - // An endpoint must not append these on a session that did not negotiate the - // extension, and we know no other parameter here, so any is a violation. - decode_params!(r, version,); - return Ok(None); - } - +) -> Result { decode_params!(r, version, cluster::HOP_PATH => hops: Option, cluster::ROUTE_COST => cost: Option, ); - Ok(Some(cluster::Advert { + Ok(cluster::Advert { hops: hops.ok_or(DecodeError::InvalidValue)?, cost: cost.unwrap_or(0), - })) + }) } /// PublishNamespaceOk message (0x07) diff --git a/rs/moq-net/src/ietf/publisher.rs b/rs/moq-net/src/ietf/publisher.rs index b6bbe92179..c11b8881d4 100644 --- a/rs/moq-net/src/ietf/publisher.rs +++ b/rs/moq-net/src/ietf/publisher.rs @@ -474,6 +474,21 @@ where tracing::info!(id = %request_id, broadcast = %absolute, track = %track_name, "subscribe started"); + // Legal requests we can't honor are refused one at a time, never by closing the + // session. A subscription that forwards nothing is only useful to a subscriber + // that later turns forwarding on, and serving a Range Filter unfiltered would + // deliver objects the subscriber excluded. + if !msg.forward { + return self + .reject_subscribe(stream, request_id, &Error::Unsupported, "FORWARD=0 not supported") + .await; + } + if msg.range_filters { + return self + .reject_subscribe(stream, request_id, &Error::Unsupported, "range filters not supported") + .await; + } + // Stats (subscriptions, viewer refcount, groups/frames/bytes) are counted in // the model, through the tagged `origin::Consumer` the broadcast resolves from. @@ -1029,13 +1044,21 @@ where /// Serve the current-group prefix for a relative joining FETCH with offset zero. async fn run_fetch_stream(mut self, mut stream: Stream, msg: ietf::Fetch<'_>) -> Result<(), Error> { + // Draft-20 removed joining FETCH, the only form we serve. if Filter::is_draft20(self.version) { + return self + .reject_fetch(stream, msg.request_id, &Error::Unsupported, "FETCH not supported") + .await; + } + + // Serving a Range Filter unfiltered would deliver objects the subscriber excluded. + if msg.range_filters { return self .reject_fetch( stream, msg.request_id, &Error::Unsupported, - "joining FETCH removed in draft-20", + "range filters not supported", ) .await; } @@ -1057,7 +1080,7 @@ where } subscriber_request_id } - FetchType::AbsoluteJoining { .. } => { + FetchType::AbsoluteJoining { .. } | FetchType::Filtered { .. } => { return self .reject_fetch(stream, msg.request_id, &Error::Unsupported, "not supported") .await; @@ -2803,6 +2826,8 @@ mod serve_tests { filter, fill, properties_wanted: true, + forward: true, + range_filters: false, } } @@ -2975,6 +3000,100 @@ mod serve_tests { } } + /// Dispatch one request stream, returning the refusal code it wrote. + /// + /// `handle_stream` returning `Ok` is what keeps the session open: the dispatch loop + /// closes the session over an `Err`. + async fn refusal(version: Version, id: u64, body: Vec) -> u64 { + let h = serve(version); + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + let mark = h.log.writes.lock().unwrap().len(); + h.publisher + .clone() + .handle_stream(id, bytes::Bytes::from(body), stream) + .unwrap_or_else(|e| panic!("{version}: the request closed the session: {e}")) + .await; + assert!(h.log.resets().is_empty(), "{version}: refusal was reset"); + + let mut buf = bytes::Bytes::from(h.log.writes.lock().unwrap()[mark..].to_vec()); + match version { + Version::Draft14 => { + // SUBSCRIBE_ERROR and FETCH_ERROR share their layout. + let _id = u64::decode(&mut buf, version).unwrap(); + ietf::SubscribeError::decode(&mut buf, version).unwrap().error_code + } + _ => { + assert_eq!(u64::decode(&mut buf, version).unwrap(), ietf::RequestError::ID); + ietf::RequestError::decode(&mut buf, version).unwrap().error_code + } + } + } + + /// Legal requests we don't serve are refused NOT_SUPPORTED one at a time, and the + /// session stays open for the next one. + #[tokio::test] + async fn legal_requests_we_do_not_serve_are_refused_per_request() { + const NOT_SUPPORTED: u64 = 0x3; + + for version in [ + Version::Draft14, + Version::Draft15, + Version::Draft16, + Version::Draft17, + Version::Draft18, + Version::Draft19, + Version::Draft20, + Version::Draft21, + Version::Draft22, + ] { + let mut paused = subscribe(Filter::NextObject, None); + paused.forward = false; + let mut body = bytes::BytesMut::new(); + paused.encode_msg(&mut body, version).unwrap(); + assert_eq!( + refusal(version, ietf::Subscribe::ID, body.to_vec()).await, + NOT_SUPPORTED, + "{version}: SUBSCRIBE with FORWARD=0" + ); + } + + // Draft-20 bytes from the figures: Request ID 0x2B, Track Namespace ("room"), + // Track Name ("video"), then the parameters. + let head: &[u8] = &[ + 0x2B, 0x01, 0x04, b'r', b'o', b'o', b'm', 0x05, b'v', b'i', b'd', b'e', b'o', + ]; + + #[rustfmt::skip] + let cases: [(&str, u64, &[u8]); 3] = [ + ("SUBSCRIBE with a Range Filter", ietf::Subscribe::ID, &[ + 0x02, // Number of Parameters + 0x03, 0x03, 0x03, 0x00, 0xAA, // AUTHORIZATION TOKEN + 0x23, 0x02, 0x00, 0x05, // OBJECTID_FILTER (0x26): SetID 0, from 5 + ]), + ("FETCH", ietf::Fetch::ID, &[ + 0x03, // Number of Parameters + 0x0A, 0x00, // FILL_TIMEOUT = 0 + 0x17, 0x01, 0x01, // LOCATION_FILTER (0x21): from one group back + 0x14, 0x01, // INCLUDE_PROPERTIES (0x35) = 1 + ]), + ("TRACK_STATUS with parameters", ietf::TrackStatus::ID, &[ + 0x02, // Number of Parameters + 0x03, 0x03, 0x03, 0x00, 0xAA, // AUTHORIZATION TOKEN + 0x32, 0x00, // INCLUDE_PROPERTIES (0x35) = 0 + ]), + ]; + + for version in [Version::Draft20, Version::Draft21, Version::Draft22] { + for (label, id, params) in cases { + assert_eq!( + refusal(version, id, [head, params].concat()).await, + NOT_SUPPORTED, + "{version}: {label}" + ); + } + } + } + /// The draft's canonical current-group join: a Next Object subscription plus a /// StartGroup=1 fill. The published head arrives exactly once, on a fetch stream, /// and the subscription starts past the snapshot, so nothing is duplicated and @@ -3149,6 +3268,7 @@ mod serve_tests { subscriber_request_id: RequestId(REQUEST_ID), group_offset: 0, }, + range_filters: false, }, ) .await?; @@ -5024,6 +5144,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: true, + forward: true, + range_filters: false, }, ) .await @@ -5047,6 +5169,7 @@ mod tests { subscriber_priority: 128, group_order: GroupOrder::Descending, fetch_type, + range_filters: false, }, ) .await @@ -5377,6 +5500,8 @@ mod range_tests { filter, fill: None, properties_wanted: true, + forward: true, + range_filters: false, } } diff --git a/rs/moq-net/src/ietf/subscribe.rs b/rs/moq-net/src/ietf/subscribe.rs index 07facd099e..83be556125 100644 --- a/rs/moq-net/src/ietf/subscribe.rs +++ b/rs/moq-net/src/ietf/subscribe.rs @@ -5,7 +5,7 @@ use std::borrow::Cow; use crate::{ Path, coding::*, - ietf::{Fill, Filter, GroupOrder, Location, Param, Parameters, Properties, RequestId}, + ietf::{Fill, Filter, GroupOrder, Location, Opaque, Parameters, Properties, RequestId}, }; use super::Message; @@ -13,31 +13,6 @@ use super::namespace::{decode_namespace, encode_namespace}; use super::Version; -/// The INCLUDE_PROPERTIES parameter (0x35), draft-20's opt-out from Track Properties. -/// -/// Length prefixed despite holding a single byte. The Key-Value-Pair rule keys the framing -/// off the parameter id's parity and 0x35 is odd, so a Length is present even though the -/// parameter's own section calls the value a uint8. Parity is what the generic parser uses -/// to skip a parameter it does not know, so following the prose instead would desync every -/// parameter after this one. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub struct IncludeProperties(pub bool); - -impl Param for IncludeProperties { - fn param_encode(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { - vec![u8::from(self.0)].encode(w, version) - } - - fn param_decode(r: &mut R, version: Version) -> Result { - match Vec::::decode(r, version)?[..] { - // The draft allows exactly 0 or 1; anything else is a protocol violation. - [0] => Ok(Self(false)), - [1] => Ok(Self(true)), - _ => Err(DecodeError::InvalidValue), - } - } -} - /// Subscribe message (0x03) /// Sent by the subscriber to request all future objects for the given track. #[derive(Clone, Debug)] @@ -53,6 +28,26 @@ pub struct Subscribe<'a> { pub fill: Option, /// Whether the subscriber wants Track Properties on the response (draft-20). pub properties_wanted: bool, + /// Whether Objects are forwarded (FORWARD). We only serve forwarding subscriptions, + /// so a subscriber that asks for none is refused. + pub forward: bool, + /// Whether the request carried a Range Filter (0x25-0x28). We advertise no + /// MAX_FILTER_RANGES, so the request is refused rather than served unfiltered. + /// Never encoded; we send no range filters. + pub range_filters: bool, +} + +/// NEW_GROUP_REQUEST (0x32) arrived in draft-16. +fn has_new_group_request(version: Version) -> bool { + !matches!(version, Version::Draft14 | Version::Draft15) +} + +/// The Range Filters (0x25-0x29) arrived in draft-19. +pub(super) fn has_range_filters(version: Version) -> bool { + !matches!( + version, + Version::Draft14 | Version::Draft15 | Version::Draft16 | Version::Draft17 | Version::Draft18 + ) } impl Message for Subscribe<'_> { @@ -72,13 +67,9 @@ impl Message for Subscribe<'_> { let group_order = GroupOrder::decode(r, version)?; let forward = bool::decode(r, version)?; - if !forward { - return Err(DecodeError::Unsupported); - } - let filter = Filter::decode(r, version)?; - let _params = Parameters::decode(r, version)?; + Parameters::skip(r, version)?; Ok(Self { request_id, @@ -89,11 +80,17 @@ impl Message for Subscribe<'_> { filter, fill: None, properties_wanted: true, + forward, + range_filters: false, }) } _ => { + // The token is ignored: the session's grant is what authorizes the request. + // NEW_GROUP_REQUEST is ignored too, as the draft lets a publisher without + // dynamic groups do. decode_params!(r, version, 0x02 => _object_delivery_timeout: Option, + 0x03 => _authorization_token: Vec, 0x04 => rendezvous_timeout: Option, 0x06 => _subgroup_delivery_timeout: Option, 0x10 => forward: Option, @@ -101,18 +98,34 @@ impl Message for Subscribe<'_> { 0x21 => filter: Option, 0x22 => group_order: Option, 0x23 => fill: Option, - 0x35 => include_properties: Option, + 0x25 => subgroup_filter: Vec, + 0x26 => object_id_filter: Vec, + 0x27 => priority_filter: Vec, + 0x28 => object_property_filter: Vec, + 0x32 => new_group_request: Option, + 0x35 => include_properties: Option, ); - // FILL_PARAMETERS and INCLUDE_PROPERTIES are draft-20 additions. An unknown - // message parameter is a protocol violation, so they stay rejected on the - // drafts that predate them rather than being quietly tolerated. - if (fill.is_some() || include_properties.is_some()) && !Filter::is_draft20(version) { + let range_filters = [ + subgroup_filter, + object_id_filter, + priority_filter, + object_property_filter, + ] + .iter() + .any(|filter| !filter.is_empty()); + + // An unknown message parameter is a protocol violation, so each stays rejected + // on the drafts that predate it rather than being quietly tolerated. + if ((fill.is_some() || include_properties.is_some()) && !Filter::is_draft20(version)) + || (range_filters && !has_range_filters(version)) + || (new_group_request.is_some() && !has_new_group_request(version)) + { return Err(DecodeError::InvalidValue); } // Defaults to 1, so an absent parameter means the subscriber wants them. - let properties_wanted = include_properties.is_none_or(|p| p.0); + let properties_wanted = include_properties.unwrap_or(true); // RENDEZVOUS_TIMEOUT arrived in draft-17; 0x04 means MAX_CACHE_DURATION in // draft-15, which is a publisher parameter with no business in a SUBSCRIBE. @@ -126,10 +139,6 @@ impl Message for Subscribe<'_> { // We still have to parse it, or the parameter alone would kill the session. let _ = rendezvous_timeout; - if forward == Some(false) { - return Err(DecodeError::Unsupported); - } - let subscriber_priority = subscriber_priority.unwrap_or(128); let group_order = group_order.unwrap_or(GroupOrder::Descending); // An absent LOCATION_FILTER means the subscription is unfiltered. @@ -144,6 +153,9 @@ impl Message for Subscribe<'_> { filter, fill, properties_wanted, + // Absent means 1. + forward: forward.unwrap_or(true), + range_filters, }) } } @@ -161,7 +173,7 @@ impl Message for Subscribe<'_> { Version::Draft14 => { self.subscriber_priority.encode(w, version)?; self.group_order.encode(w, version)?; - true.encode(w, version)?; // forward + self.forward.encode(w, version)?; self.filter.encode(w, version)?; 0u8.encode(w, version)?; // no parameters @@ -175,11 +187,10 @@ impl Message for Subscribe<'_> { // INCLUDE_PROPERTIES defaults to 1, so only the opt-out is worth bytes. It // arrived in draft-20, and an older peer would read it as an unknown // parameter, which is a protocol violation. - let include_properties = - (!self.properties_wanted && Filter::is_draft20(version)).then_some(IncludeProperties(false)); + let include_properties = (!self.properties_wanted && Filter::is_draft20(version)).then_some(false); encode_params!(w, version, - 0x10 => true, + 0x10 => self.forward, 0x20 => self.subscriber_priority, 0x21 => self.filter, 0x22 => self.group_order, @@ -279,7 +290,7 @@ impl Message for SubscribeOk { largest = Some(Location::decode(r, version)?); } - let _params = Parameters::decode(r, version)?; + Parameters::skip(r, version)?; } _ => { // GROUP_ORDER is only legal here through draft-15, but keep accepting it so a @@ -432,7 +443,7 @@ impl Message for SubscribeUpdate { let end_group = u64::decode(r, version)?; let subscriber_priority = u8::decode(r, version)?; let forward = bool::decode(r, version)?; - let _parameters = Parameters::decode(r, version)?; + Parameters::skip(r, version)?; Ok(Self { request_id, @@ -448,12 +459,19 @@ impl Message for SubscribeUpdate { let subscription_request_id = Some(RequestId::decode(r, version)?); decode_params!(r, version, 0x02 => _object_delivery_timeout: Option, + 0x03 => _authorization_token: Vec, 0x06 => _subgroup_delivery_timeout: Option, 0x10 => forward: Option, 0x20 => subscriber_priority: Option, 0x21 => _filter: Option, + 0x32 => new_group_request: Option, ); + // NEW_GROUP_REQUEST arrived in draft-16. + if new_group_request.is_some() && !has_new_group_request(version) { + return Err(DecodeError::InvalidValue); + } + let subscriber_priority = subscriber_priority.unwrap_or(128); let forward = forward.unwrap_or(true); @@ -472,18 +490,38 @@ impl Message for SubscribeUpdate { if matches!(version, Version::Draft17) { let _required_request_id_delta = u64::decode(r, version)?; } + // Nothing reads an update's Range Filters yet; they are consumed so a legal + // update does not fail the session. TRACK_PROPERTY_FILTER (0x29) is legal + // only on an update to SUBSCRIBE_TRACKS, which the message alone can't tell. decode_params!(r, version, 0x02 => _object_delivery_timeout: Option, + 0x03 => _authorization_token: Vec, 0x06 => _subgroup_delivery_timeout: Option, 0x10 => forward: Option, 0x20 => subscriber_priority: Option, 0x21 => _filter: Option, 0x23 => fill: Option, + 0x25 => subgroup_filter: Vec, + 0x26 => object_id_filter: Vec, + 0x27 => priority_filter: Vec, + 0x28 => object_property_filter: Vec, + 0x29 => track_property_filter: Vec, + 0x32 => _new_group_request: Option, ); - // FILL_PARAMETERS is a draft-20 addition, so an earlier peer sending one is - // still the protocol violation it was. - if fill.is_some() && !Filter::is_draft20(version) { + let range_filters = [ + subgroup_filter, + object_id_filter, + priority_filter, + object_property_filter, + track_property_filter, + ] + .iter() + .any(|filter| !filter.is_empty()); + + // FILL_PARAMETERS and the Range Filters postdate draft-17, so an earlier peer + // sending one is still the protocol violation it was. + if (fill.is_some() && !Filter::is_draft20(version)) || (range_filters && !has_range_filters(version)) { return Err(DecodeError::InvalidValue); } @@ -530,6 +568,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: true, + forward: true, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft14); @@ -552,6 +592,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: true, + forward: true, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft15); @@ -642,6 +684,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: true, + forward: true, + range_filters: false, }; for version in [Version::Draft17, Version::Draft18, Version::Draft19, Version::Draft20] { @@ -671,6 +715,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: true, + forward: true, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft14); @@ -854,6 +900,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: wanted, + forward: true, + range_filters: false, }; let encoded = encode_message(&msg, version); @@ -976,6 +1024,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: true, + forward: true, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft17); @@ -1034,6 +1084,8 @@ mod tests { filter: Filter::NextObject, fill: None, properties_wanted: true, + forward: true, + range_filters: false, }; let encoded = encode_message(&msg, Version::Draft18); @@ -1242,4 +1294,168 @@ mod tests { let v18 = encode_message(&v18_msg, Version::Draft18); assert_eq!(v17.len(), v18.len() + 1); } + + /// The head of a SUBSCRIBE (draft-20 Figure 10) up to its Number of Parameters: + /// Request ID 1, Track Namespace ("live"), Track Name ("video"). Every value is + /// below 128, so each varint is one byte in the drafts before and after 17 alike. + const SUBSCRIBE_HEAD: &[u8] = &[ + 0x01, 0x01, 0x04, b'l', b'i', b'v', b'e', 0x05, b'v', b'i', b'd', b'e', b'o', + ]; + + fn subscribe_body(params: &[u8]) -> Vec { + [SUBSCRIBE_HEAD, params].concat() + } + + /// Every parameter draft-20 lets a SUBSCRIBE carry that we don't act on still decodes, + /// with what we have to refuse recorded rather than failing the session. + #[test] + fn subscribe_decodes_every_draft20_parameter() { + #[rustfmt::skip] + let body = subscribe_body(&[ + 0x06, // Number of Parameters + 0x03, 0x03, 0x03, 0x00, 0xAA, // AUTHORIZATION TOKEN: USE_VALUE, type 0, one byte + 0x00, 0x03, 0x03, 0x00, 0xBB, // a second one, which the draft allows + 0x0D, 0x00, // FORWARD (0x10) = 0, a uint8 + 0x15, 0x03, 0x00, 0x00, 0x01, // SUBGROUP_FILTER (0x25): SetID 0, range 0-1 + 0x0D, 0x07, // NEW_GROUP_REQUEST (0x32), a varint + 0x03, 0x00, // INCLUDE_PROPERTIES (0x35) = 0, a uint8 + ]); + + for version in [Version::Draft20, Version::Draft21, Version::Draft22] { + let mut buf = bytes::Bytes::from(body.clone()); + let msg = Subscribe::decode_msg(&mut buf, version).unwrap_or_else(|e| panic!("{version}: {e}")); + assert!(buf.is_empty(), "{version}: trailing bytes"); + assert!(!msg.forward, "{version}"); + assert!(msg.range_filters, "{version}"); + assert!(!msg.properties_wanted, "{version}"); + } + } + + /// INCLUDE_PROPERTIES is a uint8, one raw byte with no Length. Reading a Length would + /// swallow the parameter after it, or fail on the value. + #[test] + fn include_properties_is_a_uint8() { + let msg = Subscribe { + request_id: RequestId(1), + track_namespace: crate::Path::new("live"), + track_name: "video".into(), + subscriber_priority: 128, + group_order: GroupOrder::Descending, + filter: Filter::Unfiltered, + fill: None, + properties_wanted: false, + forward: true, + range_filters: false, + }; + + #[rustfmt::skip] + let expected = subscribe_body(&[ + 0x04, // Number of Parameters + 0x10, 0x01, // FORWARD = 1 + 0x10, 0x80, // SUBSCRIBER_PRIORITY (0x20) = 128, a raw byte + 0x02, 0x02, // GROUP_ORDER (0x22) = Descending + 0x13, 0x00, // INCLUDE_PROPERTIES (0x35) = 0 + ]); + assert_eq!(encode_message(&msg, Version::Draft20), expected); + + // Anything but 0 or 1 is a protocol violation. + let body = subscribe_body(&[0x01, 0x35, 0x02]); + assert!(decode_message::(&body, Version::Draft20).is_err()); + } + + /// FORWARD=0 is legal on every draft. It decodes so the request can be refused, + /// rather than closing the session. + #[test] + fn subscribe_decodes_forward_zero() { + for version in [ + Version::Draft14, + Version::Draft15, + Version::Draft16, + Version::Draft17, + Version::Draft18, + Version::Draft19, + Version::Draft20, + ] { + let msg = Subscribe { + request_id: RequestId(1), + track_namespace: crate::Path::new("live"), + track_name: "video".into(), + subscriber_priority: 128, + group_order: GroupOrder::Descending, + filter: Filter::NextObject, + fill: None, + properties_wanted: true, + forward: false, + range_filters: false, + }; + let decoded: Subscribe = decode_message(&encode_message(&msg, version), version).unwrap(); + assert!(!decoded.forward, "{version}"); + } + } + + /// AUTHORIZATION TOKEN is a Key-Value-Pair with an odd type through draft-16, so its + /// value is length prefixed there too. + #[test] + fn subscribe_accepts_a_token_on_every_draft() { + #[rustfmt::skip] + let body = subscribe_body(&[ + 0x01, // Number of Parameters + 0x03, 0x03, 0x03, 0x00, 0xAA, // AUTHORIZATION TOKEN + ]); + for version in [ + Version::Draft15, + Version::Draft16, + Version::Draft17, + Version::Draft18, + Version::Draft19, + Version::Draft20, + ] { + let mut body = body.clone(); + if version == Version::Draft17 { + // Required Request ID Delta, draft-17 only. + body.insert(1, 0x00); + } + let mut buf = bytes::Bytes::from(body); + Subscribe::decode_msg(&mut buf, version).unwrap_or_else(|e| panic!("{version}: {e}")); + assert!(buf.is_empty(), "{version}: trailing bytes"); + } + } + + /// A parameter from a later draft is still an unknown parameter, which the draft makes + /// a protocol violation. + #[test] + fn subscribe_refuses_parameters_before_their_draft() { + for (version, params) in [ + // NEW_GROUP_REQUEST arrived in draft-16. + (Version::Draft15, &[0x01, 0x32, 0x00][..]), + // The Range Filters arrived in draft-19. + (Version::Draft18, &[0x01, 0x25, 0x00][..]), + // INCLUDE_PROPERTIES arrived in draft-20. + (Version::Draft19, &[0x01, 0x35, 0x00][..]), + ] { + let body = subscribe_body(params); + assert!( + decode_message::(&body, version).is_err(), + "{version}: {params:x?}" + ); + } + } + + /// REQUEST_UPDATE (draft-20 Figure 12) carries a token, NEW_GROUP_REQUEST and the + /// Range Filters, TRACK_PROPERTY_FILTER included, since the update may be for a + /// SUBSCRIBE_TRACKS. + #[test] + fn request_update_decodes_every_draft20_parameter() { + #[rustfmt::skip] + let body = [ + 0x02, // Request ID + 0x04, // Number of Parameters + 0x03, 0x03, 0x03, 0x00, 0xAA, // AUTHORIZATION TOKEN + 0x22, 0x00, // SUBGROUP_FILTER (0x25), empty: remove the filter + 0x04, 0x04, 0x00, 0x02, 0x00, 0x01, // TRACK_PROPERTY_FILTER (0x29): SetID 0, type 2, range 0-1 + 0x09, 0x00, // NEW_GROUP_REQUEST (0x32) + ]; + let decoded: SubscribeUpdate = decode_message(&body, Version::Draft20).unwrap(); + assert_eq!(decoded.request_id, RequestId(2)); + } } diff --git a/rs/moq-net/src/ietf/subscribe_namespace.rs b/rs/moq-net/src/ietf/subscribe_namespace.rs index bf57ee0bb7..b5993438db 100644 --- a/rs/moq-net/src/ietf/subscribe_namespace.rs +++ b/rs/moq-net/src/ietf/subscribe_namespace.rs @@ -76,7 +76,11 @@ impl Message for SubscribeNamespace<'_> { } let request_id = RequestId::decode(r, version)?; let namespace = decode_namespace(r, version)?; - decode_params!(r, version, HIDDEN_PARAM => hidden: Option); + // The token is ignored: the session's grant is what authorizes the request. + decode_params!(r, version, + 0x03 => _authorization_token: Vec, + HIDDEN_PARAM => hidden: Option, + ); Ok(Self { request_id, @@ -135,7 +139,11 @@ impl Message for SubscribeNamespaceLegacy<'_> { _ => 0x01, }; - decode_params!(r, version, HIDDEN_PARAM => hidden: Option); + // The token is ignored: the session's grant is what authorizes the request. + decode_params!(r, version, + 0x03 => _authorization_token: Vec, + HIDDEN_PARAM => hidden: Option, + ); Ok(Self { request_id, @@ -241,7 +249,7 @@ impl Namespace<'_> { // The base form has no Parameters field at all, so there is nothing to read // (and nothing to reject) unless the extension is on. let cluster = match negotiated { - true => super::publish_namespace::decode_cluster_params(r, version, true)?, + true => Some(super::publish_namespace::decode_cluster_params(r, version)?), false => None, }; diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index 4fc5a5c68b..072d8170c8 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -1953,6 +1953,8 @@ where filter: join.filter, fill: join.fill, properties_wanted: true, + forward: true, + range_filters: false, }) .await?; Ok(()) @@ -2021,6 +2023,7 @@ where ), group_order: GroupOrder::Ascending, fetch_type, + range_filters: false, }) .await?; Ok::<(), Error>(()) diff --git a/rs/moq-net/src/ietf/track.rs b/rs/moq-net/src/ietf/track.rs index 2c600c331c..00793702c4 100644 --- a/rs/moq-net/src/ietf/track.rs +++ b/rs/moq-net/src/ietf/track.rs @@ -7,11 +7,11 @@ use num_enum::{IntoPrimitive, TryFromPrimitive}; use crate::{ Path, coding::*, - ietf::{Filter, GroupOrder, Parameters, RequestId}, + ietf::{Filter, GroupOrder, RequestId, Subscribe}, }; use super::Message; -use super::namespace::{decode_namespace, encode_namespace}; +use super::namespace::encode_namespace; use super::Version; @@ -51,31 +51,14 @@ impl Message for TrackStatus<'_> { Ok(()) } + /// Every draft defines TRACK_STATUS as identical to SUBSCRIBE, so it decodes as one and + /// keeps only what names the track. We refuse the request, so the rest goes unread. fn decode_msg(r: &mut R, version: Version) -> Result { - let request_id = RequestId::decode(r, version)?; - if version == Version::Draft17 { - let _required_request_id_delta = u64::decode(r, version)?; - } - let track_namespace = decode_namespace(r, version)?; - let track_name = Cow::::decode(r, version)?; - - match version { - Version::Draft14 => { - let _subscriber_priority = u8::decode(r, version)?; - let _group_order = GroupOrder::decode(r, version)?; - let _forward = bool::decode(r, version)?; - let _filter_type = u64::decode(r, version)?; - let _params = Parameters::decode(r, version)?; - } - _ => { - decode_params!(r, version,); - } - } - + let subscribe = Subscribe::decode_msg(r, version)?; Ok(Self { - request_id, - track_namespace, - track_name, + request_id: subscribe.request_id, + track_namespace: subscribe.track_namespace, + track_name: subscribe.track_name, }) } } @@ -197,4 +180,44 @@ mod tests { assert_eq!(decoded.track_namespace.as_str(), "test/ns"); assert_eq!(decoded.track_name, "video"); } + + /// TRACK_STATUS is identical to SUBSCRIBE on every draft, so it carries whatever + /// fields and parameters a SUBSCRIBE can. A peer that sends them must still reach + /// our refusal rather than having its session closed. + #[test] + fn test_track_status_carries_subscribe_fields() { + // Request ID 1, Track Namespace ("live"), Track Name ("video"). + let head: &[u8] = &[ + 0x01, 0x01, 0x04, b'l', b'i', b'v', b'e', 0x05, b'v', b'i', b'd', b'e', b'o', + ]; + + #[rustfmt::skip] + let cases: [(Version, &[u8]); 3] = [ + (Version::Draft14, &[ + 0x80, // Subscriber Priority + 0x02, // Group Order + 0x00, // Forward + 0x03, 0x05, 0x01, // Filter Type AbsoluteStart, at {5, 1} + 0x00, // Number of Parameters + ]), + (Version::Draft15, &[ + 0x02, // Number of Parameters + 0x10, 0x00, // FORWARD = 0 + 0x20, 0x01, // SUBSCRIBER_PRIORITY = 1 + ]), + (Version::Draft20, &[ + 0x02, // Number of Parameters + 0x03, 0x03, 0x03, 0x00, 0xAA, // AUTHORIZATION TOKEN + 0x32, 0x00, // INCLUDE_PROPERTIES (0x35) = 0 + ]), + ]; + + for (version, rest) in cases { + let body = [head, rest].concat(); + let mut buf = bytes::Bytes::from(body); + let msg = TrackStatus::decode_msg(&mut buf, version).unwrap_or_else(|e| panic!("{version}: {e}")); + assert!(buf.is_empty(), "{version}: trailing bytes"); + assert_eq!(msg.track_name, "video", "{version}"); + } + } } diff --git a/rs/moq-net/src/ietf/version.rs b/rs/moq-net/src/ietf/version.rs index c78d2bbb13..0e79c1e17b 100644 --- a/rs/moq-net/src/ietf/version.rs +++ b/rs/moq-net/src/ietf/version.rs @@ -53,8 +53,8 @@ mod tests { use super::*; use crate::coding::Encode; use crate::ietf::{ - Fetch, FetchType, Fill, Filter, GroupFlags, GroupHeader, GroupOrder, Location, Message, Properties, Publish, - RequestId, Subscribe, SubscribeOk, + EndLocation, Fetch, FetchType, Fill, Filter, GroupFlags, GroupHeader, GroupOrder, Location, Message, + Properties, Publish, RequestId, Subscribe, SubscribeOk, }; use crate::{Path, Timescale}; @@ -94,6 +94,8 @@ mod tests { range_filters: false, }), properties_wanted: false, + forward: true, + range_filters: false, }; let subscribe_ok = SubscribeOk { @@ -118,12 +120,18 @@ mod tests { request_id: RequestId(3), subscriber_priority: 64, group_order: GroupOrder::Ascending, - fetch_type: FetchType::Standalone { + fetch_type: FetchType::Filtered { namespace: Path::new("broadcast"), track: "video".into(), - start: Location { group: 1, object: 0 }, - end: Location { group: 2, object: 0 }, + filter: Filter::Absolute { + start: Location { group: 1, object: 0 }, + end: Some(EndLocation { + group: 2, + object: Some(0), + }), + }, }, + range_filters: false, }; let group = GroupHeader { From 5fb33beffaf5d2987da9ec4680a00dc9d168bc62 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Wed, 30 Sep 2026 12:02:51 -0700 Subject: [PATCH 3/4] fix(ietf): refuse FETCH with FILL_TIMEOUT, reject 0x29 on JS SUBSCRIBE/FETCH A served FETCH ignored FILL_TIMEOUT, so a cache-only request could wait on upstream. JS treated TRACK_PROPERTY_FILTER on SUBSCRIBE and FETCH as a refusable Range Filter; it stays a protocol violation there, as in Rust. Co-Authored-By: Claude Opus 5.5 --- doc/concept/standard.md | 3 ++- js/net/src/ietf/fetch.ts | 3 +++ js/net/src/ietf/ietf.test.ts | 10 ++++++++ js/net/src/ietf/parameters.ts | 21 +++++++++++++---- js/net/src/ietf/subscribe.ts | 3 +++ rs/moq-net/src/fuzz.rs | 1 + rs/moq-net/src/ietf/fetch.rs | 39 +++++++++++++++++++++++++------ rs/moq-net/src/ietf/publisher.rs | 38 ++++++++++++++++++++++++++++++ rs/moq-net/src/ietf/subscriber.rs | 2 ++ rs/moq-net/src/ietf/version.rs | 1 + 10 files changed, 108 insertions(+), 13 deletions(-) diff --git a/doc/concept/standard.md b/doc/concept/standard.md index 9ffe7f51be..f96904d2ef 100644 --- a/doc/concept/standard.md +++ b/doc/concept/standard.md @@ -75,7 +75,8 @@ the session's credential is what authorizes it. A legal request that is not served is refused on its own with `NOT_SUPPORTED`, leaving the session open: a `SUBSCRIBE` with `FORWARD=0`, a `SUBSCRIBE` or -`FETCH` carrying Range Filters (no `MAX_FILTER_RANGES` is advertised), +`FETCH` carrying Range Filters (no `MAX_FILTER_RANGES` is advertised), a +`FETCH` carrying `FILL_TIMEOUT` (Timed-Out gaps are not written), `TRACK_STATUS`, and the `FETCH` forms above. `NEW_GROUP_REQUEST` is ignored, as the draft allows a publisher without dynamic groups to do. A parameter the negotiated draft does not define still closes the session with diff --git a/js/net/src/ietf/fetch.ts b/js/net/src/ietf/fetch.ts index 1263555689..e0a5159873 100644 --- a/js/net/src/ietf/fetch.ts +++ b/js/net/src/ietf/fetch.ts @@ -93,6 +93,9 @@ export class Fetch { if (params.rangeFilters && !hasRangeFilters(version)) { throw new Error("Range Filters need draft-19"); } + if (params.trackPropertyFilter) { + throw new Error("TRACK_PROPERTY_FILTER is not allowed on FETCH"); + } return new Fetch({ requestId }); } diff --git a/js/net/src/ietf/ietf.test.ts b/js/net/src/ietf/ietf.test.ts index 12f74c5220..6f223c8ba1 100644 --- a/js/net/src/ietf/ietf.test.ts +++ b/js/net/src/ietf/ietf.test.ts @@ -1900,3 +1900,13 @@ test("TrackStatusRequest: carries SUBSCRIBE fields", async () => { expect(msg.trackName).toBe("video"); } }); + +// TRACK_PROPERTY_FILTER is legal only on SUBSCRIBE_TRACKS and its updates, so on a +// SUBSCRIBE or FETCH it stays a protocol violation rather than a per-request refusal. +test("TRACK_PROPERTY_FILTER is rejected on SUBSCRIBE and FETCH", async () => { + // One parameter: TRACK_PROPERTY_FILTER (0x29), SetID 0, property type 2, from 0. + const params = [0x01, 0x29, 0x03, 0x00, 0x02, 0x00]; + const body = framed([...TRACK_HEAD, ...params]); + await expect(decodeVersioned(body, Subscribe.Subscribe.decode, Version.DRAFT_20)).rejects.toThrow(); + await expect(decodeVersioned(body, Fetch.decode, Version.DRAFT_20)).rejects.toThrow(); +}); diff --git a/js/net/src/ietf/parameters.ts b/js/net/src/ietf/parameters.ts index 6b7ccf841b..e86df0239f 100644 --- a/js/net/src/ietf/parameters.ts +++ b/js/net/src/ietf/parameters.ts @@ -223,12 +223,18 @@ const MSG_PARAM_FILL_PARAMETERS = 0x23n; /// HOP_PATH, from the MoQ Cluster extension. See `cluster.ts`. const MSG_PARAM_HOP_PATH = 0x40b57n; -/// The Range Filters (draft-19): SUBGROUP, OBJECTID, PRIORITY, OBJECT_PROPERTY and -/// TRACK_PROPERTY. Each is length prefixed whatever the parity of its id. -const MSG_PARAM_RANGE_FILTERS: readonly bigint[] = [0x25n, 0x26n, 0x27n, 0x28n, 0x29n]; +/// The object Range Filters (draft-19): SUBGROUP, OBJECTID, PRIORITY and OBJECT_PROPERTY. +/// Each is length prefixed whatever the parity of its id. +const MSG_PARAM_RANGE_FILTERS: readonly bigint[] = [0x25n, 0x26n, 0x27n, 0x28n]; +/// TRACK_PROPERTY_FILTER, the Range Filter legal only on SUBSCRIBE_TRACKS and its updates. +const MSG_PARAM_TRACK_PROPERTY_FILTER = 0x29n; /// The parameters whose definitions let them repeat within one message. -const MSG_PARAM_REPEATABLE: readonly bigint[] = [MSG_PARAM_AUTHORIZATION_TOKEN, ...MSG_PARAM_RANGE_FILTERS]; +const MSG_PARAM_REPEATABLE: readonly bigint[] = [ + MSG_PARAM_AUTHORIZATION_TOKEN, + ...MSG_PARAM_RANGE_FILTERS, + MSG_PARAM_TRACK_PROPERTY_FILTER, +]; type MessageParamKind = "varint" | "uint8" | "bool" | "location" | "bytes"; /** A `{Group, Object}` pair carried by a message parameter, such as LARGEST_OBJECT. */ @@ -260,7 +266,7 @@ function getMessageParamKind(id: bigint): MessageParamKind { case MSG_PARAM_HOP_PATH: return "bytes"; default: - if (MSG_PARAM_RANGE_FILTERS.includes(id)) return "bytes"; + if (MSG_PARAM_REPEATABLE.includes(id)) return "bytes"; throw new Error(`unknown message parameter id: ${id.toString()}`); } } @@ -317,6 +323,11 @@ export class Parameters { return MSG_PARAM_RANGE_FILTERS.some((id) => this.#repeated.has(id)); } + /** Whether the message carried TRACK_PROPERTY_FILTER, which only a SUBSCRIBE_TRACKS may. */ + get trackPropertyFilter(): boolean { + return this.#repeated.has(MSG_PARAM_TRACK_PROPERTY_FILTER); + } + // --- Numeric accessors --- get subscriberPriority(): number | undefined { diff --git a/js/net/src/ietf/subscribe.ts b/js/net/src/ietf/subscribe.ts index 84ec371890..32cf1dba1e 100644 --- a/js/net/src/ietf/subscribe.ts +++ b/js/net/src/ietf/subscribe.ts @@ -171,6 +171,9 @@ export class Subscribe { if (params.rangeFilters && !Filter.hasRangeFilters(version)) { throw new Error("Range Filters need draft-19"); } + if (params.trackPropertyFilter) { + throw new Error("TRACK_PROPERTY_FILTER is not allowed on SUBSCRIBE"); + } // An absent LOCATION_FILTER means the subscription is unfiltered. const raw = params.subscriptionFilter; diff --git a/rs/moq-net/src/fuzz.rs b/rs/moq-net/src/fuzz.rs index c08c40ad0a..fe5c5917a6 100644 --- a/rs/moq-net/src/fuzz.rs +++ b/rs/moq-net/src/fuzz.rs @@ -714,6 +714,7 @@ pub fn seeds() -> Vec { group_order: ietf::GroupOrder::Ascending, fetch_type, range_filters: false, + fill_timeout: false, }; fetch.encode_bytes(*version).ok() }) else { diff --git a/rs/moq-net/src/ietf/fetch.rs b/rs/moq-net/src/ietf/fetch.rs index 71e96bc194..e0e1b09730 100644 --- a/rs/moq-net/src/ietf/fetch.rs +++ b/rs/moq-net/src/ietf/fetch.rs @@ -129,6 +129,10 @@ pub struct Fetch<'a> { /// MAX_FILTER_RANGES, so the request is refused rather than served unfiltered. /// Never encoded; we send no range filters. pub range_filters: bool, + /// Whether the request carried FILL_TIMEOUT (0x0A), a budget for waiting on upstream + /// that ends in Timed-Out gaps we cannot write, so the request is refused rather than + /// left to wait without it. Never encoded; we send no fill timeout. + pub fill_timeout: bool, } impl Message for Fetch<'_> { @@ -186,15 +190,15 @@ impl Message for Fetch<'_> { let _required_request_id_delta = u64::decode(buf, version)?; } - // The token is ignored: the session's grant is what authorizes the request. We - // refuse or serve a FETCH the same whatever its FILL_TIMEOUT or INCLUDE_PROPERTIES. - let (fetch_type, subscriber_priority, group_order, range_filters) = match version { + // The token is ignored: the session's grant is what authorizes the request, and + // INCLUDE_PROPERTIES only shapes a FETCH_OK we don't send on draft-20. + let (fetch_type, subscriber_priority, group_order, range_filters, fill_timeout) = match version { Version::Draft14 => { let subscriber_priority = u8::decode(buf, version)?; let group_order = GroupOrder::decode(buf, version)?; let fetch_type = FetchType::decode(buf, version)?; Parameters::skip(buf, version)?; - (fetch_type, Some(subscriber_priority), Some(group_order), false) + (fetch_type, Some(subscriber_priority), Some(group_order), false, false) } Version::Draft15 | Version::Draft16 | Version::Draft17 | Version::Draft18 | Version::Draft19 => { let fetch_type = FetchType::decode(buf, version)?; @@ -224,7 +228,13 @@ impl Message for Fetch<'_> { return Err(DecodeError::InvalidValue); } - (fetch_type, subscriber_priority, group_order, range_filters) + ( + fetch_type, + subscriber_priority, + group_order, + range_filters, + fill_timeout.is_some(), + ) } // Draft-20 names the track up front and moves the range into LOCATION_FILTER. _ => { @@ -232,7 +242,7 @@ impl Message for Fetch<'_> { let track = Cow::::decode(buf, version)?; decode_params!(buf, version, 0x03 => _authorization_token: Vec, - 0x0A => _fill_timeout: Option, + 0x0A => fill_timeout: Option, 0x20 => subscriber_priority: Option, 0x21 => filter: Option, 0x22 => group_order: Option, @@ -257,7 +267,13 @@ impl Message for Fetch<'_> { // An absent LOCATION_FILTER fetches the whole track. filter: filter.unwrap_or(Filter::Unfiltered), }; - (fetch_type, subscriber_priority, group_order, range_filters) + ( + fetch_type, + subscriber_priority, + group_order, + range_filters, + fill_timeout.is_some(), + ) } }; @@ -268,6 +284,7 @@ impl Message for Fetch<'_> { group_order: group_order.unwrap_or(GroupOrder::Any), fetch_type, range_filters, + fill_timeout, }) } } @@ -649,6 +666,7 @@ mod tests { end: Location { group: 10, object: 5 }, }, range_filters: false, + fill_timeout: false, }; let encoded = encode_message(&msg, Version::Draft14); @@ -671,6 +689,7 @@ mod tests { end: Location { group: 10, object: 5 }, }, range_filters: false, + fill_timeout: false, }; let encoded = encode_message(&msg, Version::Draft15); @@ -710,6 +729,7 @@ mod tests { end: Location { group: 10, object: 5 }, }, range_filters: false, + fill_timeout: false, }; let encoded = encode_message(&msg, Version::Draft16); @@ -732,6 +752,7 @@ mod tests { end: Location { group: 10, object: 5 }, }, range_filters: false, + fill_timeout: false, }; let encoded = encode_message(&msg, Version::Draft17); @@ -805,6 +826,7 @@ mod tests { end: Location { group: 10, object: 5 }, }, range_filters: false, + fill_timeout: false, }; let encoded = encode_message(&msg, Version::Draft18); @@ -881,6 +903,7 @@ mod tests { assert_eq!(fetch.subscriber_priority, 64); assert_eq!(fetch.group_order, GroupOrder::Ascending); assert!(fetch.range_filters, "{version}"); + assert!(fetch.fill_timeout, "{version}"); assert_eq!( fetch.fetch_type, FetchType::Filtered { @@ -910,6 +933,7 @@ mod tests { filter: Filter::Unfiltered, }, range_filters: false, + fill_timeout: false, }; #[rustfmt::skip] @@ -951,6 +975,7 @@ mod tests { assert_eq!(decoded.is_ok(), ok, "{version}: {body:x?}"); if let Ok(fetch) = decoded { assert_eq!(fetch.range_filters, body == &range_filter, "{version}"); + assert_eq!(fetch.fill_timeout, body == &fill_timeout, "{version}"); } } } diff --git a/rs/moq-net/src/ietf/publisher.rs b/rs/moq-net/src/ietf/publisher.rs index d02156729c..77c75b3d26 100644 --- a/rs/moq-net/src/ietf/publisher.rs +++ b/rs/moq-net/src/ietf/publisher.rs @@ -1164,6 +1164,19 @@ where .await; } + // FILL_TIMEOUT=0 asks for cache only, and any budget ends in Timed-Out gaps we + // don't write, so waiting on upstream regardless would ignore what was asked. + if msg.fill_timeout { + return self + .reject_fetch( + stream, + msg.request_id, + &Error::Unsupported, + "FILL_TIMEOUT not supported", + ) + .await; + } + let (track, start, end, timescale, joined) = match msg.fetch_type { FetchType::Standalone { namespace, @@ -3309,6 +3322,26 @@ mod serve_tests { ]), ]; + // A standalone FETCH we would otherwise serve, carrying FILL_TIMEOUT=0: cache only, + // with gaps reported as Timed-Out, which we can't write. + #[rustfmt::skip] + let fill_timeout = [ + 0x2B, // Request ID + 0x01, // Standalone + 0x01, 0x04, b'r', b'o', b'o', b'm', 0x05, b'v', b'i', b'd', b'e', b'o', + 0x00, 0x00, // Start Location + 0x00, 0x00, // End Location: the whole group 0 + 0x01, // Number of Parameters + 0x0A, 0x00, // FILL_TIMEOUT = 0 + ]; + for version in [Version::Draft18, Version::Draft19] { + assert_eq!( + refusal(version, ietf::Fetch::ID, fill_timeout.to_vec()).await, + NOT_SUPPORTED, + "{version}: FETCH with FILL_TIMEOUT" + ); + } + for version in [Version::Draft20, Version::Draft21, Version::Draft22] { for (label, id, params) in cases { assert_eq!( @@ -3495,6 +3528,7 @@ mod serve_tests { group_offset: 0, }, range_filters: false, + fill_timeout: false, }, ) .await?; @@ -3901,6 +3935,7 @@ mod serve_tests { end, }, range_filters: false, + fill_timeout: false, }, ) .await @@ -4142,6 +4177,7 @@ mod serve_tests { group_order: GroupOrder::Ascending, fetch_type, range_filters: false, + fill_timeout: false, }, ) .await @@ -4184,6 +4220,7 @@ mod serve_tests { group_id: 5, }, range_filters: false, + fill_timeout: false, }, ) .await @@ -5730,6 +5767,7 @@ mod tests { group_order: GroupOrder::Descending, fetch_type, range_filters: false, + fill_timeout: false, }, ) .await diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index 4b863f7279..246e7f43ec 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -2093,6 +2093,7 @@ where group_order: GroupOrder::Ascending, fetch_type, range_filters: false, + fill_timeout: false, }) .await?; Ok::<(), Error>(()) @@ -3034,6 +3035,7 @@ where }, }, range_filters: false, + fill_timeout: false, }) .await?; self.read_group_fetch_response(&mut stream).await diff --git a/rs/moq-net/src/ietf/version.rs b/rs/moq-net/src/ietf/version.rs index 0e79c1e17b..9e1923a055 100644 --- a/rs/moq-net/src/ietf/version.rs +++ b/rs/moq-net/src/ietf/version.rs @@ -132,6 +132,7 @@ mod tests { }, }, range_filters: false, + fill_timeout: false, }; let group = GroupHeader { From 3f632387646816251230541c562e1bd72847f61a Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Wed, 30 Sep 2026 12:37:04 -0700 Subject: [PATCH 4/4] fix(net): keep draft-16 INCLUDE_PROPERTIES a protocol violation The uint8 accessor only read the draft-17+ storage, so a draft-16 KVP-framed 0x35 read as absent and slipped past the draft-20 gate. Co-Authored-By: Claude Opus 5.5 --- js/net/src/ietf/ietf.test.ts | 8 ++++++++ js/net/src/ietf/parameters.ts | 5 ++++- 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/js/net/src/ietf/ietf.test.ts b/js/net/src/ietf/ietf.test.ts index 6f223c8ba1..93ef4815a7 100644 --- a/js/net/src/ietf/ietf.test.ts +++ b/js/net/src/ietf/ietf.test.ts @@ -1910,3 +1910,11 @@ test("TRACK_PROPERTY_FILTER is rejected on SUBSCRIBE and FETCH", async () => { await expect(decodeVersioned(body, Subscribe.Subscribe.decode, Version.DRAFT_20)).rejects.toThrow(); await expect(decodeVersioned(body, Fetch.decode, Version.DRAFT_20)).rejects.toThrow(); }); + +// INCLUDE_PROPERTIES postdates draft-16, where it arrives as a Key-Value-Pair; it must +// still read as present so the draft gate refuses it. +test("Subscribe v16: rejects INCLUDE_PROPERTIES", async () => { + // One parameter: INCLUDE_PROPERTIES (0x35), length 1, value 0. + const body = framed([...TRACK_HEAD, 0x01, 0x35, 0x01, 0x00]); + await expect(decodeVersioned(body, Subscribe.Subscribe.decode, Version.DRAFT_16)).rejects.toThrow(); +}); diff --git a/js/net/src/ietf/parameters.ts b/js/net/src/ietf/parameters.ts index e86df0239f..0530e27917 100644 --- a/js/net/src/ietf/parameters.ts +++ b/js/net/src/ietf/parameters.ts @@ -439,7 +439,10 @@ export class Parameters { /** INCLUDE_PROPERTIES: whether the peer wants Track Properties on the response. */ get includeProperties(): boolean | undefined { - const v = this.vars.get(MSG_PARAM_INCLUDE_PROPERTIES); + // Draft-16 and earlier frame every parameter as a Key-Value-Pair, so an odd id lands + // in `bytes`. It still has to read as present, so the draft-20 gate can refuse it. + const legacy = this.bytes.get(MSG_PARAM_INCLUDE_PROPERTIES); + const v = legacy ? (legacy.length === 1 ? BigInt(legacy[0]) : 2n) : this.vars.get(MSG_PARAM_INCLUDE_PROPERTIES); if (v === undefined) return undefined; // The draft allows exactly 0 or 1; anything else is a protocol violation. if (v > 1n) throw new Error(`invalid INCLUDE_PROPERTIES value: ${v}`);