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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 12 additions & 2 deletions doc/concept/standard.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ On drafts 14–19, the Rust publisher also serves relative and absolute joining
saved prefix only, while the subscription delivers later objects. One reaching
back to earlier groups is refused with `NOT_SUPPORTED`. Draft-20 uses
subscription fills instead. JavaScript
publishing does not yet serve `FETCH`;
publishing refuses every `FETCH` with `NOT_SUPPORTED`;
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
Expand All @@ -70,7 +70,17 @@ 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), 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
`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,
Expand Down
6 changes: 6 additions & 0 deletions js/net/src/ietf/connection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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: {
Expand Down
100 changes: 55 additions & 45 deletions js/net/src/ietf/fetch.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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<void> {
Expand All @@ -82,12 +52,52 @@ export class Fetch {
return Message.encode(w, this.#encode.bind(this));
}

static async decode(r: Reader, _version: IetfVersion): Promise<Fetch> {
return Message.decode(r, Fetch.#decode);
}

static async #decode(_r: Reader): Promise<Fetch> {
throw new Error("FETCH messages are not supported");
static async decode(r: Reader, version: IetfVersion): Promise<Fetch> {
return Message.decode(r, (mr) => Fetch.#decode(mr, version));
}

static async #decode(r: Reader, version: IetfVersion): Promise<Fetch> {
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);

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Enforce the version-specific FETCH parameter allowlist

When a peer sends a parameter that is recognized globally but is illegal on FETCH, Parameters.decode accepts it and these checks let the request reach the normal NOT_SUPPORTED refusal instead of treating it as a protocol violation. For example, draft-20 FETCH accepts FORWARD or NEW_GROUP_REQUEST, and draft-16/17 FETCH accepts the draft-18 FILL_TIMEOUT. The new explicit 0x29 check is fresh evidence that the same per-message allowlist remains incomplete; validate every decoded parameter against the negotiated draft's FETCH set.

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

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Declining for this PR. js/net has always validated message parameters against one global id table rather than per message and per draft, and SUBSCRIBE and PUBLISH have the same gap today, so fixing FETCH alone would leave it inconsistent. The cases here are low risk: JS refuses every FETCH with NOT_SUPPORTED, so a misplaced parameter never changes what gets served. The checks this PR adds only cover what Rust also gates (0x29, Range Filters before draft-19). Moving js/net to per-message allowlists that mirror decode_params! is worth doing as its own change. I have listed it as a follow-up.

(Written by Claude Opus 5.5)

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 });
}
}

Expand Down
5 changes: 5 additions & 0 deletions js/net/src/ietf/filter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
121 changes: 121 additions & 0 deletions js/net/src/ietf/ietf.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -1797,3 +1798,123 @@ 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");
}
});

// 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();
});

// 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();
});
Loading
Loading