Skip to content
5 changes: 5 additions & 0 deletions doc/concept/standard.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,11 @@ This is an implementation limit, not a limit in the IETF draft. A larger
declared block stops its subgroup stream with `MALFORMED_TRACK` before reading
the block; other groups and the session stay open.

Rust and JavaScript read an incoming padding stream (draft-18 and later) to
the end and discard it, without sending `STOP_SENDING`. A unidirectional stream type the negotiated draft does not
define, or a `SUBGROUP_HEADER` type it marks invalid, closes the session with
`PROTOCOL_VIOLATION`, as the draft requires.

An IETF publisher declares the track's default priority in `SUBSCRIBE_OK` or
`PUBLISH` when that draft carries track properties. Groups without a priority
flag inherit it. If the property is absent, the IETF wire default of 128 maps
Expand Down
78 changes: 78 additions & 0 deletions js/net/src/ietf/connection.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
import { expect, spyOn, test } from "bun:test";
import { SessionCode } from "../error.ts";
import { createMockTransportPair } from "../mock.ts";
import { Stream, Writer } from "../stream.ts";
import { Connection } from "./connection.ts";
import { ALPN, type IetfVersion, Version } from "./version.ts";

const PADDING = 0x132b3e28n;

/** A server-side session over `version`, and a client uni stream that sends only `type`. */
async function openUni(
version: IetfVersion,
alpn: string,
type: bigint,
): Promise<{ pair: ReturnType<typeof createMockTransportPair>; connection: Connection; writer: Writer }> {
const pair = createMockTransportPair(alpn);
const control = await Stream.open(pair.server, { version });
const connection = new Connection({
url: new URL("https://example.com"),
quic: pair.server,
control,
maxRequestId: 100n,
version,
client: false,
});

const writer = new Writer(await pair.client.createUnidirectionalStream(), version);
await writer.u62(type);

return { pair, connection, writer };
}

/** Padding is read to the end and dropped: no STOP_SENDING, and the session stays up. */
test("a padding stream is discarded", async () => {
const { pair, connection, writer } = await openUni(Version.DRAFT_19, ALPN.DRAFT_19, PADDING);
let closed = false;
void pair.server.closed.then(() => {
closed = true;
});

try {
await writer.write(new Uint8Array(4096));
writer.close();

// A STOP_SENDING would reject this; a clean close means every byte was read.
await writer.closed;

// A session close would come from the same handler that read the stream.
await Bun.sleep(0);
expect(closed).toBe(false);
} finally {
connection.close();
}
});

/** An unknown or invalid stream type MUST close the session with PROTOCOL_VIOLATION. */
test("an unknown uni stream type closes the session", async () => {
for (const [version, alpn, type] of [
[Version.DRAFT_19, ALPN.DRAFT_19, 0n],
// Past 2^53, so it cannot be read as a number.
[Version.DRAFT_19, ALPN.DRAFT_19, 2n ** 53n],
// Padding arrived in draft-18.
[Version.DRAFT_17, ALPN.DRAFT_17, PADDING],
// A SUBGROUP_HEADER with the reserved SUBGROUP_ID_MODE (0b11).
[Version.DRAFT_19, ALPN.DRAFT_19, 0x56n],
] as const) {
const logged = spyOn(console, "error").mockImplementation(() => void 0);
const { pair, connection } = await openUni(version, alpn, type);

try {
const info = await pair.client.closed;
expect(info.closeCode).toBe(SessionCode.ProtocolViolation);
} finally {
logged.mockRestore();
connection.close();
}
}
});
58 changes: 47 additions & 11 deletions js/net/src/ietf/connection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,14 @@ import type * as announce from "../announced.ts";
import type { Established } from "../connection/established.ts";
import { type Probe, type Stats, transportStats } from "../connection/stats.ts";
import { type Transport, transportOf } from "../connection/transport.ts";
import { error, fromClose, ProtocolViolation, StreamCode, StreamError } from "../error.ts";
import { error, fromClose, ProtocolViolation, SessionCode, StreamCode, StreamError } from "../error.ts";
import type { Consumer as OriginConsumer } from "../origin.ts";
import type * as Path from "../path.ts";
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 { Fetch, FetchHeader } from "./fetch.ts";
import { GoAway } from "./goaway.ts";
import { Group } from "./object.ts";
import { Publish } from "./publish.ts";
Expand All @@ -22,6 +22,9 @@ import { Subscriber } from "./subscriber.ts";
import { TrackStatusRequest } from "./track.ts";
import { type IetfVersion, Version, versionName } from "./version.ts";

// The PADDING stream type (draft-18+): bytes a peer sends to probe for bandwidth.
const PADDING = 0x132b3e28n;

/**
* Represents a connection to a MoQ server using moq-transport protocol.
*
Expand Down Expand Up @@ -158,17 +161,29 @@ export class Connection implements Established {
* Closes the connection.
*/
close() {
this.#close();
}

// Close with the session code the peer should see, a clean close by default.
#close(info?: WebTransportCloseInfo) {
if (this.#closed) return;

this.#closed = true;

this.#session.close();

// Before the session, whose own close would send a clean code first.
try {
this.#quic.close();
this.#quic.close(info);
} catch {
// ignore
}

this.#session.close();
}

// The peer broke the protocol, so losing the stream is not enough: nothing stops it
// repeating the violation on the next one.
#violated(err: ProtocolViolation) {
this.#close({ closeCode: SessionCode.ProtocolViolation, reason: err.message });
}

async #run(): Promise<void> {
Expand Down Expand Up @@ -199,10 +214,7 @@ export class Connection implements Established {
void this.#runBidi(stream).catch((err: unknown) => {
console.error("error processing bidi stream", err);
stream.abort(new Error("bidi stream error"));

// The peer broke the protocol, so losing the stream is not enough: nothing
// stops it repeating the violation on the next one.
if (err instanceof ProtocolViolation) this.close();
if (err instanceof ProtocolViolation) this.#violated(err);
});
}
}
Expand Down Expand Up @@ -313,13 +325,37 @@ export class Connection implements Established {
.catch((err: unknown) => {
console.error("error processing object stream", err);
stream.stop(err);

// An unknown or invalid stream type MUST close the session, not just the stream.
if (err instanceof ProtocolViolation) this.#violated(err);
});
}
}

async #runUni(stream: Reader) {
const header = await Group.decode(stream, this.#session.version);
await this.#subscriber.handleGroup(header, stream);
const version = this.#session.version;
// Full width, so an unknown type past 2^53 is still classified rather than thrown.
const type = await stream.u62();

// SUBGROUP_HEADER types match 0b0XX1XXXX; Group.decode validates the bits per draft.
if (type <= 0xffn && (type & 0x90n) === 0x10n) {
const header = await Group.decode(stream, version, Number(type));
await this.#subscriber.handleGroup(header, stream);
return;
}

// The receiver MUST discard padding. We read it to the end rather than cancel,
// so a peer probing for bandwidth gets the throughput it is measuring.
if (type === PADDING && version >= Version.DRAFT_18) {
await stream.discard();
return;
}

// We never FETCH, so a fetch response answers nothing of ours.
if (type === BigInt(FetchHeader.type)) throw new Error("unexpected fetch stream");

// Anything else is unknown, and a second SETUP is a violation too.
throw new ProtocolViolation(`unknown uni stream type: 0x${type.toString(16)}`);
}

/**
Expand Down
26 changes: 13 additions & 13 deletions js/net/src/ietf/object.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { StreamCode, Stream as StreamError } from "../error.ts";
import { ProtocolViolation, StreamCode, Stream as StreamError } from "../error.ts";
import { asIetf, type Cursor, type Reader, Writer } from "../stream.ts";
import { Timescale, Timestamp } from "../time.ts";
import { type IetfVersion, Version } from "./version.ts";
Expand Down Expand Up @@ -213,27 +213,27 @@ export class Group {
}
}

static async decode(r: Reader, version: IetfVersion): Promise<Group> {
const raw = await r.u53();
/** Decode a SUBGROUP_HEADER. Pass `type` when the caller already read it to classify the stream. */
static async decode(r: Reader, version: IetfVersion, type?: number): Promise<Group> {
const raw = type ?? (await r.u53());
// Strip the draft-18 FIRST_OBJECT bit before the range check, but keep the value:
// it is the only signal that a subgroup starts partway through. Drafts that predate
// it carry no such signal, so they are taken at their word.
const legacy = !hasFirstObjectBit(version);
const firstObject = legacy || (raw & FIRST_OBJECT_BIT) !== 0;
const id = legacy ? raw : raw & ~FIRST_OBJECT_BIT;

let hasPriority: boolean;
let baseId: number;
if (id >= 0x10 && id <= 0x1f) {
hasPriority = true;
baseId = id;
} else if (id >= 0x30 && id <= 0x3f) {
hasPriority = false;
baseId = id - (0x30 - 0x10);
} else {
throw new Error(`Unsupported group type: ${id}`);
// An invalid type MUST close the session, including the reserved SUBGROUP_ID_MODE
// 0b11 (draft-21 section 11.3.1).
const known = (id >= 0x10 && id <= 0x1f) || (id >= 0x30 && id <= 0x3f);
if (!known || (id & 0x06) === 0x06) {
throw new ProtocolViolation(`Unsupported group type: ${raw}`);
}

// 0x30-0x3F omit the priority and inherit it from the control message.
const hasPriority = id < 0x30;
const baseId = hasPriority ? id : id - (0x30 - 0x10);

const flags: GroupFlags = {
hasExtensions: (baseId & 0x01) !== 0,
hasSubgroupObject: (baseId & 0x02) !== 0,
Expand Down
9 changes: 9 additions & 0 deletions js/net/src/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -407,6 +407,15 @@ export class Reader {
return this.#slice(this.#buffer.byteLength);
}

// Reads to the end of the stream, dropping every byte instead of buffering it.
async discard(): Promise<void> {
this.#buffer = new Uint8Array();
do {
this.#chunks = [];
this.#chunked = 0;
} while (await this.#fill());
}

async string(): Promise<string> {
return this.decode(STRING);
}
Expand Down
1 change: 0 additions & 1 deletion quest/m0/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,6 @@ Published API or wire breaks still land on dev; each quest's Plan says so.
## Required

- [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
- [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
- [Request caps](/quest/m0/request-caps.md) - lite message sizes, IETF request IDs, and per-session announces and subscriptions are bounded
- [quest check everywhere](/quest/m0/quest-check-everywhere.md) - `quest check` guards `main`, `dev`, and the line branches on push and PR, not only PRs into `main`
- [noq reassembly cap](/quest/m0/noq-reassembly-cap.md) - noq carries quinn's stream reassembly cap and the connection receive window is finite by default
Expand Down
45 changes: 0 additions & 45 deletions quest/m0/ietf-uni-stream-types.md

This file was deleted.

13 changes: 9 additions & 4 deletions rs/moq-net/src/ietf/group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -702,9 +702,10 @@ mod tests {
assert!(flags.has_end);
}

/// Regression: a publisher-emitted Draft18 GroupHeader byte must satisfy the
/// subscriber's uni-stream classifier mask `(byte & 0x90) == 0x10`. Otherwise
/// the uni stream is dropped as UnexpectedStream and the data plane stalls.
/// Regression: a publisher-emitted Draft18 GroupHeader byte must match the
/// SUBGROUP_HEADER form `(byte & 0x90) == 0x10` and decode as flags, which is how
/// the uni-stream classifier recognizes it. Otherwise the uni stream is an unknown
/// type, which closes the session.
#[test]
fn test_draft18_group_header_passes_stream_classifier() {
let header = GroupHeader {
Expand All @@ -719,10 +720,14 @@ mod tests {
header.encode(&mut buf, Version::Draft18).unwrap();
let type_byte = buf[0] as u64;

// The check in session.rs::run_uni_group.
assert_eq!(
type_byte & 0x90,
0x10,
"draft-18 SUBGROUP_HEADER type 0x{type_byte:02x} is outside the SUBGROUP_HEADER form",
);
// The check in session.rs::UniType::classify.
assert!(
GroupFlags::decode(type_byte, Version::Draft18).is_ok(),
"draft-18 SUBGROUP_HEADER type 0x{type_byte:02x} not recognized by uni-stream classifier",
);
}
Expand Down
Loading
Loading