diff --git a/.github/workflows/nightly.yml b/.github/workflows/nightly.yml index 36731b7b5d..4d8ee86658 100644 --- a/.github/workflows/nightly.yml +++ b/.github/workflows/nightly.yml @@ -103,6 +103,12 @@ jobs: if: ${{ !cancelled() }} run: nix develop --command bun js/net/bench/varint.ts + # Bytes and time for lite-06 against lite-07 frames, groups, and requests. No + # threshold: it only has to keep running. + - name: JS lite varint benchmark + if: ${{ !cancelled() }} + run: nix develop --command bun js/net/bench/lite-varint.ts + # Fails if publishing a group costs more as the track retains more groups, # which means the latency guard or the cache went back to scanning them all. - name: JS track retention benchmark diff --git a/doc/concept/moq-lite.md b/doc/concept/moq-lite.md index 393ee60648..e627c18f21 100644 --- a/doc/concept/moq-lite.md +++ b/doc/concept/moq-lite.md @@ -32,7 +32,10 @@ and newer, each side also sends a `SETUP` message with its capabilities. Rust and TypeScript speak moq-lite 01 through 06 and moq-transport drafts 14 through 22. Clients offer `moq-lite-06` first by default. moq-lite 07 is still in progress: it negotiates as `moq-lite-07-wip`, and only when both -sides explicitly enable it. +sides explicitly enable it. moq-lite 07 also switches every varint from QUIC's +two-bit length prefix to moq-transport's leading-ones form, so values up to 127 +take one byte instead of up to 63, and the range widens from 62 to 64 bits. Rust +still refuses lite-07 values above 2^62-1 until its `VarInt` widens. ## Subscription completion diff --git a/drafts/draft-lcurley-moq-lite.md b/drafts/draft-lcurley-moq-lite.md index b153781239..38621fa966 100644 --- a/drafts/draft-lcurley-moq-lite.md +++ b/drafts/draft-lcurley-moq-lite.md @@ -29,6 +29,7 @@ normative: date: false RFC3986: RFC6455: + RFC9000: RFC9002: informative: @@ -682,6 +683,12 @@ A group whose frame does not fit is simply not eligible for datagram delivery. # Encoding This section covers the encoding of each message. +## Variable-Length Integers {#varint} +A field marked `(i)` is a variable-length integer. +moq-lite-07 uses the leading-ones encoding of [moqt]: the number of leading 1 bits in the first byte gives the length, from 1 byte carrying 7 bits to 9 bytes carrying 64, and every length is valid. +Earlier versions use the two-bit length prefix of [RFC9000], Section 16, which carries at most 2^62-1. +A relay that cannot encode a value for a peer on an earlier version MUST NOT truncate it. + ## Message Length Most messages are prefixed with a variable-length integer indicating the number of bytes in the message payload that follows. This length field does not include the length of the varint length itself. @@ -1340,6 +1347,7 @@ The `Message Length` describes the payload size on the wire. ## moq-lite-07 - Assigned `moq-lite-07-wip` as this draft's protocol identifier until it is finalized as `moq-lite-07`. +- Switched every variable-length integer, including SETUP parameter values, from QUIC's two-bit length prefix to moq-transport's leading-ones encoding, widening the range to 64 bits. - Hid routes with a `.`-prefixed segment below the requested prefix from announce discovery, and added the ANNOUNCE_REQUEST `Hidden` field to opt in. - Added `Stream Count` to SUBSCRIBE_END: the number of Group Streams opened for the subscription. SUBSCRIBE_END is now sent once every counted Group Stream has opened, rather than as soon as the final group is known. - Removed SUBSCRIBE_DROP and its type 0x2; a group without a Group Stream is not counted. diff --git a/js/net/bench/frames.ts b/js/net/bench/frames.ts index 34752d5262..bff942ce72 100644 --- a/js/net/bench/frames.ts +++ b/js/net/bench/frames.ts @@ -7,6 +7,7 @@ import { Subscribe, SubscribeOk } from "../src/ietf/subscribe.ts"; import { Subscriber as IetfSubscriber } from "../src/ietf/subscriber.ts"; import { ALPN, Version as IetfVersion } from "../src/ietf/version.ts"; import { readFrames } from "../src/lite/group.ts"; +import { Version as LiteVersion } from "../src/lite/version.ts"; import { createMockTransportPair } from "../src/mock.ts"; import * as Path from "../src/path.ts"; import { Reader, Stream } from "../src/stream.ts"; @@ -48,7 +49,9 @@ const lite: Protocol = { async subscribe() { const open: Open = (stream) => { const producer = new Producer(0); - const done = readFrames(new Reader(stream), producer, 1_000_000).then(() => producer.close()); + const done = readFrames(new Reader(stream, undefined, LiteVersion.DRAFT_06), producer, 1_000_000).then(() => + producer.close(), + ); return { group: Promise.resolve(producer.consume()), done }; }; return { open, close: () => {} }; diff --git a/js/net/bench/lite-varint.ts b/js/net/bench/lite-varint.ts new file mode 100644 index 0000000000..4338a448e5 --- /dev/null +++ b/js/net/bench/lite-varint.ts @@ -0,0 +1,146 @@ +/** Lite-07 leading-ones varints against lite-06 QUIC varints, for what a publisher writes per frame, group, and request. */ +import { randomHop } from "../src/hop.ts"; +import { Datagram } from "../src/lite/datagram.ts"; +import { Group } from "../src/lite/group.ts"; +import { ProbeLevel, Setup } from "../src/lite/setup.ts"; +import { Subscribe } from "../src/lite/subscribe.ts"; +import { Version } from "../src/lite/version.ts"; +import * as Path from "../src/path.ts"; +import { type Cursor, Reader, Writer } from "../src/stream.ts"; + +const versions = [Version.DRAFT_06, Version.DRAFT_07]; +const minMs = 200; // Run each case at least this long. +let checksum = 0; + +/** A lite object: how the publisher writes it and how the subscriber reads it back. */ +interface Sample { + name: string; + encode(w: Writer, version: Version): Promise; + decode(r: Reader, version: Version): Promise; +} + +// A FRAME header: the zigzag timestamp delta, then the size. +const header = (c: Cursor) => { + c.u62(); + return c.u53(); +}; + +// FRAME headers with the payloads left out, as (timestamp delta in µs, size) pairs. +function frames(name: string, headers: [bigint, number][]): Sample { + const zigzag = (d: bigint) => (d << 1n) ^ (d >> 63n); + return { + name, + async encode(w) { + for (const [delta, size] of headers) { + await w.u62(zigzag(delta)); + await w.u53(size); + } + }, + async decode(r) { + // The subscriber's read loop, one synchronous decode per buffered frame, minus the payload. + let n = 0; + while (r.tryDecode(header) !== undefined) n++; + return n; + }, + }; +} + +const video = frames( + "Video", + Array.from({ length: 60 }, (_, n): [bigint, number] => { + if (n === 0) return [0n, 60_000]; + return [33_333n, n % 10 === 0 ? 17_000 : 8_000]; + }), +); +const audio = frames( + "Audio", + Array.from({ length: 50 }, (_, n): [bigint, number] => [n === 0 ? 0n : 20_000n, 160]), +); + +const group: Sample = { + name: "Group", + encode: (w, version) => new Group({ subscribe: 3n, sequence: 1_234 }).encode(w, version), + decode: (r, version) => Group.decode(r, version), +}; + +const subscribe: Sample = { + name: "Subscribe", + encode: (w, version) => + new Subscribe({ + id: 3n, + broadcast: Path.from("room/alice"), + track: "video", + priority: 2, + maxAge: 10_000, + }).encode(w, version), + decode: (r, version) => Subscribe.decode(r, version), +}; + +const datagram: Sample = { + name: "Datagram", + encode: (w, version) => w.write(new Datagram(3n, 1_234, 1_234_567_890, new Uint8Array()).encode(version)), + decode: (r, version) => r.readAll().then((data) => Datagram.decode(data, version)), +}; + +const hop = randomHop(); +const setup: Sample = { + name: "Setup", + encode: (w, version) => new Setup({ probe: ProbeLevel.Report, hop }).encode(w, version), + decode: (r, version) => Setup.decode(r, version), +}; + +/** Collect what one encode writes. */ +async function wire(sample: Sample, version: Version): Promise { + const chunks: Uint8Array[] = []; + const w = new Writer(new WritableStream({ write: (c) => void chunks.push(c.slice()) }), version); + await sample.encode(w, version); + w.close(); + await w.closed; + const out = new Uint8Array(chunks.reduce((n, c) => n + c.byteLength, 0)); + let offset = 0; + for (const c of chunks) { + out.set(c, offset); + offset += c.byteLength; + } + return out; +} + +/** Time `f` until `minMs` has passed, returning ns per call. */ +async function time(f: () => Promise): Promise { + for (let i = 0; i < 1_000; i++) await f(); // warm up + let calls = 0; + const start = performance.now(); + let elapsed = 0; + while (elapsed < minMs) { + for (let i = 0; i < 1_000; i++) await f(); + calls += 1_000; + elapsed = performance.now() - start; + } + return (elapsed * 1e6) / calls; +} + +// A Writer whose sink discards, so encode is timed without collecting bytes. +const sink = () => + new WritableStream({ + write: (c) => { + checksum += c.byteLength; + }, + }); + +console.log("sample,version,bytes,encode_ns,decode_ns"); +for (const sample of [video, audio, group, subscribe, datagram, setup]) { + for (const version of versions) { + const bytes = await wire(sample, version); + const encode = await time(async () => { + const w = new Writer(sink(), version); + await sample.encode(w, version); + }); + const decode = await time(async () => { + checksum += Number((await sample.decode(new Reader(undefined, bytes, version), version)) !== undefined); + }); + console.log( + `${sample.name},${version.toString(16)},${bytes.byteLength},${encode.toFixed(0)},${decode.toFixed(0)}`, + ); + } +} +if (checksum === 0) throw new Error("benchmark did no work"); diff --git a/js/net/bench/reader.ts b/js/net/bench/reader.ts index 76eac2e4c9..acdede13c6 100644 --- a/js/net/bench/reader.ts +++ b/js/net/bench/reader.ts @@ -1,4 +1,5 @@ /** Sweep frame size and chunk size for one Reader.read of a fragmented frame. */ +import { Version } from "../src/lite/version.ts"; import { Reader } from "../src/stream.ts"; const frameSizes = [16 * 1024, 256 * 1024, 1024 * 1024]; @@ -31,7 +32,7 @@ for (const frameSize of frameSizes) { controller.close(); }, }); - const reader = new Reader(stream); + const reader = new Reader(stream, undefined, Version.DRAFT_06); const start = performance.now(); const read = await reader.read(frameSize); diff --git a/js/net/bench/varint.ts b/js/net/bench/varint.ts index d8322a1c1e..1cc8812e85 100644 --- a/js/net/bench/varint.ts +++ b/js/net/bench/varint.ts @@ -15,7 +15,7 @@ let checksum = 0; const formats = [ { name: "quic", - version: undefined, + version: Version.DRAFT_16, encode: Varint.encodeTo, values: [2 ** 6 - 1, 2 ** 14 - 1, 2 ** 30 - 1, Number.MAX_SAFE_INTEGER], }, diff --git a/js/net/src/connection/accept.ts b/js/net/src/connection/accept.ts index d0f4f2c297..86a9ac33ad 100644 --- a/js/net/src/connection/accept.ts +++ b/js/net/src/connection/accept.ts @@ -142,7 +142,7 @@ async function acceptSetup( wiring: SessionProps, ): Promise { // Accept bidi, read ClientSetup, write ServerSetup - const stream = await Stream.accept(transport); + const stream = await Stream.accept(transport, version); if (!stream) throw new Error("no incoming bidi stream for SETUP"); const clientCompat = await stream.reader.u53(); @@ -187,7 +187,7 @@ async function acceptNegotiated( ): Promise { const setupVersion = Ietf.Version.DRAFT_14; - const stream = await Stream.accept(transport); + const stream = await Stream.accept(transport, setupVersion); if (!stream) throw new Error("no incoming bidi stream for SETUP"); const clientCompat = await stream.reader.u53(); diff --git a/js/net/src/connection/connect.ts b/js/net/src/connection/connect.ts index 759c3cceb0..493b101610 100644 --- a/js/net/src/connection/connect.ts +++ b/js/net/src/connection/connect.ts @@ -297,7 +297,7 @@ async function negotiate(url: URL, session: WebTransport, wiring: SessionProps): throw new Error(`unsupported WebTransport protocol: ${protocol}`); } - const stream = await Stream.open(session); + const stream = await Stream.open(session, { version: setupVersion }); await stream.writer.u53(Lite.StreamId.ClientCompat); const encoder = new TextEncoder(); diff --git a/js/net/src/ietf/adapter.ts b/js/net/src/ietf/adapter.ts index b51e48d8da..780b6a1b39 100644 --- a/js/net/src/ietf/adapter.ts +++ b/js/net/src/ietf/adapter.ts @@ -212,9 +212,7 @@ export class ControlStreamAdapter implements Session { }, }); - const stream = new Stream({ readable, writable: sendWritable }); - stream.reader.version = this.version; - stream.writer.version = this.version; + const stream = new Stream({ readable, writable: sendWritable, version: this.version }); return stream; } @@ -339,9 +337,7 @@ export class ControlStreamAdapter implements Session { const sendWritable = this.#createSendWritable(); - const stream = new Stream({ readable, writable: sendWritable }); - stream.reader.version = this.version; - stream.writer.version = this.version; + const stream = new Stream({ readable, writable: sendWritable, version: this.version }); this.#streams.set(requestId, { controller }); diff --git a/js/net/src/ietf/ietf.test.ts b/js/net/src/ietf/ietf.test.ts index ec6d906c20..cfd7d29525 100644 --- a/js/net/src/ietf/ietf.test.ts +++ b/js/net/src/ietf/ietf.test.ts @@ -1068,7 +1068,7 @@ test("TrackStatusRequest v17: round trip with requiredRequestIdDelta", async () // Helper to encode a namespace to raw bytes async function encodeNamespace(namespace: Path.Valid): Promise { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, Version.DRAFT_14); await Namespace.encode(writer, namespace); writer.close(); await writer.closed; @@ -1077,14 +1077,14 @@ async function encodeNamespace(namespace: Path.Valid): Promise { // Helper to decode a namespace from raw bytes async function decodeNamespace(bytes: Uint8Array): Promise { - const reader = new Reader(undefined, bytes); + const reader = new Reader(undefined, bytes, Version.DRAFT_14); return await Namespace.decode(reader); } // Helper to encode raw IETF namespace tuple fields async function encodeNamespaceTuple(parts: string[]): Promise { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, Version.DRAFT_14); await writer.u53(parts.length); for (const part of parts) await writer.string(part); writer.close(); diff --git a/js/net/src/ietf/object.ts b/js/net/src/ietf/object.ts index a5584da033..9d6a63b883 100644 --- a/js/net/src/ietf/object.ts +++ b/js/net/src/ietf/object.ts @@ -1,5 +1,5 @@ import { StreamCode, Stream as StreamError } from "../error.ts"; -import { type Cursor, type Reader, Writer } from "../stream.ts"; +import { asIetf, type Cursor, type Reader, Writer } from "../stream.ts"; import { Timescale, Timestamp } from "../time.ts"; import { type IetfVersion, Version } from "./version.ts"; @@ -47,12 +47,7 @@ function hasDeltaObjectPropertyTypes(version: IetfVersion | undefined): boolean } } -async function encodeObjectPropertyType( - w: Writer, - id: bigint, - prev: bigint, - version: IetfVersion | undefined, -): Promise { +async function encodeObjectPropertyType(w: Writer, id: bigint, prev: bigint, version: IetfVersion): Promise { const encoded = hasDeltaObjectPropertyTypes(version) ? id - prev : id; await w.u62(encoded); } @@ -65,7 +60,7 @@ async function encodeObjectTime( w: Writer, timestamp: Timestamp, timescale: Timescale, - version: IetfVersion | undefined, + version: IetfVersion, ): Promise { const value = Math.round((timestamp.value * timescale) / timestamp.scale); await encodeObjectPropertyType(w, PROP_TIMESTAMP, 0n, version); @@ -75,7 +70,7 @@ async function encodeObjectTime( async function encodeObjectExtensions( timestamp: Timestamp | undefined, timescale: Timescale, - version: IetfVersion | undefined, + version: IetfVersion, ): Promise { if (timestamp === undefined) { return new Uint8Array(); @@ -112,7 +107,7 @@ function decodeObjectTime(c: Cursor, timescale: Timescale): Timestamp | undefine while (c.remaining > 0) { const step = c.u62(); - const id = !hasDeltaObjectPropertyTypes(c.version) || first ? step : prevType + step; + const id = !hasDeltaObjectPropertyTypes(asIetf(c.version)) || first ? step : prevType + step; first = false; prevType = id; @@ -285,7 +280,7 @@ export class Frame { * `idDelta` is the first object's absolute Object ID and zero for every later one, so a * group whose head was trimmed by a filter still puts the true numbering on the wire. */ - async encode(w: Writer, flags: GroupFlags, timescale: Timescale, version = w.version, idDelta = 0): Promise { + async encode(w: Writer, flags: GroupFlags, timescale: Timescale, version: IetfVersion, idDelta = 0): Promise { await w.u53(idDelta); if (flags.hasExtensions) { @@ -398,7 +393,7 @@ export class FetchFrame { } /** Encode this object at `position`, stamping it in the track's timescale. */ - async encode(w: Writer, position: FetchPosition, timescale: Timescale, version = w.version): Promise { + async encode(w: Writer, position: FetchPosition, timescale: Timescale, version: IetfVersion): Promise { if (position.first) { // Include the priority too: "same as the prior object" has no prior to refer to. const properties = this.timestamp !== undefined ? FETCH_PROPERTIES : 0; diff --git a/js/net/src/lite/announce.test.ts b/js/net/src/lite/announce.test.ts index 69b34cb9e5..86b212d308 100644 --- a/js/net/src/lite/announce.test.ts +++ b/js/net/src/lite/announce.test.ts @@ -25,10 +25,11 @@ function concat(chunks: Uint8Array[]): Uint8Array { return out; } -async function bytes(f: (w: Writer) => Promise): Promise { +async function bytes(f: (w: Writer) => Promise, version: Version): Promise { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await f(writer); writer.close(); @@ -37,7 +38,11 @@ async function bytes(f: (w: Writer) => Promise): Promise { } async function roundTrip(msg: AnnounceBroadcast, version: Version): Promise { - const reader = new Reader(undefined, await bytes((w) => encodeAnnounceBroadcast(w, msg, version))); + const reader = new Reader( + undefined, + await bytes((w) => encodeAnnounceBroadcast(w, msg, version), version), + version, + ); return decodeAnnounceBroadcast(reader, version); } @@ -74,8 +79,8 @@ test("AnnounceBroadcast skips an unknown type on draft-06", async () => { await w.u53(4); await w.u53(1); await w.u8(0); - }); - const reader = new Reader(undefined, encoded); + }, Version.DRAFT_06); + const reader = new Reader(undefined, encoded, Version.DRAFT_06); expect(await decodeAnnounceBroadcast(reader, Version.DRAFT_06)).toEqual({ status: "skipped" }); }); @@ -91,37 +96,49 @@ test("AnnounceBroadcast drops the route cost before draft-06", async () => { test("AnnounceBroadcast rejects cross-version forms", async () => { await expect( - bytes((w) => encodeAnnounceBroadcast(w, { status: "endedId", id: 1n }, Version.DRAFT_05)), + bytes((w) => encodeAnnounceBroadcast(w, { status: "endedId", id: 1n }, Version.DRAFT_05), Version.DRAFT_05), ).rejects.toThrow(); await expect( - bytes((w) => encodeAnnounceBroadcast(w, { status: "restart", id: 1n, hops: [] }, Version.DRAFT_05)), + bytes( + (w) => encodeAnnounceBroadcast(w, { status: "restart", id: 1n, hops: [] }, Version.DRAFT_05), + Version.DRAFT_05, + ), ).rejects.toThrow(); await expect( - bytes((w) => encodeAnnounceBroadcast(w, { status: "ended", suffix: Path.from("room/cam") }, Version.DRAFT_06)), + bytes( + (w) => encodeAnnounceBroadcast(w, { status: "ended", suffix: Path.from("room/cam") }, Version.DRAFT_06), + Version.DRAFT_06, + ), ).rejects.toThrow(); }); test("AnnounceBroadcast accepts explicit restart status on draft-05", async () => { - const wire = await bytes((w) => - encodeAnnounceBroadcast(w, { status: "active", suffix: Path.from("room/cam"), hops: [] }, Version.DRAFT_05), + const wire = await bytes( + (w) => + encodeAnnounceBroadcast(w, { status: "active", suffix: Path.from("room/cam"), hops: [] }, Version.DRAFT_05), + Version.DRAFT_05, ); wire[1] = 2; - const got = await decodeAnnounceBroadcast(new Reader(undefined, wire), Version.DRAFT_05); + const got = await decodeAnnounceBroadcast(new Reader(undefined, wire, Version.DRAFT_05), Version.DRAFT_05); expect(got).toEqual({ status: "active", suffix: Path.from("room/cam"), hops: [] }); }); test("AnnounceBroadcast rejects explicit restart status before draft-05", async () => { - const wire = await bytes((w) => - encodeAnnounceBroadcast(w, { status: "active", suffix: Path.from("room/cam"), hops: [] }, Version.DRAFT_04), + const wire = await bytes( + (w) => + encodeAnnounceBroadcast(w, { status: "active", suffix: Path.from("room/cam"), hops: [] }, Version.DRAFT_04), + Version.DRAFT_04, ); wire[1] = 2; - await expect(decodeAnnounceBroadcast(new Reader(undefined, wire), Version.DRAFT_04)).rejects.toThrow(); + await expect( + decodeAnnounceBroadcast(new Reader(undefined, wire, Version.DRAFT_04), Version.DRAFT_04), + ).rejects.toThrow(); }); async function requestRoundTrip(msg: AnnounceRequest, version: Version): Promise { - const reader = new Reader(undefined, await bytes((w) => msg.encode(w, version))); + const reader = new Reader(undefined, await bytes((w) => msg.encode(w, version), version), version); return AnnounceRequest.decode(reader, version); } @@ -139,8 +156,8 @@ test("AnnounceRequest drops excludeHop on draft-06", async () => { const got = await requestRoundTrip(msg, Version.DRAFT_06); expect(got.excludeHop).toBe(0n); - const with05 = await bytes((w) => msg.encode(w, Version.DRAFT_05)); - const with06 = await bytes((w) => msg.encode(w, Version.DRAFT_06)); + const with05 = await bytes((w) => msg.encode(w, Version.DRAFT_05), Version.DRAFT_05); + const with06 = await bytes((w) => msg.encode(w, Version.DRAFT_06), Version.DRAFT_06); expect(with06.byteLength).toBeLessThan(with05.byteLength); }); @@ -158,7 +175,11 @@ test("AnnounceRequest carries hidden from draft-07", async () => { // conforming publisher. test("AnnounceOk accepts the reserved unknown origin", async () => { const msg = new AnnounceOk(UNKNOWN_HOP, 3); - const reader = new Reader(undefined, await bytes((w) => msg.encode(w, Version.DRAFT_05))); + const reader = new Reader( + undefined, + await bytes((w) => msg.encode(w, Version.DRAFT_05), Version.DRAFT_05), + Version.DRAFT_05, + ); const got = await AnnounceOk.decode(reader, Version.DRAFT_05); expect(got.hop).toBe(UNKNOWN_HOP); expect(got.active).toBe(3); @@ -166,7 +187,11 @@ test("AnnounceOk accepts the reserved unknown origin", async () => { test("AnnounceOk round-trips a declared origin", async () => { const msg = new AnnounceOk(HopSchema.parse(42n), 1); - const reader = new Reader(undefined, await bytes((w) => msg.encode(w, Version.DRAFT_05))); + const reader = new Reader( + undefined, + await bytes((w) => msg.encode(w, Version.DRAFT_05), Version.DRAFT_05), + Version.DRAFT_05, + ); const got = await AnnounceOk.decode(reader, Version.DRAFT_05); expect(got.hop).toBe(HopSchema.parse(42n)); expect(got.active).toBe(1); @@ -179,7 +204,9 @@ test("a hop chain that revisits a hop is refused in both directions", async () = // Outbound: refused before it reaches the wire. A receiver must close the session over // a repeated Hop ID, so sending one costs someone else their session. - await expect(bytes((w) => encodeAnnounceBroadcast(w, looped, Version.DRAFT_06))).rejects.toThrow("appears twice"); + await expect(bytes((w) => encodeAnnounceBroadcast(w, looped, Version.DRAFT_06), Version.DRAFT_06)).rejects.toThrow( + "appears twice", + ); // Inbound: encode a chain that is legal, then rewrite its last hop to repeat the // first. Only a non-conforming sender produces these bytes, which is why they have to @@ -189,7 +216,7 @@ test("a hop chain that revisits a hop is refused in both directions", async () = suffix: Path.from("room"), hops: [four, eight, HopSchema.parse(9n)], }; - const forged = await bytes((w) => encodeAnnounceBroadcast(w, legal, Version.DRAFT_06)); + const forged = await bytes((w) => encodeAnnounceBroadcast(w, legal, Version.DRAFT_06), Version.DRAFT_06); const nine = forged.lastIndexOf(9); expect(nine).toBeGreaterThan(0); forged[nine] = 4; @@ -197,12 +224,12 @@ test("a hop chain that revisits a hop is refused in both directions", async () = // The type carries the consequence, not just the text: the subscriber's dispatch closes // the session on `instanceof ProtocolViolation`, so a plain Error here would reset the // stream and leave a nonconforming peer free to repeat itself. - await expect(decodeAnnounceBroadcast(new Reader(undefined, forged), Version.DRAFT_06)).rejects.toThrow( - ProtocolViolation, - ); - await expect(decodeAnnounceBroadcast(new Reader(undefined, forged), Version.DRAFT_06)).rejects.toThrow( - "appears twice", - ); + await expect( + decodeAnnounceBroadcast(new Reader(undefined, forged, Version.DRAFT_06), Version.DRAFT_06), + ).rejects.toThrow(ProtocolViolation); + await expect( + decodeAnnounceBroadcast(new Reader(undefined, forged, Version.DRAFT_06), Version.DRAFT_06), + ).rejects.toThrow("appears twice"); // Repeated unknowns are not a loop: 0 identifies nothing, so any number of hops may // be unknown. A lite-03 announcement is nothing but these. @@ -225,7 +252,7 @@ function unhex(text: string): Uint8Array { // Resolve every announcement on a lite-07 stream, as the subscriber does. async function resolveStream(data: Uint8Array) { - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, Version.DRAFT_07); const history = new AnnounceHistory(); const out: unknown[] = []; for (;;) { @@ -260,7 +287,7 @@ const GOLDEN_RESOLVED = [ // Pinned from the Rust encoder (`lite::compress::tests::golden_stream_is_pinned`), so the // JS decoder is checked against real compressed output. const GOLDEN = - "001600000a726f6f6d2f612f63616d000251116222000000000d0102036d696301017333010000020a00010180004444010000010101000d02010162020180005555010000"; + "001600000a726f6f6d2f612f63616d00029111a222000000000d0102036d69630101b3330100000209000101c04444010000010101000c020101620201c05555010000"; test("AnnounceHistory resolves the Rust encoder's compressed stream", async () => { expect(await resolveStream(unhex(GOLDEN))).toEqual(GOLDEN_RESOLVED); @@ -268,7 +295,7 @@ test("AnnounceHistory resolves the Rust encoder's compressed stream", async () = // JS always encodes literally; Rust decodes these bytes too (`js_literal_stream_decodes`). const JS_LITERAL = - "001600000a726f6f6d2f612f63616d000251116222000000001600000a726f6f6d2f612f6d6963000273336222000000020c0000028000444462220000000101010014000006726f6f6d2f620002800055556222000000"; + "001600000a726f6f6d2f612f63616d00029111a222000000001600000a726f6f6d2f612f6d69630002b333a222000000020b000002c04444a2220000000101010013000006726f6f6d2f620002c05555a222000000"; test("the literal draft-07 stream matches what Rust decodes", async () => { const wire = await bytes(async (w) => { @@ -291,7 +318,7 @@ test("the literal draft-07 stream matches what Rust decodes", async () => { { status: "active", suffix: Path.from("room/b"), hops: [hop(0x5555n), relay], cost }, v, ); - }); + }, Version.DRAFT_07); expect(hex(wire)).toBe(JS_LITERAL); expect(await resolveStream(wire)).toEqual(GOLDEN_RESOLVED); }); @@ -335,10 +362,10 @@ test("a keep without a base is a violation on draft-07", async () => { await w.u53(0); await w.u53(8); for (const b of [0, 1, 0, 0, 0, 0, 0, 0]) await w.u8(b); - }); - await expect(decodeAnnounceBroadcast(new Reader(undefined, wire), Version.DRAFT_07)).rejects.toThrow( - ProtocolViolation, - ); + }, Version.DRAFT_07); + await expect( + decodeAnnounceBroadcast(new Reader(undefined, wire, Version.DRAFT_07), Version.DRAFT_07), + ).rejects.toThrow(ProtocolViolation); }); test("draft-06 has no room for a base", async () => { @@ -348,5 +375,31 @@ test("draft-06 has no room for a base", async () => { hops: [], pathBase: { distance: 1n, keep: 1 }, }; - await expect(bytes((w) => encodeAnnounceBroadcast(w, msg, Version.DRAFT_06))).rejects.toThrow(); + await expect(bytes((w) => encodeAnnounceBroadcast(w, msg, Version.DRAFT_06), Version.DRAFT_06)).rejects.toThrow(); +}); + +// Costs saturate at 2^62-1 on every version: a larger one from a lite-07 peer reads as the +// ceiling, and one set locally goes out as the ceiling, so a cost always forwards to lite-06. +test("route costs saturate at 2^62-1 on every version", async () => { + const ceiling = 2n ** 62n - 1n; + const huge = 2n ** 64n - 1n; + for (const version of [Version.DRAFT_06, Version.DRAFT_07]) { + const got = await roundTrip( + { status: "active", suffix: Path.from("x"), hops: [], cost: { warm: huge, cold: huge } }, + version, + ); + expect(got).toMatchObject({ cost: { warm: ceiling, cold: ceiling } }); + } + + // ANNOUNCE_START: path base, path keep, empty suffix, hop base, no hops, hop keep, then + // warm and cold at 2^64-1, which only lite-07's varints can carry. + const wire = await bytes(async (w) => { + await w.u53(0); + await w.u53(24); + for (const b of [0, 0, 0, 0, 0, 0]) await w.u8(b); + await w.u62(huge); + await w.u62(huge); + }, Version.DRAFT_07); + const got = await decodeAnnounceBroadcast(new Reader(undefined, wire, Version.DRAFT_07), Version.DRAFT_07); + expect(got).toMatchObject({ cost: { warm: ceiling, cold: ceiling } }); }); diff --git a/js/net/src/lite/announce.ts b/js/net/src/lite/announce.ts index da874389fa..69ee0c0944 100644 --- a/js/net/src/lite/announce.ts +++ b/js/net/src/lite/announce.ts @@ -176,16 +176,20 @@ async function decodeHopsBlock(r: Reader, version: Version): Promise<{ hops: Hop } // The route cost rides lite-06+ announcements as two varints, warm then cold; older -// versions carry neither. +// versions carry neither. Costs saturate at 2^62-1 on every version, even where lite-07's +// varints could carry more, so a cost always forwards to a peer on an older version. +const MAX_COST = 2n ** 62n - 1n; +const saturate = (v: bigint) => (v > MAX_COST ? MAX_COST : v); + async function encodeRouteCost(w: Writer, version: Version, cost: Cost | undefined) { if (!hasRouteCost(version)) return; - await w.u62(cost?.warm ?? 0n); - await w.u62(cost?.cold ?? 0n); + await w.u62(saturate(cost?.warm ?? 0n)); + await w.u62(saturate(cost?.cold ?? 0n)); } async function decodeRouteCost(r: Reader, version: Version): Promise { if (!hasRouteCost(version)) return undefined; - return { warm: await r.u62(), cold: await r.u62() }; + return { warm: saturate(await r.u62()), cold: saturate(await r.u62()) }; } // lite-06 message body (no discriminator; the type is carried outside the length prefix). diff --git a/js/net/src/lite/connection.ts b/js/net/src/lite/connection.ts index 4aec63c3a3..855ce50af8 100644 --- a/js/net/src/lite/connection.ts +++ b/js/net/src/lite/connection.ts @@ -116,6 +116,11 @@ export class Connection implements Established { this.url = url; this.#quic = quic; this.#session = session; + // The session stream was opened at the SETUP exchange's version; the rest of it speaks this one. + if (session) { + session.reader.version = version; + session.writer.version = version; + } this.version = versionName(version); this.#version = version; this.transport = transportOf(quic); @@ -204,7 +209,7 @@ export class Connection implements Established { // our session identity so the peer can filter // reflected announcements (lite-06 removed ANNOUNCE_REQUEST's exclude_hop for it). async #sendSetup(): Promise { - const writer = await Writer.open(this.#quic); + const writer = await Writer.open(this.#quic, { version: this.#version }); try { await writer.u53(DataType.Setup); const probe = await probeLevel(this.#quic, this.#version); @@ -218,7 +223,7 @@ export class Connection implements Established { async #runBidis() { for (;;) { - const stream = await Stream.accept(this.#quic); + const stream = await Stream.accept(this.#quic, this.#version); if (!stream) break; this.#runBidi(stream) @@ -259,7 +264,7 @@ export class Connection implements Established { } async #runUnis() { - const readers = new Readers(this.#quic); + const readers = new Readers(this.#quic, this.#version); for (;;) { const stream = await readers.next(); diff --git a/js/net/src/lite/datagram.test.ts b/js/net/src/lite/datagram.test.ts index ff2b8d04f7..80abce58db 100644 --- a/js/net/src/lite/datagram.test.ts +++ b/js/net/src/lite/datagram.test.ts @@ -1,28 +1,38 @@ import { expect, test } from "bun:test"; import { Datagram } from "./datagram.ts"; +import { Version } from "./version.ts"; const enc = new TextEncoder(); const dec = new TextDecoder(); test("datagram body round-trips", async () => { - const dg = new Datagram(7n, 42, 1000, enc.encode("hello")); - const decoded = await Datagram.decode(dg.encode()); - expect(decoded.subscribe).toBe(7n); - expect(decoded.sequence).toBe(42); - expect(decoded.timestamp).toBe(1000); - expect(dec.decode(decoded.payload)).toBe("hello"); + for (const version of [Version.DRAFT_06, Version.DRAFT_07]) { + const dg = new Datagram(7n, 42, 1000, enc.encode("hello")); + const decoded = await Datagram.decode(dg.encode(version), version); + expect(decoded.subscribe).toBe(7n); + expect(decoded.sequence).toBe(42); + expect(decoded.timestamp).toBe(1000); + expect(dec.decode(decoded.payload)).toBe("hello"); + } }); test("datagram body has no inner length prefix", () => { const dg = new Datagram(1n, 2, 3, enc.encode("world")); - const body = dg.encode(); + const body = dg.encode(Version.DRAFT_06); // Three single-byte varints (values < 64) followed by the raw 5-byte payload. expect(body.byteLength).toBe(8); expect(dec.decode(body.slice(3))).toBe("world"); }); +test("datagram body varints follow the version's encoding", () => { + // 100 takes the two-byte QUIC form on lite-06 and one leading-ones byte on lite-07. + const dg = new Datagram(100n, 100, 100, new Uint8Array()); + expect([...dg.encode(Version.DRAFT_06)]).toEqual([0x40, 0x64, 0x40, 0x64, 0x40, 0x64]); + expect([...dg.encode(Version.DRAFT_07)]).toEqual([0x64, 0x64, 0x64]); +}); + test("datagram body round-trips an empty payload", async () => { const dg = new Datagram(0n, 0, 0, new Uint8Array()); - const decoded = await Datagram.decode(dg.encode()); + const decoded = await Datagram.decode(dg.encode(Version.DRAFT_07), Version.DRAFT_07); expect(decoded.payload.byteLength).toBe(0); }); diff --git a/js/net/src/lite/datagram.ts b/js/net/src/lite/datagram.ts index 1f76707e17..6fbb0df9d0 100644 --- a/js/net/src/lite/datagram.ts +++ b/js/net/src/lite/datagram.ts @@ -3,8 +3,8 @@ * * @module */ -import { Reader } from "../stream.ts"; -import * as Varint from "../varint.ts"; +import { encodeVarint, Reader } from "../stream.ts"; +import type { Version } from "./version.ts"; /** * A QUIC datagram body: `subscribe (i) | sequence (i) | timestamp (i) | payload (b)`. @@ -31,10 +31,10 @@ export class Datagram { } /** Encode the body to a single `Uint8Array` (no length prefix; the datagram boundary delimits it). */ - encode(): Uint8Array { - const subscribe = Varint.encodeTo(new ArrayBuffer(8), this.subscribe); - const sequence = Varint.encodeTo(new ArrayBuffer(8), this.sequence); - const timestamp = Varint.encodeTo(new ArrayBuffer(8), this.timestamp); + encode(version: Version): Uint8Array { + const subscribe = encodeVarint(this.subscribe, version); + const sequence = encodeVarint(this.sequence, version); + const timestamp = encodeVarint(this.timestamp, version); const out = new Uint8Array( subscribe.byteLength + sequence.byteLength + timestamp.byteLength + this.payload.byteLength, @@ -51,8 +51,8 @@ export class Datagram { } /** Decode a datagram body from the raw bytes of one QUIC datagram. */ - static async decode(data: Uint8Array): Promise { - const r = new Reader(undefined, data); + static async decode(data: Uint8Array, version: Version): Promise { + const r = new Reader(undefined, data, version); const subscribe = await r.u62(); const sequence = await r.u53(); const timestamp = await r.u53(); diff --git a/js/net/src/lite/fetch.test.ts b/js/net/src/lite/fetch.test.ts index 0eb06dca94..c073a51cd9 100644 --- a/js/net/src/lite/fetch.test.ts +++ b/js/net/src/lite/fetch.test.ts @@ -19,6 +19,7 @@ async function encode(version: Version, fetch: Fetch): Promise { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await fetch.encode(writer, version); writer.close(); @@ -27,7 +28,7 @@ async function encode(version: Version, fetch: Fetch): Promise { } async function roundtrip(version: Version, fetch: Fetch): Promise { - const reader = new Reader(undefined, await encode(version, fetch)); + const reader = new Reader(undefined, await encode(version, fetch), version); return Fetch.decode(reader, version); } diff --git a/js/net/src/lite/goaway.test.ts b/js/net/src/lite/goaway.test.ts index 17d7a43800..a051b010ce 100644 --- a/js/net/src/lite/goaway.test.ts +++ b/js/net/src/lite/goaway.test.ts @@ -18,6 +18,7 @@ async function encode(msg: Goaway, version: Version): Promise { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await msg.encode(writer, version); writer.close(); @@ -26,7 +27,7 @@ async function encode(msg: Goaway, version: Version): Promise { } async function decode(bytes: Uint8Array, version: Version): Promise { - const reader = new Reader(undefined, bytes); + const reader = new Reader(undefined, bytes, version); return await Goaway.decode(reader, version); } diff --git a/js/net/src/lite/group.test.ts b/js/net/src/lite/group.test.ts index e0df4212ed..44e2c9b64e 100644 --- a/js/net/src/lite/group.test.ts +++ b/js/net/src/lite/group.test.ts @@ -3,9 +3,13 @@ import { Producer } from "../group.ts"; import { Reader } from "../stream.ts"; import * as Varint from "../varint.ts"; import { readFrames } from "./group.ts"; +import { Version } from "./version.ts"; const SCALE = 1000; +// The frames below are QUIC varints, so the streams read them as lite-06. +const VERSION = Version.DRAFT_06; + /** Frames as a group stream carries them: a zigzag timestamp delta when timestamped, then the sized payload. */ function encode(frames: { delta?: number; payload: number[] }[]): Uint8Array { const bytes: number[] = []; @@ -24,6 +28,8 @@ function streamOf(chunks: Uint8Array[]): Reader { controller.close(); }, }), + undefined, + VERSION, ); } @@ -86,7 +92,7 @@ test("a stream that ends inside a frame rejects", async () => { test("stops once the group closes", async () => { const producer = new Producer(0); // Never ends: only the close can stop the read. - const stream = new Reader(new ReadableStream()); + const stream = new Reader(new ReadableStream(), undefined, VERSION); const done = readFrames(stream, producer, SCALE); producer.close(); await done; diff --git a/js/net/src/lite/message.ts b/js/net/src/lite/message.ts index ba36ab80a7..b48552e554 100644 --- a/js/net/src/lite/message.ts +++ b/js/net/src/lite/message.ts @@ -31,6 +31,7 @@ export async function encode(writer: Writer, f: (w: Writer) => Promise, ma } }, }), + writer.version, ); await f(temp); @@ -55,7 +56,7 @@ export async function decode(reader: Reader, f: (r: Reader) => Promise, ma } const data = await reader.read(size); - const limit = new Reader(undefined, data); + const limit = new Reader(undefined, data, reader.version); const msg = await f(limit); // Check that we consumed exactly the right number of bytes diff --git a/js/net/src/lite/priority.test.ts b/js/net/src/lite/priority.test.ts index 5adf91897e..25733b381d 100644 --- a/js/net/src/lite/priority.test.ts +++ b/js/net/src/lite/priority.test.ts @@ -3,6 +3,7 @@ import { Producer as BroadcastProducer } from "../broadcast.ts"; import { type SendStream, Writer } from "../stream.ts"; import { Priority, sendOrder } from "./priority.ts"; +import { Version } from "./version.ts"; // The last position that still fits below the track priority. const MAX_POSITION = 2 ** 44 - 1; @@ -64,7 +65,7 @@ function ranking(options: { priority?: number; ordered?: boolean } = {}) { const open = () => { const stream = new WritableStream() as SendStream; streams.push(stream); - return new Writer(stream); + return new Writer(stream, Version.DRAFT_06); }; return { priority: new Priority(subscriber), open, streams, broadcast }; diff --git a/js/net/src/lite/probe.test.ts b/js/net/src/lite/probe.test.ts index 150f96c8c7..ff44a4f748 100644 --- a/js/net/src/lite/probe.test.ts +++ b/js/net/src/lite/probe.test.ts @@ -14,10 +14,11 @@ function concat(chunks: Uint8Array[]): Uint8Array { return out; } -async function bytes(f: (w: Writer) => Promise): Promise { +async function bytes(f: (w: Writer) => Promise, version: Version): Promise { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await f(writer); writer.close(); @@ -26,7 +27,7 @@ async function bytes(f: (w: Writer) => Promise): Promise { } async function roundTrip(msg: Probe, version: Version = Version.DRAFT_05): Promise { - const reader = new Reader(undefined, await bytes((w) => msg.encode(w, version))); + const reader = new Reader(undefined, await bytes((w) => msg.encode(w, version), version), version); const got = await Probe.decode(reader, version); expect(await reader.done()).toBe(true); return got; @@ -66,7 +67,9 @@ test("lite-03 carries the bitrate alone", async () => { // varint encoder converts with `BigInt`, which throws on one, and the publisher's // catch would then close the probe stream for the rest of the session. test("a fractional RTT must not reach the encoder", async () => { - await expect(bytes((w) => new Probe({ bitrate: 1_000, rtt: 12.34 }).encode(w, Version.DRAFT_05))).rejects.toThrow(); + await expect( + bytes((w) => new Probe({ bitrate: 1_000, rtt: 12.34 }).encode(w, Version.DRAFT_05), Version.DRAFT_05), + ).rejects.toThrow(); // Rounding first is what the publisher does, and it encodes cleanly. const got = await roundTrip(new Probe({ bitrate: 1_000, rtt: Math.round(12.34) })); diff --git a/js/net/src/lite/publisher.test.ts b/js/net/src/lite/publisher.test.ts index 2466166cd2..64d50dd980 100644 --- a/js/net/src/lite/publisher.test.ts +++ b/js/net/src/lite/publisher.test.ts @@ -45,6 +45,7 @@ test.each([Version.DRAFT_01, Version.DRAFT_03, Version.DRAFT_06])( await Promise.resolve(); const failure = new Error("peer stopped receiving announcements"); const stream = new Stream({ + version: version, readable: new ReadableStream(), writable: new WritableStream({ write() { @@ -96,6 +97,7 @@ test.each([ const written: Uint8Array[] = []; const stream = new Stream({ + version: version, readable: new ReadableStream(), writable: new WritableStream({ write(chunk) { @@ -132,8 +134,8 @@ async function subscribeEnd(sequences: number[], version: Version = Version.DRAF const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version }); + const server = await Stream.accept(pair.server, version); if (!server) throw new Error("publisher never accepted the subscribe stream"); const msg = new Subscribe({ id: 0n, broadcast: Path.from("test"), track: "video", priority: 0 }); @@ -195,8 +197,8 @@ async function groupSendOrders(options: { priority: number; sequences: number[]; const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the subscribe stream"); const msg = replaySubscribe({ @@ -279,8 +281,8 @@ test("lite draft-05: a subscribe update re-ranks a group already on the wire", a const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the subscribe stream"); const msg = replaySubscribe({ @@ -343,8 +345,8 @@ test("lite draft-05: a subscribe update during the stream open still ranks the g return open(options); }; - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the subscribe stream"); const msg = replaySubscribe({ @@ -392,8 +394,8 @@ test("lite draft-05: many concurrent groups share one subscription listener", as const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the subscribe stream"); const msg = replaySubscribe({ @@ -451,7 +453,7 @@ test("lite draft-05: the fetch response ranks the publisher's own writes", async group.close(); track.writeGroup(group); - const client = await Stream.open(pair.client); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); // Accept by hand rather than via Stream.accept, so the test keeps the writable the // publisher ranks (a real WebTransportSendStream takes the same assignment). @@ -459,7 +461,7 @@ test("lite draft-05: the fetch response ranks the publisher's own writes", async const accepted = await incoming.read(); incoming.releaseLock(); if (accepted.done) throw new Error("publisher never saw the fetch stream"); - const server = new Stream(accepted.value); + const server = new Stream({ ...accepted.value, version: Version.DRAFT_05 }); const msg = new Fetch({ broadcast: Path.from("test"), track: "video", priority: 3, group: 7 }); try { @@ -545,7 +547,7 @@ async function servedSubscription( const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video", { maxAge: options.maxAge }); - const client = await Stream.open(pair.client); + const client = await Stream.open(pair.client, { version: version }); // Accept by hand rather than via Stream.accept, so the gate sits between the publisher // and the wire. @@ -555,7 +557,7 @@ async function servedSubscription( if (accepted.done) throw new Error("publisher never accepted the subscribe stream"); const gate = gateWrites(accepted.value.writable, options.gated ?? false); - const server = new Stream({ readable: accepted.value.readable, writable: gate.writable }); + const server = new Stream({ readable: accepted.value.readable, writable: gate.writable, version: version }); const msg = replaySubscribe({ id: 0n, @@ -597,7 +599,7 @@ async function servedSubscription( clearTimeout(timer); if (!next || next.done) return undefined; - const reader = new Reader(next.value); + const reader = new Reader(next.value, undefined, version); await reader.u53(); // stream type const header = await GroupMessage.decode(reader, version); @@ -1046,8 +1048,8 @@ test("lite draft-07: subscribe end waits for groups below a declared finish", as const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_07 }); + const server = await Stream.accept(pair.server, Version.DRAFT_07); if (!server) throw new Error("publisher never accepted the subscribe stream"); void publisher.runSubscribe( new Subscribe({ id: 0n, broadcast: Path.from("test"), track: "video", priority: 0 }), @@ -1101,8 +1103,8 @@ async function heldOpenEnd() { return createUni(options); }); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_07 }); + const server = await Stream.accept(pair.server, Version.DRAFT_07); if (!server) throw new Error("publisher never accepted the subscribe stream"); void publisher.runSubscribe( new Subscribe({ id: 0n, broadcast: Path.from("test"), track: "video", priority: 0 }), @@ -1178,8 +1180,8 @@ async function serve( const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_06 }); + const server = await Stream.accept(pair.server, Version.DRAFT_06); if (!server) throw new Error("publisher never accepted the subscribe stream"); const msg = new Subscribe({ @@ -1221,7 +1223,7 @@ async function serve( clearTimeout(timer); if (!next || next.done) break; - const stream = new Reader(next.value); + const stream = new Reader(next.value, undefined, Version.DRAFT_06); await stream.u53(); // stream type const header = await GroupMessage.decode(stream, Version.DRAFT_06); @@ -1367,8 +1369,8 @@ async function saturatedGroup() { const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the subscribe stream"); const msg = new Subscribe({ id: 0n, broadcast: Path.from("test"), track: "video", priority: 0 }); @@ -1446,8 +1448,8 @@ test("lite draft-05: a blocked group header is reset when the group expires", as const publisher = new Publisher(pair.server, Version.DRAFT_05, randomHop(), origin.consume()); const broadcast = publish(origin, Path.from("test")); const track = broadcast.createTrack("video"); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the subscribe stream"); try { @@ -1533,8 +1535,8 @@ test("runProbe rounds a fractional smoothedRtt instead of killing the stream", a const publisher = new Publisher(pair.server, Version.DRAFT_05, randomHop()); // The subscriber opens the probe stream; the publisher only replies on it. - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the probe stream"); // `runProbe` loops until the stream closes, so close it rather than leaving the @@ -1603,8 +1605,8 @@ test("lite draft-05: a group that goes stale while its stream opens writes nothi return stale; }); - const client = await Stream.open(pair.client); - const server = await Stream.accept(pair.server); + const client = await Stream.open(pair.client, { version: Version.DRAFT_05 }); + const server = await Stream.accept(pair.server, Version.DRAFT_05); if (!server) throw new Error("publisher never accepted the subscribe stream"); try { diff --git a/js/net/src/lite/publisher.ts b/js/net/src/lite/publisher.ts index 7ed9622222..d8d5c57004 100644 --- a/js/net/src/lite/publisher.ts +++ b/js/net/src/lite/publisher.ts @@ -974,7 +974,7 @@ export class Publisher { // Convert the timestamp to the track's advertised timescale, matching #serveGroup. const ts = Math.round(datagram.timestamp.as(timescale)); - const body = new DatagramMessage(sub, datagram.sequence, ts, datagram.payload).encode(); + const body = new DatagramMessage(sub, datagram.sequence, ts, datagram.payload).encode(this.version); // No group fallback: drop anything that doesn't fit a single datagram. if (body.byteLength > maxSize) { @@ -1038,6 +1038,7 @@ export class Publisher { // in the order we asked, which is oldest-first, exactly backwards for live media. // Failing here drops the group and lets the next one compete for the next slot. const stream = await Writer.tryOpen(this.#quic, { + version: this.version, sendOrder: priority.rank(group.sequence), cancel: unsubscribed, waitUntilAvailable: false, diff --git a/js/net/src/lite/setup.test.ts b/js/net/src/lite/setup.test.ts index 77d2332ba1..7a7b805142 100644 --- a/js/net/src/lite/setup.test.ts +++ b/js/net/src/lite/setup.test.ts @@ -23,10 +23,11 @@ function concat(chunks: Uint8Array[]): Uint8Array { return out; } -async function bytes(f: (w: Writer) => Promise): Promise { +async function bytes(f: (w: Writer) => Promise, version: Version): Promise { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await f(writer); writer.close(); @@ -35,7 +36,11 @@ async function bytes(f: (w: Writer) => Promise): Promise { } async function roundTrip(msg: Setup): Promise { - const reader = new Reader(undefined, await bytes((w) => msg.encode(w, Version.DRAFT_05))); + const reader = new Reader( + undefined, + await bytes((w) => msg.encode(w, Version.DRAFT_05), Version.DRAFT_05), + Version.DRAFT_05, + ); const got = await Setup.decode(reader, Version.DRAFT_05); expect(await reader.done()).toBe(true); return got; @@ -49,14 +54,14 @@ async function decodeParam(id: bigint, value: Uint8Array): Promise { await w.u62(id); await w.u53(value.byteLength); if (value.byteLength > 0) await w.write(value); - }); + }, Version.DRAFT_05); const framed = await bytes(async (w) => { await w.u53(body.byteLength); // Message size prefix await w.write(body); - }); + }, Version.DRAFT_05); - return await Setup.decode(new Reader(undefined, framed), Version.DRAFT_05); + return await Setup.decode(new Reader(undefined, framed, Version.DRAFT_05), Version.DRAFT_05); } test("empty SETUP round-trips on draft-05", async () => { @@ -98,15 +103,18 @@ test("Both is the wire default, so it is omitted rather than encoded", async () // Both must be the absence of the parameter, so a server that predates it decodes a // Both client back to Both. Assert the empty bag directly: comparing against another // default-constructed Setup would agree with itself even if we encoded Both. - const both = await bytes((w) => new Setup({ role: Role.Both }).encode(w, Version.DRAFT_05)); + const both = await bytes((w) => new Setup({ role: Role.Both }).encode(w, Version.DRAFT_05), Version.DRAFT_05); const emptyBag = await bytes(async (w) => { await w.u53(1); // Message size prefix: one byte of body follows await w.u53(0); // parameter count - }); + }, Version.DRAFT_05); expect(both).toEqual(emptyBag); // A directional role is still encoded, so the check above can actually fail. - const publisher = await bytes((w) => new Setup({ role: Role.Publisher }).encode(w, Version.DRAFT_05)); + const publisher = await bytes( + (w) => new Setup({ role: Role.Publisher }).encode(w, Version.DRAFT_05), + Version.DRAFT_05, + ); expect(publisher.byteLength).toBeGreaterThan(both.byteLength); }); @@ -140,12 +148,12 @@ test("hop 0 decodes as absent", async () => { }); test("SETUP is rejected before draft-05", async () => { - await expect(bytes((w) => new Setup().encode(w, Version.DRAFT_04))).rejects.toThrow(); + await expect(bytes((w) => new Setup().encode(w, Version.DRAFT_04), Version.DRAFT_04)).rejects.toThrow(); }); test("SETUP decode is rejected before draft-05", async () => { - const framed = await bytes((w) => new Setup().encode(w, Version.DRAFT_05)); - await expect(Setup.decode(new Reader(undefined, framed), Version.DRAFT_04)).rejects.toThrow(); + const framed = await bytes((w) => new Setup().encode(w, Version.DRAFT_05), Version.DRAFT_05); + await expect(Setup.decode(new Reader(undefined, framed, Version.DRAFT_04), Version.DRAFT_04)).rejects.toThrow(); }); test("empty path decodes as the root", async () => { @@ -165,11 +173,13 @@ test("SETUP encode stops at the 64 KiB receive limit", async () => { expect((await roundTrip(msg)).path).toBe(msg.path); const over = new Setup({ path: "a".repeat(atLimit + 1) }); - await expect(bytes((w) => over.encode(w, Version.DRAFT_05))).rejects.toThrow("too large"); + await expect(bytes((w) => over.encode(w, Version.DRAFT_05), Version.DRAFT_05)).rejects.toThrow("too large"); }); test("SETUP decode refuses an oversized length before reading the body", async () => { // Only the prefix is present, so reading the body would fail with a different error. - const prefix = await bytes((w) => w.u53(64 * 1024 + 1)); - await expect(Setup.decode(new Reader(undefined, prefix), Version.DRAFT_05)).rejects.toThrow("too large"); + const prefix = await bytes((w) => w.u53(64 * 1024 + 1), Version.DRAFT_05); + await expect(Setup.decode(new Reader(undefined, prefix, Version.DRAFT_05), Version.DRAFT_05)).rejects.toThrow( + "too large", + ); }); diff --git a/js/net/src/lite/setup.ts b/js/net/src/lite/setup.ts index 37b724a4d0..22b61b4268 100644 --- a/js/net/src/lite/setup.ts +++ b/js/net/src/lite/setup.ts @@ -6,9 +6,8 @@ */ import { type Hop, HopSchema } from "../hop.ts"; -import type { Reader, Writer } from "../stream.ts"; +import { encodeVarint, Reader, type Writer } from "../stream.ts"; import { decodeUtf8 } from "../util/utf8.ts"; -import * as Varint from "../varint.ts"; import * as Message from "./message.ts"; import { hasSetupStream, type Version } from "./version.ts"; @@ -114,20 +113,21 @@ class Parameters { return this.#entries.get(id); } - /** Set a parameter to a varint value, replacing any existing entry. */ - setVarint(id: bigint, value: number | bigint) { - this.#entries.set(id, Varint.encode(Number(value))); + /** Set a parameter to a varint value in `version`'s encoding, replacing any existing entry. */ + setVarint(id: bigint, value: number | bigint, version: Version) { + this.#entries.set(id, encodeVarint(value, version)); } - /** Decode a parameter as a single varint, if present. Throws if trailing bytes remain. */ - getVarint(id: bigint): bigint | undefined { + /** Decode a parameter as a single varint in `version`'s encoding, if present. Throws if trailing bytes remain. */ + async getVarint(id: bigint, version: Version): Promise { const bytes = this.#entries.get(id); if (bytes === undefined) return undefined; - const [value, remain] = Varint.decode(bytes); - if (remain.byteLength !== 0) { + const r = new Reader(undefined, bytes, version); + const value = await r.u62(); + if (!(await r.done())) { throw new Error("trailing bytes after varint parameter"); } - return BigInt(value); + return value; } async encode(w: Writer) { @@ -225,29 +225,29 @@ export class Setup { } } - async #encode(w: Writer) { + async #encode(w: Writer, version: Version) { const params = new Parameters(); // None is the wire default, so omit it to keep the message empty when nothing is set. if (this.probe !== ProbeLevel.None) { - params.setVarint(PARAM_PROBE, this.probe); + params.setVarint(PARAM_PROBE, this.probe, version); } if (this.path !== undefined) { params.setBytes(PARAM_PATH, new TextEncoder().encode(this.path)); } // Both is the wire default, sent as the absence of the parameter. if (this.role !== Role.Both) { - params.setVarint(PARAM_ROLE, this.role); + params.setVarint(PARAM_ROLE, this.role, version); } if (this.hop !== undefined && this.hop !== 0n) { - params.setVarint(PARAM_HOP, this.hop); + params.setVarint(PARAM_HOP, this.hop, version); } await params.encode(w); } - static async #decode(r: Reader): Promise { + static async #decode(r: Reader, version: Version): Promise { const params = await Parameters.decode(r); - const probeCode = params.getVarint(PARAM_PROBE); + const probeCode = await params.getVarint(PARAM_PROBE, version); const probe = probeCode === undefined ? ProbeLevel.None : probeFromCode(probeCode); // An empty path is valid and means the same as omitting the parameter, so a @@ -255,11 +255,11 @@ export class Setup { const pathBytes = params.getBytes(PARAM_PATH); const path = pathBytes === undefined ? undefined : decodeUtf8(pathBytes); - const roleCode = params.getVarint(PARAM_ROLE); + const roleCode = await params.getVarint(PARAM_ROLE, version); const role = roleCode === undefined ? Role.Both : roleFromCode(roleCode); // 0 carries no identity (it cannot be excluded), so it decodes as absent. - const hopRaw = params.getVarint(PARAM_HOP); + const hopRaw = await params.getVarint(PARAM_HOP, version); const hop = hopRaw === undefined || hopRaw === 0n ? undefined : HopSchema.parse(hopRaw); return new Setup({ probe, path, role, hop }); @@ -268,12 +268,12 @@ export class Setup { /** Encode the SETUP message with its size prefix. Throws on pre-lite-05 versions. */ async encode(w: Writer, version: Version): Promise { Setup.#guard(version); - return Message.encode(w, this.#encode.bind(this), MAX_SETUP_SIZE); + return Message.encode(w, (w) => this.#encode(w, version), MAX_SETUP_SIZE); } /** Decode a SETUP message with its size prefix. Throws on pre-lite-05 versions. */ static async decode(r: Reader, version: Version): Promise { Setup.#guard(version); - return Message.decode(r, Setup.#decode, MAX_SETUP_SIZE); + return Message.decode(r, (r) => Setup.#decode(r, version), MAX_SETUP_SIZE); } } diff --git a/js/net/src/lite/subscribe.test.ts b/js/net/src/lite/subscribe.test.ts index afd4c7c0a6..4da1af0918 100644 --- a/js/net/src/lite/subscribe.test.ts +++ b/js/net/src/lite/subscribe.test.ts @@ -32,6 +32,7 @@ async function encode(version: Version, resp: SubscribeResponse): Promise({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await encodeSubscribeResponse(writer, resp, version); writer.close(); @@ -40,12 +41,12 @@ async function encode(version: Version, resp: SubscribeResponse): Promise { - const reader = new Reader(undefined, await encode(version, resp)); + const reader = new Reader(undefined, await encode(version, resp), version); return decodeSubscribeResponse(reader, version); } async function encodeSubscribe(msg: Subscribe): Promise { - const writer = new Writer(new WritableStream()); + const writer = new Writer(new WritableStream(), Version.DRAFT_06); try { await msg.encode(writer, Version.DRAFT_06); } finally { @@ -60,6 +61,7 @@ async function encodeMessage( const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await message.encode(writer, version); writer.close(); @@ -90,7 +92,7 @@ test("Subscribe round-trips every option including startGroup 0", async () => { // Lite-06 carries the raw floor, and a floor of 0 with no frame offset is the same // absence of a constraint as no floor at all, so it canonicalizes to undefined. const got = await Subscribe.decode( - new Reader(undefined, await encodeMessage(Version.DRAFT_06, message)), + new Reader(undefined, await encodeMessage(Version.DRAFT_06, message), Version.DRAFT_06), Version.DRAFT_06, ); expect(got.priority).toBe(7); @@ -102,7 +104,7 @@ test("Subscribe round-trips every option including startGroup 0", async () => { // partway through group 0 (a catalog never leaves it). message.startFrame = 4; const resumed = await Subscribe.decode( - new Reader(undefined, await encodeMessage(Version.DRAFT_06, message)), + new Reader(undefined, await encodeMessage(Version.DRAFT_06, message), Version.DRAFT_06), Version.DRAFT_06, ); expect(resumed.startGroup).toBe(0); @@ -112,7 +114,7 @@ test("Subscribe round-trips every option including startGroup 0", async () => { // A pre-06 wire folds the vacuous floor back to absent: an explicit group 0 there // would mean "replay from the beginning", which is not what a floor of 0 asks for. const folded = await Subscribe.decode( - new Reader(undefined, await encodeMessage(Version.DRAFT_05, message)), + new Reader(undefined, await encodeMessage(Version.DRAFT_05, message), Version.DRAFT_05), Version.DRAFT_05, ); expect(folded.startGroup).toBeUndefined(); @@ -127,7 +129,7 @@ test("SubscribeUpdate round-trips every option including startGroup 0", async () endGroup: 12, }); const got = await SubscribeUpdate.decode( - new Reader(undefined, await encodeMessage(Version.DRAFT_06, message)), + new Reader(undefined, await encodeMessage(Version.DRAFT_06, message), Version.DRAFT_06), Version.DRAFT_06, ); expect(got.priority).toBe(8); @@ -137,7 +139,7 @@ test("SubscribeUpdate round-trips every option including startGroup 0", async () // The same fold as SUBSCRIBE on a pre-06 wire. const folded = await SubscribeUpdate.decode( - new Reader(undefined, await encodeMessage(Version.DRAFT_05, message)), + new Reader(undefined, await encodeMessage(Version.DRAFT_05, message), Version.DRAFT_05), Version.DRAFT_05, ); expect(folded.startGroup).toBeUndefined(); @@ -173,9 +175,9 @@ test("SubscribeDrop is gone on draft-07", async () => { // A draft-06 DROP is an unknown response type on draft-07. const wire06 = await encode(Version.DRAFT_06, drop); - await expect(decodeSubscribeResponse(new Reader(undefined, wire06), Version.DRAFT_07)).rejects.toThrow( - "unknown subscribe response type: 2", - ); + await expect( + decodeSubscribeResponse(new Reader(undefined, wire06, Version.DRAFT_07), Version.DRAFT_07), + ).rejects.toThrow("unknown subscribe response type: 2"); }); test("SubscribeDrop is type 0x2 on draft-05 and 0x1 on draft-04", async () => { diff --git a/js/net/src/lite/subscriber.test.ts b/js/net/src/lite/subscriber.test.ts index 03ba18cad5..de0190a75f 100644 --- a/js/net/src/lite/subscriber.test.ts +++ b/js/net/src/lite/subscriber.test.ts @@ -72,6 +72,7 @@ function announceHarness(version: Version, origin = 1n) { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await f(writer); writer.close(); @@ -392,6 +393,7 @@ async function probeBytes(probes: Probe[], version: Version): Promise({ write: (chunk) => void chunks.push(new Uint8Array(chunk)) }), + version, ); for (const probe of probes) await probe.encode(writer, version); writer.close(); @@ -775,6 +777,7 @@ async function answerTrackInfo(stream: FakeStream): Promise { const chunks: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void chunks.push(new Uint8Array(chunk)) }), + Version.DRAFT_05, ); await new TrackInfo({}).encode(writer, Version.DRAFT_05); for (const chunk of chunks) stream.inbound.enqueue(chunk); diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index b81d761ee1..30a6e03cdb 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -217,7 +217,7 @@ export class Subscriber { // to reset the stream, not just close our side of it. let stream: Stream; try { - stream = await Stream.open(this.#quic); + stream = await Stream.open(this.#quic, { version: this.version }); } catch (err: unknown) { announced.close(error(err)); return; @@ -663,7 +663,7 @@ export class Subscriber { }; this.#subscribes.set(id, entry); - state.stream = await Stream.open(this.#quic); + state.stream = await Stream.open(this.#quic, { version: this.version }); await state.stream.writer.u53(StreamId.Subscribe); await msg.encode(state.stream.writer, this.version); @@ -680,7 +680,7 @@ export class Subscriber { // Opens a TRACK stream, reads the single TRACK_INFO, and FINs. Lite-05+ only. async #trackInfo(broadcast: Path.Valid, track: string): Promise { - return this.#exchange(undefined, async (stream) => { + return this.#exchange({ version: this.version }, async (stream) => { await stream.writer.u53(StreamId.Track); await new TrackMessage(broadcast, track).encode(stream.writer, this.version); const info = await TrackInfo.decode(stream.reader, this.version); @@ -693,7 +693,7 @@ export class Subscriber { // Opens a stream and runs a request/response exchange on it, resetting the stream if `run` // fails. Subscriber.close() also resets it while `run` is pending, so a peer that never // answers cannot hold it open, and a stream that opens after the close is reset at once. - async #exchange(options: OpenOptions | undefined, run: (stream: Stream) => Promise): Promise { + async #exchange(options: OpenOptions, run: (stream: Stream) => Promise): Promise { const closed = this.#closed.signal; closed.throwIfAborted(); const stream = await Stream.open(this.#quic, options); @@ -820,7 +820,7 @@ export class Subscriber { group: netGroup.Producer, ): Promise<{ stream: Stream; info: TrackInfo }> { const info = await untilClosed(group, this.#trackInfo(broadcast, track)); - return this.#exchange({ sendOrder: sendOrder({ priority }) }, async (stream) => { + return this.#exchange({ sendOrder: sendOrder({ priority }), version: this.version }, async (stream) => { await stream.writer.u53(StreamId.Fetch); await new FetchMessage({ broadcast, track, priority, group: sequence }).encode(stream.writer, this.version); // A byte or an empty-group FIN accepts the fetch; a reset rejects it. @@ -1127,7 +1127,7 @@ export class Subscriber { // Decode one datagram body and hand it to the matching subscription's producer. Drops the // datagram (best-effort) if the subscription is unknown/closed or its timescale isn't resolved. async #routeDatagram(payload: Uint8Array): Promise { - const dg = await DatagramMessage.decode(payload); + const dg = await DatagramMessage.decode(payload, this.version); const entry = this.#subscribes.get(dg.subscribe); if (!entry) return; // Unknown or already-closed subscription. @@ -1180,7 +1180,7 @@ export class Subscriber { // transport hiccup) MUST NOT tear down the connection. On error, drop the // estimates so consumers know they're stale. try { - const stream = await Stream.open(this.#quic); + const stream = await Stream.open(this.#quic, { version: this.version }); await stream.writer.u53(StreamId.Probe); for (;;) { diff --git a/js/net/src/lite/tail.test.ts b/js/net/src/lite/tail.test.ts index 87e0a9a474..b6717c2860 100644 --- a/js/net/src/lite/tail.test.ts +++ b/js/net/src/lite/tail.test.ts @@ -37,7 +37,7 @@ function groupStream(subscriber: Subscriber, sequence: number) { const readable = new ReadableStream({ start: (c) => (controller = c) }); const handled = subscriber.runGroup( new GroupMessage({ subscribe: 0n, sequence }), - new Reader(readable, undefined, undefined), + new Reader(readable, undefined, subscriber.version), ); return { write: (payload: string) => controller.enqueue(frame(payload)), @@ -57,14 +57,14 @@ async function subscribed(version: Version, maxAge = GRACE, groups?: Groups) { const subscriber = new Subscriber(pair.client, version, randomHop()); const reader = subscriber.consume(Path.from("room")).track("video").subscribe({ maxAge, groups }); - const info = await Stream.accept(pair.server); + const info = await Stream.accept(pair.server, version); if (!info) throw new Error("the subscriber never asked for TRACK_INFO"); expect(await info.reader.u53()).toBe(StreamId.Track); await TrackMessage.decode(info.reader, version); await new TrackInfo({ maxAge: 60_000 }).encode(info.writer, version); info.close(); - const sub = await Stream.accept(pair.server); + const sub = await Stream.accept(pair.server, version); if (!sub) throw new Error("the subscriber never subscribed"); expect(await sub.reader.u53()).toBe(StreamId.Subscribe); await Subscribe.decode(sub.reader, version); diff --git a/js/net/src/lite/track.test.ts b/js/net/src/lite/track.test.ts index 6e12024335..ca4d25d44f 100644 --- a/js/net/src/lite/track.test.ts +++ b/js/net/src/lite/track.test.ts @@ -16,10 +16,11 @@ function concat(chunks: Uint8Array[]): Uint8Array { return out; } -async function bytes(f: (w: Writer) => Promise): Promise { +async function bytes(f: (w: Writer) => Promise, version: Version): Promise { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + version, ); await f(writer); writer.close(); @@ -33,7 +34,11 @@ test("TrackInfo round-trips on draft-05", async () => { maxAge: 2000, timescale: 90000, }); - const reader = new Reader(undefined, await bytes((w) => info.encode(w, Version.DRAFT_05))); + const reader = new Reader( + undefined, + await bytes((w) => info.encode(w, Version.DRAFT_05), Version.DRAFT_05), + Version.DRAFT_05, + ); const got = await TrackInfo.decode(reader, Version.DRAFT_05); expect(got.priority).toBe(7); expect(got.maxAge).toBe(2000); @@ -42,14 +47,18 @@ test("TrackInfo round-trips on draft-05", async () => { test("TrackInfo defaults match cross-language wire bytes", async () => { const info = new TrackInfo(infoDefaults()); - expect(await bytes((w) => info.encode(w, Version.DRAFT_05))).toEqual( + expect(await bytes((w) => info.encode(w, Version.DRAFT_05), Version.DRAFT_05)).toEqual( new Uint8Array([0x06, 0x00, 0x00, 0x53, 0x88, 0x43, 0xe8]), ); }); test("Track request round-trips on draft-05", async () => { const msg = new Track(Path.from("room"), "video"); - const reader = new Reader(undefined, await bytes((w) => msg.encode(w, Version.DRAFT_05))); + const reader = new Reader( + undefined, + await bytes((w) => msg.encode(w, Version.DRAFT_05), Version.DRAFT_05), + Version.DRAFT_05, + ); const got = await Track.decode(reader, Version.DRAFT_05); expect(got.broadcast).toBe(Path.from("room")); expect(got.track).toBe("video"); @@ -57,7 +66,7 @@ test("Track request round-trips on draft-05", async () => { test("TRACK_INFO is rejected before draft-05", async () => { const info = new TrackInfo({ timescale: 90000 }); - await expect(bytes((w) => info.encode(w, Version.DRAFT_04))).rejects.toThrow(); + await expect(bytes((w) => info.encode(w, Version.DRAFT_04), Version.DRAFT_04)).rejects.toThrow(); }); test("TrackInfo round-trips valid boundary metadata", async () => { @@ -66,7 +75,11 @@ test("TrackInfo round-trips valid boundary metadata", async () => { maxAge: Number.MAX_SAFE_INTEGER, timescale: Number.MAX_SAFE_INTEGER, }); - const reader = new Reader(undefined, await bytes((w) => info.encode(w, Version.DRAFT_05))); + const reader = new Reader( + undefined, + await bytes((w) => info.encode(w, Version.DRAFT_05), Version.DRAFT_05), + Version.DRAFT_05, + ); const got = await TrackInfo.decode(reader, Version.DRAFT_05); expect(got.priority).toBe(255); expect(got.maxAge).toBe(Number.MAX_SAFE_INTEGER); @@ -115,6 +128,7 @@ test("mutated TrackInfo fields emit no bytes on encode", async () => { const written: Uint8Array[] = []; const writer = new Writer( new WritableStream({ write: (chunk) => void written.push(new Uint8Array(chunk)) }), + Version.DRAFT_05, ); await expect(info.encode(writer, Version.DRAFT_05)).rejects.toThrow(RangeError); writer.close(); diff --git a/js/net/src/stream.test.ts b/js/net/src/stream.test.ts index 1c54ff92c8..c2b770f2f3 100644 --- a/js/net/src/stream.test.ts +++ b/js/net/src/stream.test.ts @@ -11,10 +11,14 @@ import { StreamError, } from "./error.ts"; import { Version } from "./ietf/version.ts"; -import { type Cursor, Reader, Stream, Writer } from "./stream.ts"; +import { Version as Lite } from "./lite/version.ts"; +import { type Cursor, Reader, Stream, type StreamVersion, Writer } from "./stream.ts"; import { TimeoutError } from "./util/timeout.ts"; import { U64 } from "./util/u64.ts"; +// The generic tests below read and write QUIC varints on the moq-lite stream code registry. +const QUIC = Lite.DRAFT_06; + // Helper to create a writable stream that captures written data function createTestWritableStream(): { stream: WritableStream; written: Uint8Array[] } { const written: Uint8Array[] = []; @@ -40,7 +44,7 @@ function concatChunks(chunks: Uint8Array[]): Uint8Array { test("Writer u8", async () => { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await writer.u8(42); await writer.u8(255); @@ -56,7 +60,7 @@ test("Writer u8", async () => { test("Writer u8 refuses out of range values before emitting bytes", async () => { for (const value of [-1, 1.5, 256, Number.NaN, Number.POSITIVE_INFINITY]) { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await expect(writer.u8(value)).rejects.toThrow(RangeError); writer.close(); await writer.closed; @@ -66,7 +70,7 @@ test("Writer u8 refuses out of range values before emitting bytes", async () => test("Writer i32", async () => { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await writer.i32(0); await writer.i32(-1); @@ -87,7 +91,7 @@ test("Writer i32", async () => { test("Writer u53", async () => { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await writer.u53(0); await writer.u53(63); // MAX_U6 @@ -119,7 +123,7 @@ test("Writer u53 refuses unsafe values before emitting bytes", async () => { 2 ** 62 + 1, ]) { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await expect(writer.u53(value)).rejects.toThrow(RangeError); writer.close(); await writer.closed; @@ -129,7 +133,7 @@ test("Writer u53 refuses unsafe values before emitting bytes", async () => { test("Writer string", async () => { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await writer.string("hello"); await writer.string("🎉"); @@ -140,7 +144,7 @@ test("Writer string", async () => { const result = concatChunks(written); // Create a reader to parse the result - const reader = new Reader(undefined, result); + const reader = new Reader(undefined, result, QUIC); const str1 = await reader.string(); const str2 = await reader.string(); @@ -150,14 +154,14 @@ test("Writer string", async () => { }); test("Reader string rejects malformed UTF-8", async () => { - const reader = new Reader(undefined, new Uint8Array([2, 0xc3, 0x28])); + const reader = new Reader(undefined, new Uint8Array([2, 0xc3, 0x28]), QUIC); await expect(reader.string()).rejects.toThrow(); }); test("Reader u8", async () => { const data = new Uint8Array([42, 255, 0, 128]); - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, QUIC); expect(await reader.u8()).toBe(42); expect(await reader.u8()).toBe(255); @@ -169,7 +173,7 @@ test("Reader u8", async () => { test("Reader read with exact sizes", async () => { const data = new Uint8Array([1, 2, 3, 4, 5, 6, 7, 8]); - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, QUIC); const chunk1 = await reader.read(3); expect(chunk1).toEqual(new Uint8Array([1, 2, 3])); @@ -185,7 +189,7 @@ test("Reader read with exact sizes", async () => { test("Reader read with zero size", async () => { const data = new Uint8Array([1, 2, 3]); - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, QUIC); const chunk = await reader.read(0); expect(chunk).toEqual(new Uint8Array([])); @@ -197,7 +201,7 @@ test("Reader read with zero size", async () => { test("Reader readAll", async () => { const data = new Uint8Array([1, 2, 3, 4, 5]); - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, QUIC); // Read some data first await reader.read(2); @@ -214,7 +218,7 @@ test("Reader u53 varint decoding", async () => { const testValues = [0, 63, 64, 16383, 16384, 1073741823, 1073741824, Number.MAX_SAFE_INTEGER]; const { stream, written } = createTestWritableStream(); - const testWriter = new Writer(stream); + const testWriter = new Writer(stream, QUIC); for (const value of testValues) { await testWriter.u53(value); @@ -224,7 +228,7 @@ test("Reader u53 varint decoding", async () => { await testWriter.closed; const data = concatChunks(written); - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, QUIC); for (const expectedValue of testValues) { const actualValue = await reader.u53(); @@ -241,13 +245,13 @@ test("Reader u53 rejects integers that cannot be represented exactly", async () for (const value of wireValues) { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await writer.u62(value); writer.close(); await writer.closed; // A failed decode consumes nothing; the stream is unusable after it anyway. - const reader = new Reader(undefined, concatChunks(written)); + const reader = new Reader(undefined, concatChunks(written), QUIC); await expect(reader.u53()).rejects.toThrow(`value larger than 53-bits: ${value}`); expect(await reader.done()).toBe(false); } @@ -257,7 +261,7 @@ test("Reader u62 varint decoding", async () => { const testValues = [0n, 63n, 64n, 16383n, 16384n, 1073741823n, 1073741824n, 9007199254740991n]; // MAX_U53 const { stream, written } = createTestWritableStream(); - const testWriter = new Writer(stream); + const testWriter = new Writer(stream, QUIC); for (const value of testValues) { await testWriter.u62(value); @@ -267,7 +271,7 @@ test("Reader u62 varint decoding", async () => { await testWriter.closed; const data = concatChunks(written); - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, QUIC); for (const expectedValue of testValues) { const actualValue = await reader.u62(); @@ -281,7 +285,7 @@ test("Reader string decoding", async () => { const testStrings = ["hello", "🎉", "", "world with spaces", "multi\nline\nstring"]; const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); for (const str of testStrings) { await writer.string(str); @@ -291,7 +295,7 @@ test("Reader string decoding", async () => { await writer.closed; const data = concatChunks(written); - const reader = new Reader(undefined, data); + const reader = new Reader(undefined, data, QUIC); for (const expectedString of testStrings) { const actualString = await reader.string(); @@ -313,7 +317,7 @@ test("Reader from stream", async () => { }, }); - const reader = new Reader(stream); + const reader = new Reader(stream, undefined, QUIC); // Read all data const result = await reader.readAll(); @@ -332,7 +336,7 @@ test("Reader stream with partial reads", async () => { }, }); - const reader = new Reader(stream); + const reader = new Reader(stream, undefined, QUIC); // Read specific amounts that cross chunk boundaries const first = await reader.read(3); // Should span first two chunks @@ -357,7 +361,7 @@ test("Reader preserves returned views across fills", async () => { controller.close(); }, }); - const reader = new Reader(stream); + const reader = new Reader(stream, undefined, QUIC); const head = await reader.read(1); expect(head).toEqual(new Uint8Array([1])); const joined = await reader.read(3); @@ -376,7 +380,7 @@ test("Reader returns a view of a chunk that already holds the read", async () => controller.close(); }, }); - const reader = new Reader(stream); + const reader = new Reader(stream, undefined, QUIC); const read = await reader.read(3); expect(read.buffer).toBe(chunk.buffer); expect(read).toEqual(new Uint8Array([1, 2, 3])); @@ -390,7 +394,7 @@ test("Reader joins every chunk a read spans", async () => { controller.close(); }, }); - const reader = new Reader(stream); + const reader = new Reader(stream, undefined, QUIC); expect(await reader.u8()).toBe(0); expect(await reader.read(98)).toEqual(Uint8Array.from({ length: 98 }, (_, index) => index + 1)); expect(await reader.done()).toBe(false); @@ -400,13 +404,13 @@ test("Reader joins every chunk a read spans", async () => { test("Reader u53 decodes two-byte stream type prefixes", async () => { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, QUIC); await writer.u53(0x40); writer.close(); await writer.closed; - const reader = new Reader(undefined, concatChunks(written)); + const reader = new Reader(undefined, concatChunks(written), QUIC); expect(await reader.u53()).toBe(0x40); expect(await reader.done()).toBe(true); }); @@ -416,7 +420,7 @@ const sized = (c: Cursor) => c.read(c.u53()); test("Reader tryDecode drains every buffered message, then consumes nothing from a partial one", async () => { let controller!: ReadableStreamDefaultController; - const reader = new Reader(new ReadableStream({ start: (c) => (controller = c) })); + const reader = new Reader(new ReadableStream({ start: (c) => (controller = c) }), undefined, QUIC); controller.enqueue(new Uint8Array([1, 0xa, 2, 0xb, 0xc, 3, 0xd])); expect(await reader.done()).toBe(false); @@ -434,13 +438,13 @@ test("Reader tryDecode drains every buffered message, then consumes nothing from }); test("Reader tryDecode holds only the decode that ran short to the bytes it needs", () => { - const reader = new Reader(undefined, new Uint8Array([2, 0xa])); + const reader = new Reader(undefined, new Uint8Array([2, 0xa]), QUIC); expect(reader.tryDecode(sized)).toBeUndefined(); expect(reader.tryDecode((c) => c.u8())).toBe(2); }); test("Reader decode rejects a stream that ends inside a message", async () => { - const reader = new Reader(undefined, new Uint8Array([3, 0xa])); + const reader = new Reader(undefined, new Uint8Array([3, 0xa]), QUIC); expect(reader.tryDecode(sized)).toBeUndefined(); await expect(reader.decode(sized)).rejects.toThrow("unexpected end of stream"); }); @@ -449,13 +453,13 @@ test("Reader refuses an oversized value even when it is already buffered", async const size = 64 * 1024 * 1024 + 1; const buffer = new Uint8Array(4 + size); buffer.set([0x84, 0x00, 0x00, 0x01]); // the 4-byte varint for size - await expect(new Reader(undefined, buffer).string()).rejects.toThrow("exceeds max size"); - await expect(new Reader(undefined, buffer.subarray(4)).read(size)).rejects.toThrow("exceeds max size"); + await expect(new Reader(undefined, buffer, QUIC).string()).rejects.toThrow("exceeds max size"); + await expect(new Reader(undefined, buffer.subarray(4), QUIC).read(size)).rejects.toThrow("exceeds max size"); }); test("Reader refuses a buffered decode whose fields together exceed the max size", async () => { const half = 32 * 1024 * 1024; - const reader = new Reader(undefined, new Uint8Array(2 * half + 1)); + const reader = new Reader(undefined, new Uint8Array(2 * half + 1), QUIC); await expect(reader.decode((c) => [c.read(half), c.read(half + 1)])).rejects.toThrow("exceeds max size"); }); @@ -477,6 +481,8 @@ test("Reader closed rejects with the decoded reset code", async () => { new ReadableStream({ start: (controller) => controller.error(new Reset(2)), }), + undefined, + QUIC, ); const err = await reader.closed.then( @@ -492,6 +498,7 @@ test("Writer closed rejects with the decoded reset code", async () => { new WritableStream({ start: (controller) => controller.error(new Reset(31)), }), + QUIC, ); const err = await writer.closed.then( @@ -504,7 +511,7 @@ test("Writer closed rejects with the decoded reset code", async () => { test("Writer reset forwards only a stream error code", async () => { const session = new SessionError(SessionCode.Unauthorized); - const sessionWriter = new Writer(new WritableStream()); + const sessionWriter = new Writer(new WritableStream(), QUIC); sessionWriter.reset(session); const sessionResult = await sessionWriter.closed.then( () => undefined, @@ -513,7 +520,7 @@ test("Writer reset forwards only a stream error code", async () => { expect(sessionResult).toBe(session); const stream = new StreamError(StreamCode.DeliveryTimeout); - const streamWriter = new Writer(new WritableStream()); + const streamWriter = new Writer(new WritableStream(), QUIC); streamWriter.reset(stream); const streamResult = await streamWriter.closed.then( () => undefined, @@ -524,7 +531,7 @@ test("Writer reset forwards only a stream error code", async () => { }); test("closed is stable, so racing it per frame does not allocate", async () => { - const reader = new Reader(new ReadableStream()); + const reader = new Reader(new ReadableStream(), undefined, QUIC); expect(reader.closed).toBe(reader.closed); }); @@ -592,7 +599,7 @@ function stalledBidiTransport() { test("open gives up on a peer that never frees a bidi slot", async () => { const { quic, freeSlot, discarded } = stalledBidiTransport(); - await expect(Stream.open(quic, { timeout: EXPIRES_MS })).rejects.toThrow(/timed out/); + await expect(Stream.open(quic, { timeout: EXPIRES_MS, version: QUIC })).rejects.toThrow(/timed out/); freeSlot(); await discarded; @@ -601,7 +608,7 @@ test("open gives up on a peer that never frees a bidi slot", async () => { test("open waits for a bidi slot rather than failing on a busy session", async () => { const { quic, freeSlot } = stalledBidiTransport(); - const opening = Stream.open(quic, { timeout: OUTLASTS_TEST_MS }); + const opening = Stream.open(quic, { timeout: OUTLASTS_TEST_MS, version: QUIC }); expect(await Promise.race([opening.then(() => "opened"), Promise.resolve("waiting")])).toBe("waiting"); freeSlot(); @@ -614,7 +621,7 @@ test("tryOpen gives up on a cancel that settled before the open", async () => { const { quic, freeSlot, aborted } = stalledTransport(); freeSlot(); - expect(await Writer.tryOpen(quic, { cancel: Promise.resolve() })).toBeUndefined(); + expect(await Writer.tryOpen(quic, { cancel: Promise.resolve(), version: QUIC })).toBeUndefined(); await aborted; }); @@ -626,7 +633,7 @@ test("tryOpen gives up when cancelled, resetting a stream that opens afterwards" cancel = resolve; }); - const opening = Writer.tryOpen(quic, { cancel: cancelled }); + const opening = Writer.tryOpen(quic, { cancel: cancelled, version: QUIC }); cancel(); expect(await opening).toBeUndefined(); @@ -642,7 +649,7 @@ test("tryOpen rethrows a rejected cancel, resetting a stream that opens afterwar fail = reject; }); - const opening = Writer.tryOpen(quic, { cancel: cancelled }); + const opening = Writer.tryOpen(quic, { cancel: cancelled, version: QUIC }); fail(new Error("stop sending")); await expect(opening).rejects.toThrow("stop sending"); @@ -655,7 +662,9 @@ test("tryOpen rethrows a rejected cancel, resetting a stream that opens afterwar test("tryOpen gives up when the peer never frees a slot", async () => { const { quic, freeSlot, aborted } = stalledTransport(); - expect(await Writer.tryOpen(quic, { cancel: new Promise(() => {}), timeout: EXPIRES_MS })).toBeUndefined(); + expect( + await Writer.tryOpen(quic, { cancel: new Promise(() => {}), timeout: EXPIRES_MS, version: QUIC }), + ).toBeUndefined(); freeSlot(); await aborted; @@ -664,7 +673,7 @@ test("tryOpen gives up when the peer never frees a slot", async () => { test("tryOpen returns the stream when a slot is available", async () => { const { quic, freeSlot } = stalledTransport(); - const opening = Writer.tryOpen(quic, { cancel: new Promise(() => {}), timeout: OUTLASTS_TEST_MS }); + const opening = Writer.tryOpen(quic, { cancel: new Promise(() => {}), timeout: OUTLASTS_TEST_MS, version: QUIC }); freeSlot(); expect(await opening).toBeInstanceOf(Writer); @@ -680,7 +689,7 @@ test("tryOpen shares one reaction on a cancel reused across opens", async () => const cancel = new Promise(() => {}); const reactions = spyOn(cancel, "then"); for (let i = 0; i < 100; i++) { - expect(await Writer.tryOpen(quic, { cancel })).toBeInstanceOf(Writer); + expect(await Writer.tryOpen(quic, { cancel, version: QUIC })).toBeInstanceOf(Writer); } expect(reactions).toHaveBeenCalledTimes(1); }); @@ -698,10 +707,10 @@ test("open waits for a stream slot instead of rejecting once the peer's limit is }, } as unknown as WebTransport; - await Stream.open(quic, { sendOrder: 7 }); - await Writer.open(quic); + await Stream.open(quic, { sendOrder: 7, version: QUIC }); + await Writer.open(quic, { version: QUIC }); // The one path that opens faster than a peer can retire streams opts out. - await Writer.open(quic, { waitUntilAvailable: false }); + await Writer.open(quic, { waitUntilAvailable: false, version: QUIC }); expect(options).toEqual([ { sendOrder: 7, waitUntilAvailable: true }, @@ -714,7 +723,7 @@ test("open waits for a stream slot instead of rejecting once the peer's limit is // each row says what it costs on a draft that predates the registration. TOO_FAR_BEHIND // arrived in draft-17, and moq-lite's own 48-63 codes are in no draft at all. for (const [version, tooFarBehind] of [ - [undefined, StreamCode.TooFarBehind], + [QUIC, StreamCode.TooFarBehind], [Version.DRAFT_14, StreamCode.Internal], [Version.DRAFT_19, StreamCode.TooFarBehind], [Version.DRAFT_20, StreamCode.TooFarBehind], @@ -723,9 +732,9 @@ for (const [version, tooFarBehind] of [ for (const [reason, expected] of [ [new Lagged(), tooFarBehind], [new Reset(5), tooFarBehind], - [new FrameTooLarge(), version === undefined ? StreamCode.FrameTooLarge : StreamCode.Internal], - [new GroupTooLarge(), version === undefined ? StreamCode.GroupTooLarge : StreamCode.Internal], - [new NotFound("broadcast"), version === undefined ? StreamCode.NotFound : StreamCode.Internal], + [new FrameTooLarge(), version === QUIC ? StreamCode.FrameTooLarge : StreamCode.Internal], + [new GroupTooLarge(), version === QUIC ? StreamCode.GroupTooLarge : StreamCode.Internal], + [new NotFound("broadcast"), version === QUIC ? StreamCode.NotFound : StreamCode.Internal], // Assigned by every draft, so these survive the translation intact. [new TimeoutError("open"), StreamCode.DeliveryTimeout], [new ProtocolViolation("bad message"), StreamCode.SessionClosed], @@ -793,20 +802,69 @@ for (const [version, tooFarBehind] of [ ); const err = await reader.closed.catch((err: unknown) => err); expect(err).toBeInstanceOf(StreamError); - expect((err as StreamError).code).toBe(version === undefined ? StreamCode(70) : StreamCode.Internal); - if (version !== undefined) expect((err as StreamError).message).toContain("70"); + expect((err as StreamError).code).toBe(version === QUIC ? StreamCode(70) : StreamCode.Internal); + if (version !== QUIC) expect((err as StreamError).message).toContain("70"); }); } +async function written(version: Lite, f: (w: Writer) => Promise): Promise { + const { stream, written } = createTestWritableStream(); + const writer = new Writer(stream, version); + await f(writer); + writer.close(); + await writer.closed; + return [...concatChunks(written)]; +} + +test("lite-07 varints count leading ones; lite-06 keeps the QUIC form", async () => { + expect(await written(Lite.DRAFT_06, (w) => w.u53(100))).toEqual([0x40, 0x64]); + expect(await written(Lite.DRAFT_07, (w) => w.u53(100))).toEqual([0x64]); + expect(await written(Lite.DRAFT_07, (w) => w.u62(2n ** 62n - 1n))).toEqual([ + 0xff, 0x3f, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, + ]); + + expect(await new Reader(undefined, new Uint8Array([0x40, 0x64]), Lite.DRAFT_06).u53()).toBe(100); + expect(await new Reader(undefined, new Uint8Array([0x80, 0x64]), Lite.DRAFT_07).u53()).toBe(100); + // The 7-byte form, reserved only on draft-17. + const seven = new Uint8Array([0xfd, 0x23, 0x45, 0x67, 0x89, 0xab, 0xcd]); + expect(await new Reader(undefined, seven, Lite.DRAFT_07).u62()).toBe(0x1_2345_6789_abcdn); +}); + +test("the varint range follows the version: 64 bits on lite-07, 62 on lite-06", async () => { + const max = 2n ** 64n - 1n; + const wire = await written(Lite.DRAFT_07, (w) => w.u62(max)); + expect(wire).toEqual(Array(9).fill(0xff)); + expect(await new Reader(undefined, Uint8Array.from(wire), Lite.DRAFT_07).u62()).toBe(max); + await expect(written(Lite.DRAFT_06, (w) => w.u62(2n ** 62n))).rejects.toThrow("62-bits"); +}); + +test("a lite-07 stream decodes reset codes with the moq-lite registry", async () => { + // 0x4 is GOING_AWAY on moq-lite but a foreign code on draft-17. + const reset = async (version: StreamVersion) => { + const reader = new Reader( + new ReadableStream({ + start: (controller) => controller.error(new Reset(4)), + }), + undefined, + version, + ); + const err = await reader.closed.then( + () => undefined, + (e: unknown) => e, + ); + return (err as StreamError).code; + }; + expect(await reset(Lite.DRAFT_07)).toBe(StreamCode.GoingAway); + expect(await reset(Version.DRAFT_17)).toBe(StreamCode.Internal); +}); + test("Writer and Reader varint round-trip every size in both formats", async () => { const values = [0n, 63n, 64n, 127n, 128n, 2n ** 14n, 2n ** 21n, 2n ** 28n, 2n ** 30n, 2n ** 32n, 2n ** 35n]; values.push(2n ** 42n, 2n ** 49n, 2n ** 53n, 2n ** 56n, 2n ** 62n - 1n); - for (const version of [undefined, Version.DRAFT_16, Version.DRAFT_17, Version.DRAFT_19]) { + for (const version of [Lite.DRAFT_06, Lite.DRAFT_07, Version.DRAFT_16, Version.DRAFT_17, Version.DRAFT_19]) { // Only leading-ones varints reach past 62 bits. const all = - version === Version.DRAFT_17 || version === Version.DRAFT_19 - ? [...values, 2n ** 62n, 2n ** 64n - 1n] - : values; + version === Lite.DRAFT_06 || version === Version.DRAFT_16 ? values : [...values, 2n ** 62n, 2n ** 64n - 1n]; const { stream, written } = createTestWritableStream(); const writer = new Writer(stream, version); for (const value of all) await writer.varint(U64.fromBigInt(value)); @@ -821,7 +879,7 @@ test("Writer and Reader varint round-trip every size in both formats", async () test("Writer varint refuses a QUIC varint past 62 bits before emitting bytes", async () => { const { stream, written } = createTestWritableStream(); - const writer = new Writer(stream); + const writer = new Writer(stream, Version.DRAFT_16); await expect(writer.varint(U64.fromBigInt(2n ** 62n))).rejects.toThrow(/larger than 62-bits/); await expect(writer.varint(U64.MAX)).rejects.toThrow(/larger than 62-bits/); expect(written).toEqual([]); @@ -830,7 +888,7 @@ test("Writer varint refuses a QUIC varint past 62 bits before emitting bytes", a test("Reader u53 decodes every size in both formats, one byte per chunk", async () => { const values = [0, 63, 64, 127, 128, 16383, 16384, 2 ** 21, 2 ** 28 - 1, 2 ** 28, 2 ** 30 - 1, 2 ** 30]; values.push(2 ** 35, 2 ** 42, 2 ** 49, Number.MAX_SAFE_INTEGER); - for (const version of [undefined, Version.DRAFT_17]) { + for (const version of [Lite.DRAFT_06, Lite.DRAFT_07, Version.DRAFT_17]) { const { stream, written } = createTestWritableStream(); const writer = new Writer(stream, version); for (const value of values) await writer.u53(value); diff --git a/js/net/src/stream.ts b/js/net/src/stream.ts index 47e3d374d5..c2847f395a 100644 --- a/js/net/src/stream.ts +++ b/js/net/src/stream.ts @@ -2,6 +2,7 @@ import { race } from "@moq/signals"; import { fromTransport, StreamCode, StreamError, toStreamCode, toTransport } from "./error.ts"; import type { IetfVersion } from "./ietf/version.ts"; import { Version } from "./ietf/version.ts"; +import { Version as Lite, type Version as LiteVersion } from "./lite/version.ts"; import { TimeoutError, withTimeout } from "./util/timeout.ts"; import { POW32, toBigInt, toNumber, U64 } from "./util/u64.ts"; import { decodeUtf8 } from "./util/utf8.ts"; @@ -20,7 +21,8 @@ import { // Decode raw transport errors before mapping so they cannot bypass the negotiated // registry. Ordinary errors already send 0 and retain their local identity. -function withCode(reason: unknown, version?: IetfVersion): unknown { +function withCode(reason: unknown, stream: StreamVersion): unknown { + const version = asIetf(stream); const decoded = fromTransport(reason, { version }); const code = toStreamCode(decoded, { version }); return code === StreamCode.Internal && decoded === reason ? reason : toTransport(code, decoded.message); @@ -60,13 +62,59 @@ async function openWithin(opening: Promise, timeout: number, discard: (str } } -function isLeadingOnes(version?: IetfVersion): boolean { - return ( - version !== undefined && - version !== Version.DRAFT_14 && - version !== Version.DRAFT_15 && - version !== Version.DRAFT_16 - ); +/** + * The version a stream's bytes follow: a moq-transport draft or a moq-lite draft. Required on + * every stream so the varint form never defaults silently; a stream opened before negotiation + * names the version its handshake is encoded with. + */ +export type StreamVersion = IetfVersion | LiteVersion; + +const LITE: ReadonlySet = new Set(Object.values(Lite)); + +function isLite(version: StreamVersion): version is LiteVersion { + return LITE.has(version); +} + +/** The moq-transport draft a stream follows, or undefined on moq-lite, whose stream codes are the same on every draft. */ +export function asIetf(version: StreamVersion): IetfVersion | undefined { + return isLite(version) ? undefined : version; +} + +// Every draft newer than these counts leading ones, so a new version falls forward. +function isLeadingOnes(version: StreamVersion): boolean { + switch (version) { + case Version.DRAFT_14: + case Version.DRAFT_15: + case Version.DRAFT_16: + case Lite.DRAFT_01: + case Lite.DRAFT_02: + case Lite.DRAFT_03: + case Lite.DRAFT_04: + case Lite.DRAFT_05: + case Lite.DRAFT_06: + return false; + default: + return true; + } +} + +// Encode `hi`/`lo` into `dst` in the varint form `version` uses: QUIC up to 2^62-1, leading-ones up to 2^64-1. +function encodeTo(dst: ArrayBuffer, hi: number, lo: number, version: StreamVersion): Uint8Array { + let buf: Uint8Array; + if (isLeadingOnes(version)) { + buf = new Uint8Array(dst, 0, lengthLeadingOnes(hi, lo)); + writeLeadingOnes(buf, hi, lo, buf.length); + } else { + buf = new Uint8Array(dst, 0, lengthQuic(hi, lo)); + writeQuic(buf, hi, lo, buf.length); + } + return buf; +} + +/** Encode one varint in the form `version` uses, for a body written outside a {@link Writer}. */ +export function encodeVarint(v: number | bigint, version: StreamVersion): Uint8Array { + const lo = split(v); + return encodeTo(new ArrayBuffer(9), parts.hi, lo, version); } /** @@ -83,8 +131,8 @@ export type SendStream = WritableStream & { sendOrder?: number }; /** Options for opening an outgoing stream. */ export interface OpenOptions { - /** The negotiated IETF version, which selects the varint encoding. */ - version?: IetfVersion; + /** The negotiated version, which selects the varint encoding. */ + version: StreamVersion; /** * The transport send order, where HIGHER values are transmitted first. @@ -125,7 +173,7 @@ export class Stream { constructor(props: { writable: WritableStream; readable: ReadableStream; - version?: IetfVersion; + version: StreamVersion; }); /** Pair halves that were opened separately, as the SETUP exchange does. */ constructor(props: { writer: Writer; reader: Reader }); @@ -134,17 +182,19 @@ export class Stream { readable?: ReadableStream; writer?: Writer; reader?: Reader; - version?: IetfVersion; + version?: StreamVersion; }) { - const writer = props.writer ?? (props.writable && new Writer(props.writable, props.version)); - const reader = props.reader ?? (props.readable && new Reader(props.readable, undefined, props.version)); + const version = props.version; + const writer = props.writer ?? (props.writable && version !== undefined && new Writer(props.writable, version)); + const reader = + props.reader ?? (props.readable && version !== undefined && new Reader(props.readable, undefined, version)); if (!writer || !reader) throw new Error("stream needs both halves"); this.writer = writer; this.reader = reader; } - static async accept(quic: WebTransport, version?: IetfVersion): Promise { + static async accept(quic: WebTransport, version: StreamVersion): Promise { for (;;) { const reader = quic.incomingBidirectionalStreams.getReader() as ReadableStreamDefaultReader; @@ -163,7 +213,7 @@ export class Stream { * @param options - The version its varints encode with, and the send order ranking it * against the session's other streams */ - static async open(quic: WebTransport, options?: OpenOptions): Promise { + static async open(quic: WebTransport, options: OpenOptions): Promise { const { readable, writable } = await openWithin( quic.createBidirectionalStream(sendOptions(options)), options?.timeout ?? OPEN_TIMEOUT_MS, @@ -172,7 +222,7 @@ export class Stream { void stream.readable.cancel().catch(() => void 0); }, ); - return new Stream({ readable, writable, version: options?.version }); + return new Stream({ readable, writable, version: options.version }); } close() { @@ -202,12 +252,16 @@ export class Reader { #closed?: Promise; // The decode that last ran short and how far, so a retry can wait for those bytes. #short?: { decode: (c: Cursor) => unknown; err: Short }; - version?: IetfVersion; + version: StreamVersion; // Either stream or buffer MUST be provided. - constructor(stream: ReadableStream, buffer?: Uint8Array, version?: IetfVersion); - constructor(stream: undefined, buffer: Uint8Array, version?: IetfVersion); - constructor(stream?: ReadableStream, buffer?: Uint8Array, version?: IetfVersion) { + constructor(stream: ReadableStream, buffer: Uint8Array | undefined, version: StreamVersion); + constructor(stream: undefined, buffer: Uint8Array, version: StreamVersion); + constructor( + stream: ReadableStream | undefined, + buffer: Uint8Array | undefined, + version: StreamVersion, + ) { this.#buffer = buffer ?? new Uint8Array(); this.#stream = stream; this.#reader = this.#stream?.getReader(); @@ -223,7 +277,7 @@ export class Reader { // Every read of this stream funnels through here, so decoding the peer's reset code // once is enough to keep the raw transport error out of every caller (and every app). const result = await this.#reader.read().catch((err: unknown) => { - throw fromTransport(err, { version: this.version }); + throw fromTransport(err, { version: asIetf(this.version) }); }); if (result.done) { @@ -397,7 +451,7 @@ export class Reader { // shape depending on which one won. Derived once, so racing it per frame doesn't allocate. get closed(): Promise { this.#closed ??= (this.#reader?.closed ?? Promise.resolve()).catch((err: unknown) => { - throw fromTransport(err, { version: this.version }); + throw fromTransport(err, { version: asIetf(this.version) }); }); return this.#closed; } @@ -424,7 +478,7 @@ const EMPTY = new Short(1); * anything before its last read, must not swallow what it throws, and must read at least a byte. */ export class Cursor { - readonly version?: IetfVersion; + readonly version: StreamVersion; #buffer: Uint8Array; #offset = 0; // Resolved once, since every varint read branches on it. @@ -435,7 +489,7 @@ export class Cursor { // First bytes below this are a varint of at most 4 bytes: 0xc0 for QUIC, 0xf0 for leading-ones. #word: number; - constructor(buffer: Uint8Array, version?: IetfVersion) { + constructor(buffer: Uint8Array, version: StreamVersion) { this.#buffer = buffer; this.version = version; this.#leadingOnes = isLeadingOnes(version); @@ -591,9 +645,9 @@ export class Writer { // Scratch buffer for each primitive write, sized for the longest (a 9-byte leading-ones varint). #scratch: ArrayBuffer; - version?: IetfVersion; + version: StreamVersion; - constructor(stream: WritableStream, version?: IetfVersion) { + constructor(stream: WritableStream, version: StreamVersion) { this.#stream = stream; this.#scratch = new ArrayBuffer(9); this.#writer = this.#stream.getWriter(); @@ -656,22 +710,14 @@ export class Writer { } #varint(hi: number, lo: number): Promise { - let buf: Uint8Array; - if (isLeadingOnes(this.version)) { - buf = new Uint8Array(this.#scratch, 0, lengthLeadingOnes(hi, lo)); - writeLeadingOnes(buf, hi, lo, buf.length); - } else { - buf = new Uint8Array(this.#scratch, 0, lengthQuic(hi, lo)); - writeQuic(buf, hi, lo, buf.length); - } - return this.write(buf); + return this.write(encodeTo(this.#scratch, hi, lo, this.version)); } async write(v: Uint8Array) { // Mirrors Reader.#fill: every write funnels through here, so a STOP_SENDING from the // peer surfaces as a typed code rather than the transport's own error shape. await this.#writer.write(v).catch((err: unknown) => { - throw fromTransport(err, { version: this.version }); + throw fromTransport(err, { version: asIetf(this.version) }); }); } @@ -689,7 +735,7 @@ export class Writer { // typed code it would get from a write. get closed(): Promise { this.#closed ??= this.#writer.closed.catch((err: unknown) => { - throw fromTransport(err, { version: this.version }); + throw fromTransport(err, { version: asIetf(this.version) }); }); return this.#closed; } @@ -704,14 +750,14 @@ export class Writer { * @param options - The version its varints encode with, and the send order ranking it * against the session's other streams */ - static async open(quic: WebTransport, options?: OpenOptions): Promise { + static async open(quic: WebTransport, options: OpenOptions): Promise { const writable = await openWithin( quic.createUnidirectionalStream(sendOptions(options)) as Promise>, options?.timeout ?? OPEN_TIMEOUT_MS, (stream) => void stream.abort().catch(() => void 0), ); - return new Writer(writable, options?.version); + return new Writer(writable, options.version); } /** @@ -771,9 +817,9 @@ function setInt32(dst: ArrayBuffer, v: number): Uint8Array { // Returns the next stream from the connection export class Readers { #reader: ReadableStreamDefaultReader>; - #version?: IetfVersion; + #version: StreamVersion; - constructor(quic: WebTransport, version?: IetfVersion) { + constructor(quic: WebTransport, version: StreamVersion) { this.#reader = quic.incomingUnidirectionalStreams.getReader() as ReadableStreamDefaultReader< ReadableStream >; diff --git a/quest/m1/rs2ts/varint-codec.md b/quest/m1/rs2ts/varint-codec.md index 39557cb31b..cda246b8ba 100644 --- a/quest/m1/rs2ts/varint-codec.md +++ b/quest/m1/rs2ts/varint-codec.md @@ -17,7 +17,12 @@ WASM against 6 KB for a hand-carved one. Guidance: -- Keep the 62-bit range; the wire is not bounded to 2^53. +- Widen `VarInt` to the full 64 bits that leading-ones carries on moq-lite-07 + and moq-transport draft-17+, and never bound it to 2^53. QUIC varints still + refuse to encode past 2^62-1. On lite-07 this lets Rust delete the + `BoundsExceeded` arm in the lite dispatch and flips `lite_varint_interop` from + refusal to a round trip. Also revisit `MAX_COST`: costs saturate at 2^62-1 on + every version, but the draft caps them at the largest value a varint carries. - Message types keep a local trait; the primitives become inherent methods on concrete reader and writer types (`varint`, `string`, `bytes`, ...). Make the version a concrete type rather than a generic `V` where possible. diff --git a/rs/moq-net/Cargo.toml b/rs/moq-net/Cargo.toml index 967afb259d..e4c2c85c5f 100644 --- a/rs/moq-net/Cargo.toml +++ b/rs/moq-net/Cargo.toml @@ -16,7 +16,7 @@ categories = ["multimedia", "network-programming", "web-programming"] ignored = ["getrandom"] [features] -# Exposes the wire codecs to the `fuzz/` harness and the announce bench through the +# Exposes the wire codecs to the `fuzz/` harness and the announce and varint benches through the # hidden `fuzz` module. Off by default and hidden from the docs. fuzz = [] @@ -75,6 +75,12 @@ name = "announce" harness = false required-features = ["fuzz"] +# Reaches the private lite codec through the hidden `fuzz` module. +[[bench]] +name = "varint" +harness = false +required-features = ["fuzz"] + [[bench]] name = "group" harness = false diff --git a/rs/moq-net/benches/varint.rs b/rs/moq-net/benches/varint.rs new file mode 100644 index 0000000000..01ce78a397 --- /dev/null +++ b/rs/moq-net/benches/varint.rs @@ -0,0 +1,43 @@ +//! Lite-07 leading-ones varints against lite-06 QUIC varints: the bytes and the CPU to +//! encode and decode what a publisher writes per frame, per group, and per request. +//! +//! The two versions share these layouts, so the only difference is the varint codec. +//! The byte counts print once before the timings. +//! +//! Run with `cargo bench -p moq-net --features fuzz --bench varint`. + +use std::hint::black_box; + +use criterion::{BenchmarkId, Criterion, criterion_group, criterion_main}; +use moq_net::{Version, fuzz::LiteSample}; + +fn versions() -> [(&'static str, Version); 2] { + ["moq-lite-06", "moq-lite-07-wip"].map(|name| (name, name.parse().unwrap())) +} + +fn bench(c: &mut Criterion) { + println!("{:<10} {:>8} {:>8}", "sample", "lite-06", "lite-07"); + for sample in LiteSample::ALL { + let [lite06, lite07] = versions().map(|(_, version)| sample.encode(version).len()); + println!("{:<10} {lite06:>8} {lite07:>8}", format!("{sample:?}")); + } + + let mut group = c.benchmark_group("lite_varint"); + for sample in LiteSample::ALL { + for (name, version) in versions() { + let id = format!("{sample:?}/{name}"); + group.bench_function(BenchmarkId::new("encode", &id), |b| { + b.iter(|| black_box(sample.encode(black_box(version)))) + }); + + let wire = sample.encode(version); + group.bench_function(BenchmarkId::new("decode", &id), |b| { + b.iter(|| black_box(sample.decode(black_box(version), black_box(&wire)))) + }); + } + } + group.finish(); +} + +criterion_group!(benches, bench); +criterion_main!(benches); diff --git a/rs/moq-net/src/coding/varint.rs b/rs/moq-net/src/coding/varint.rs index e51180d531..08e35cab0f 100644 --- a/rs/moq-net/src/coding/varint.rs +++ b/rs/moq-net/src/coding/varint.rs @@ -270,7 +270,7 @@ impl VarInt { /// - `1111110x` → 7 bytes, 49 usable bits (draft-18+, INVALID in draft-17 per #1595) /// - `11111110` → 8 bytes, 56 usable bits /// - `11111111` → 9 bytes, 64 usable bits - fn decode_leading_ones(r: &mut R, version: ietf::Version) -> Result { + fn decode_leading_ones(r: &mut R) -> Result { if !r.has_remaining() { return Err(DecodeError::Short); } @@ -341,9 +341,6 @@ impl VarInt { } 6 => { // 1111110x + 6 bytes, 49 bits (draft-18+, INVALID in draft-17 per #1595) - if matches!(version, ietf::Version::Draft17) { - return Err(DecodeError::InvalidValue); - } if r.remaining() < 6 { return Err(DecodeError::Short); } @@ -380,7 +377,7 @@ impl VarInt { /// Always emits the minimal canonical form. Draft-18 also accepts 7-byte form /// (`1111110x`) on decode but we never emit it because the 8-byte form is one byte /// larger but simpler and is universally valid. - fn encode_leading_ones(&self, w: &mut W, _version: ietf::Version) -> Result<(), EncodeError> { + fn encode_leading_ones(&self, w: &mut W) -> Result<(), EncodeError> { let x = self.0; let remaining = w.remaining_mut(); @@ -452,16 +449,39 @@ impl VarInt { use crate::{Version, ietf, lite}; -// All lite versions use QUIC-style varint encoding. +// Lite01-06 use QUIC-style varints; lite-07+ uses leading-ones. Lite-07 is only reached +// through its ALPN, so the codec is known before the first byte of any stream. impl Encode for VarInt { - fn encode(&self, w: &mut W, _: lite::Version) -> Result<(), EncodeError> { - self.encode_quic(w) + fn encode(&self, w: &mut W, version: lite::Version) -> Result<(), EncodeError> { + match version { + lite::Version::Lite01 + | lite::Version::Lite02 + | lite::Version::Lite03 + | lite::Version::Lite04 + | lite::Version::Lite05 + | lite::Version::Lite06 => self.encode_quic(w), + _ => self.encode_leading_ones(w), + } } } impl Decode for VarInt { - fn decode(r: &mut R, _: lite::Version) -> Result { - Self::decode_quic(r) + fn decode(r: &mut R, version: lite::Version) -> Result { + match version { + lite::Version::Lite01 + | lite::Version::Lite02 + | lite::Version::Lite03 + | lite::Version::Lite04 + | lite::Version::Lite05 + | lite::Version::Lite06 => Self::decode_quic(r), + // Lite-07 values span the full 64 bits, but `VarInt` is still 62-bit here, so a + // larger value is a loud decode error, never a truncation. Known limitation, lifted + // when the VarInt codec quest (quest/m1/rs2ts/varint-codec.md) widens `VarInt`. + _ => match Self::decode_leading_ones(r)? { + x if x > Self::MAX => Err(DecodeError::BoundsExceeded), + x => Ok(x), + }, + } } } @@ -470,7 +490,7 @@ impl Encode for VarInt { fn encode(&self, w: &mut W, version: ietf::Version) -> Result<(), EncodeError> { match version { ietf::Version::Draft14 | ietf::Version::Draft15 | ietf::Version::Draft16 => self.encode_quic(w), - _ => self.encode_leading_ones(w, version), + _ => self.encode_leading_ones(w), } } } @@ -479,7 +499,11 @@ impl Decode for VarInt { fn decode(r: &mut R, version: ietf::Version) -> Result { match version { ietf::Version::Draft14 | ietf::Version::Draft15 | ietf::Version::Draft16 => Self::decode_quic(r), - _ => Self::decode_leading_ones(r, version), + // Draft-18 made 1111110x the 7-byte form; draft-17 reserves it (#1595). + ietf::Version::Draft17 if r.chunk().first().is_some_and(|b| b.leading_ones() == 6) => { + Err(DecodeError::InvalidValue) + } + _ => Self::decode_leading_ones(r), } } } @@ -592,7 +616,7 @@ mod tests { for (bytes, expected) in cases { // Test decoding let mut buf = Bytes::from(bytes.to_vec()); - let decoded = VarInt::decode_leading_ones(&mut buf, ietf::Version::Draft17).expect("decode should succeed"); + let decoded = VarInt::decode_leading_ones(&mut buf).expect("decode should succeed"); assert_eq!( decoded.into_inner(), *expected, @@ -607,9 +631,7 @@ mod tests { && (bytes.len() == 1 || *expected != 37) { let mut encoded = Vec::new(); - varint - .encode_leading_ones(&mut encoded, ietf::Version::Draft17) - .expect("encode should succeed"); + varint.encode_leading_ones(&mut encoded).expect("encode should succeed"); assert_eq!(&encoded, bytes, "encode mismatch for value {expected}"); } } @@ -621,7 +643,7 @@ mod tests { let mut buf = Bytes::from_static(&[0xFC]); assert!( matches!( - VarInt::decode_leading_ones(&mut buf, ietf::Version::Draft17), + VarInt::decode(&mut buf, ietf::Version::Draft17), Err(DecodeError::InvalidValue) ), "0xFC should be rejected as invalid on draft-17" @@ -643,7 +665,7 @@ mod tests { let varint = VarInt::from_u64(value).expect("value should be representable as VarInt"); let mut encoded = Vec::new(); varint - .encode_leading_ones(&mut encoded, ietf::Version::Draft17) + .encode_leading_ones(&mut encoded) .expect("leading-ones encode should succeed"); assert_eq!( encoded.len(), @@ -652,8 +674,7 @@ mod tests { ); let mut bytes = Bytes::from(encoded); - let decoded = VarInt::decode_leading_ones(&mut bytes, ietf::Version::Draft17) - .expect("leading-ones decode should succeed"); + let decoded = VarInt::decode_leading_ones(&mut bytes).expect("leading-ones decode should succeed"); assert_eq!(decoded.into_inner(), value, "round-trip mismatch for value {value}"); } } @@ -663,7 +684,7 @@ mod tests { // 1111110x prefix: invalid on draft-17. let bytes = Bytes::from(vec![0xFC, 0, 0, 0, 0, 0, 0]); let mut buf = bytes.clone(); - let err = VarInt::decode_leading_ones(&mut buf, ietf::Version::Draft17).unwrap_err(); + let err = VarInt::decode(&mut buf, ietf::Version::Draft17).unwrap_err(); assert!(matches!(err, DecodeError::InvalidValue)); } @@ -721,6 +742,67 @@ mod tests { } } + fn lite(value: u64, version: lite::Version) -> Vec { + let mut buf = Vec::new(); + VarInt::try_from(value).unwrap().encode(&mut buf, version).unwrap(); + buf + } + + /// Lite-01 through lite-06 keep the QUIC form byte for byte. + #[test] + fn lite06_keeps_quic_varints() { + for version in [lite::Version::Lite01, lite::Version::Lite05, lite::Version::Lite06] { + assert_eq!(lite(63, version), [0x3F]); + assert_eq!(lite(64, version), [0x40, 0x40]); + assert_eq!(lite(16_384, version), [0x80, 0x00, 0x40, 0x00]); + } + } + + /// Lite-07 switches to leading-ones, including at every length boundary. + #[test] + fn lite07_uses_leading_ones() { + let version = lite::Version::Lite07; + let cases: &[(u64, &[u8])] = &[ + (0, &[0x00]), + (127, &[0x7F]), + (128, &[0x80, 0x80]), + ((1 << 14) - 1, &[0xBF, 0xFF]), + (1 << 14, &[0xC0, 0x40, 0x00]), + ((1 << 21) - 1, &[0xDF, 0xFF, 0xFF]), + (1 << 21, &[0xE0, 0x20, 0x00, 0x00]), + ((1 << 56) - 1, &[0xFE, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF]), + (1 << 56, &[0xFF, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00]), + ( + VarInt::MAX.into_inner(), + &[0xFF, 0x3F, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF], + ), + ]; + for (value, wire) in cases { + assert_eq!(lite(*value, version), *wire, "encode {value}"); + let mut buf = *wire; + assert_eq!(VarInt::decode(&mut buf, version).unwrap().into_inner(), *value); + assert!(buf.is_empty()); + } + } + + /// The 7-byte `1111110x` form is valid on lite-07, as on draft-18+. + #[test] + fn lite07_accepts_7_byte_varint() { + let mut buf: &[u8] = &[0xFD, 0x23, 0x45, 0x67, 0x89, 0xAB, 0xCD]; + let decoded = VarInt::decode(&mut buf, lite::Version::Lite07).unwrap(); + assert_eq!(decoded.into_inner(), 0x1_2345_6789_ABCD); + } + + /// Lite-07 values above 2^62-1 are legal on the wire, but the 62-bit `VarInt` cannot + /// hold them yet, so they must fail loud rather than wrap or truncate. + #[test] + fn lite07_values_above_62_bits_fail_loud() { + for wire in [[0xFF, 0x40, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00], [0xFF; 9]] { + let err = VarInt::decode(&mut &wire[..], lite::Version::Lite07).unwrap_err(); + assert!(matches!(err, DecodeError::BoundsExceeded), "{wire:02x?}: {err:?}"); + } + } + #[test] fn draft18_accepts_7_byte_varint() { // Value 0x1234_5678_9ABC encoded as 7-byte leading-ones (1111110x | hi, +6 bytes). @@ -734,7 +816,7 @@ mod tests { bytes.push(((value >> shift) & 0xFF) as u8); } let mut buf = Bytes::from(bytes); - let decoded = VarInt::decode_leading_ones(&mut buf, ietf::Version::Draft18).unwrap(); + let decoded = VarInt::decode(&mut buf, ietf::Version::Draft18).unwrap(); assert_eq!(decoded.into_inner(), value); } } diff --git a/rs/moq-net/src/fuzz.rs b/rs/moq-net/src/fuzz.rs index 5fc9f793b6..a7dd026195 100644 --- a/rs/moq-net/src/fuzz.rs +++ b/rs/moq-net/src/fuzz.rs @@ -252,6 +252,124 @@ pub fn decode_announces(mut data: &[u8], compress: bool) -> Vec { resolved } +/// A varint-heavy lite wire object, as the varint bench times it: what a publisher +/// pays per frame, per group, and per request. +#[derive(Clone, Copy, Debug)] +pub enum LiteSample { + /// A 30 fps video group's FRAME headers: a zigzag timestamp delta in microseconds and + /// a size, 60 of them, led by a 60 KB keyframe. Payloads are left out. + Video, + /// A 50 Hz Opus group's FRAME headers: 50 of 20 ms and 160 bytes. + Audio, + /// A GROUP header, well into a long-running track. + Group, + /// A SUBSCRIBE for a track, the way a player opens one. + Subscribe, + /// A datagram's header, with an empty payload. + Datagram, + /// A SETUP declaring a probe level, a cost, and a random Hop ID. + Setup, +} + +impl LiteSample { + /// Every sample, in the order the bench reports them. + pub const ALL: [Self; 6] = [ + Self::Video, + Self::Audio, + Self::Group, + Self::Subscribe, + Self::Datagram, + Self::Setup, + ]; + + /// The FRAME headers for [`Self::Video`] and [`Self::Audio`], as (delta, size) pairs. + fn frames(self) -> Vec<(i64, u64)> { + match self { + Self::Video => (0..60) + .map(|n| match n { + 0 => (0, 60_000), + n if n % 10 == 0 => (33_333, 17_000), + _ => (33_333, 8_000), + }) + .collect(), + Self::Audio => (0..50).map(|n| (if n == 0 { 0 } else { 20_000 }, 160)).collect(), + _ => Vec::new(), + } + } + + /// Encode this sample at `version`, which must be a moq-lite version. + pub fn encode(self, version: crate::Version) -> Vec { + let version = lite::Version::try_from(version).expect("a moq-lite version"); + let mut buf = Vec::new(); + match self { + Self::Video | Self::Audio => { + for (delta, size) in self.frames() { + VarInt::from_zigzag(delta).unwrap().encode(&mut buf, version).unwrap(); + size.encode(&mut buf, version).unwrap(); + } + } + Self::Group => lite::Group { + subscribe: 3, + sequence: 1_234, + frame_start: 0, + } + .encode(&mut buf, version) + .unwrap(), + Self::Subscribe => lite::Subscribe { + id: 3, + broadcast: Path::new("room/alice"), + track: "video".into(), + priority: 2, + max_age: std::time::Duration::from_secs(10), + start_group: None, + end_group: None, + start_frame: 0, + end_frame: None, + } + .encode(&mut buf, version) + .unwrap(), + Self::Datagram => lite::Datagram { + subscribe: 3, + sequence: 1_234, + timestamp: 1_234_567_890, + payload: bytes::Bytes::new(), + } + .encode(&mut buf, version) + .unwrap(), + Self::Setup => lite::Setup { + probe: lite::ProbeLevel::Report, + cost: Some(1), + hop: Some(crate::Hop::new(0x1d_2c3b_4a59_6877).unwrap()), + ..Default::default() + } + .encode(&mut buf, version) + .unwrap(), + } + buf + } + + /// Decode what [`Self::encode`] wrote at `version`, returning how many objects it read. + pub fn decode(self, version: crate::Version, mut data: &[u8]) -> usize { + let version = lite::Version::try_from(version).expect("a moq-lite version"); + let data = &mut data; + match self { + Self::Video | Self::Audio => { + let mut frames = 0; + while !data.is_empty() { + VarInt::decode(data, version).unwrap(); + u64::decode(data, version).unwrap(); + frames += 1; + } + frames + } + Self::Group => lite::Group::decode(data, version).map(|_| 1).unwrap(), + Self::Subscribe => lite::Subscribe::decode(data, version).map(|_| 1).unwrap(), + Self::Datagram => lite::Datagram::decode(data, version).map(|_| 1).unwrap(), + Self::Setup => lite::Setup::decode(data, version).map(|_| 1).unwrap(), + } + } +} + /// Feed a lite-07 announce stream through the stateful decoder, then check that our /// encoder's compression of what it resolved reads back as the same announcements. /// @@ -324,13 +442,14 @@ pub fn ietf_wire(data: &[u8]) -> bool { /// Decode a varint with whichever codec the selected version uses. /// -/// Byte 0 picks the version, which is the whole point: moq-lite and drafts 14-16 use -/// the QUIC two-bit length tag, while draft-17+ counts leading ones, and the two -/// disagree about which byte sequences are even legal. +/// Byte 0 picks the version, which is the whole point: lite-01 to lite-06 and drafts +/// 14-16 use the QUIC two-bit length tag, while lite-07 and draft-17+ count leading +/// ones, and the two disagree about which byte sequences are even legal. /// -/// The decoded value is deliberately not asserted to be within [`VarInt::MAX`]: the -/// leading-ones form spans the full `u64` by design, so a 9-byte encoding decodes -/// above the 62-bit ceiling and only fails when re-encoded for a QUIC-form version. +/// The decoded value is deliberately not asserted to be within [`VarInt::MAX`]: on the +/// IETF wire the leading-ones form spans the full `u64` by design, so a 9-byte encoding +/// decodes above the 62-bit ceiling. Lite-07 allows the same range, but refuses it at +/// decode until `VarInt` widens to 64 bits. pub fn varint(data: &[u8]) -> bool { let Some((&selector, rest)) = data.split_first() else { return false; diff --git a/rs/moq-net/src/lite/announce.rs b/rs/moq-net/src/lite/announce.rs index 9e460d1530..415fe598d3 100644 --- a/rs/moq-net/src/lite/announce.rs +++ b/rs/moq-net/src/lite/announce.rs @@ -201,10 +201,13 @@ impl Decode for Cost { if !version.has_route_cost() { return Ok(Cost::UNKNOWN); } - Ok(Cost { + // Costs saturate at 2^62-1 on every version, so a larger one (lite-07's varints + // reach 2^64-1) reads as the ceiling and still forwards to an older peer. + let cost = Cost { warm: u64::decode(buf, version)?, cold: u64::decode(buf, version)?, - }) + }; + Ok(cost.clamped()) } } @@ -783,15 +786,19 @@ mod tests { ); } - // A peer may legally advertise the largest varint there is, and adding this - // link's price to it must not push the result out of range. + // Costs saturate at 2^62-1 on every version, lite-07's 64-bit varints included, so + // charging a link on top of the ceiling still re-encodes for a peer on any version. #[test] fn charged_cost_stays_encodable() { - let mut buf = Vec::new(); - crate::origin::Cost::MAX - .charged(1) - .encode(&mut buf, Version::Lite06) - .expect("a charged cost must stay encodable"); + let cost = Cost::new(u64::MAX).charged(1); + assert_eq!(cost, Cost::new((1 << 62) - 1)); + assert_eq!(Cost::MAX.charged(1), cost); + for version in [Version::Lite06, Version::Lite07] { + let mut buf = Vec::new(); + cost.encode(&mut buf, version) + .expect("a charged cost must stay encodable"); + assert_eq!(Cost::decode(&mut &buf[..], version).unwrap(), cost, "{version}"); + } } #[test] diff --git a/rs/moq-net/src/lite/compress.rs b/rs/moq-net/src/lite/compress.rs index e908e8cbec..eda17b3916 100644 --- a/rs/moq-net/src/lite/compress.rs +++ b/rs/moq-net/src/lite/compress.rs @@ -581,7 +581,7 @@ mod tests { /// A compressed stream pinned byte for byte: `js/net/src/lite/announce.test.ts` /// decodes the same hex, so a codec change on either side breaks both. - const GOLDEN: &str = "001600000a726f6f6d2f612f63616d000251116222000000000d0102036d696301017333010000020a00010180004444010000010101000d02010162020180005555010000"; + const GOLDEN: &str = "001600000a726f6f6d2f612f63616d00029111a222000000000d0102036d69630101b3330100000209000101c04444010000010101000c020101620201c05555010000"; fn golden() -> Vec { use crate::fuzz::Announced; @@ -604,7 +604,7 @@ mod tests { /// The literal lite-07 stream `js/net` writes for the same announcements, which /// never picks a base. - const JS_LITERAL: &str = "001600000a726f6f6d2f612f63616d000251116222000000001600000a726f6f6d2f612f6d6963000273336222000000020c0000028000444462220000000101010014000006726f6f6d2f620002800055556222000000"; + const JS_LITERAL: &str = "001600000a726f6f6d2f612f63616d00029111a222000000001600000a726f6f6d2f612f6d69630002b333a222000000020b000002c04444a2220000000101010013000006726f6f6d2f620002c05555a222000000"; #[test] fn js_literal_stream_decodes() { diff --git a/rs/moq-net/src/lite/parameters.rs b/rs/moq-net/src/lite/parameters.rs index b0148a6aae..b7930a6c74 100644 --- a/rs/moq-net/src/lite/parameters.rs +++ b/rs/moq-net/src/lite/parameters.rs @@ -23,20 +23,19 @@ impl Parameters { self.0.get(&id).map(Vec::as_slice) } - /// Set a parameter to a varint value, replacing any existing entry. - pub fn set_varint(&mut self, id: u64, value: u64) { - let mut buf = Vec::new(); - // Infallible: writing into a Vec never runs short. - value.encode(&mut buf, Version::Lite05).expect("varint encode into Vec"); - self.0.insert(id, buf); + /// Set a parameter to a varint value in `version`'s encoding, replacing any existing entry. + pub fn set_varint(&mut self, id: u64, value: u64, version: Version) -> Result<(), EncodeError> { + self.0.insert(id, value.encode_bytes(version)?.to_vec()); + Ok(()) } - /// Decode a parameter as a single varint, if present. Errors if trailing bytes remain. - pub fn get_varint(&self, id: u64) -> Result, DecodeError> { + /// Decode a parameter as a single varint in `version`'s encoding, if present. Errors if + /// trailing bytes remain. + pub fn get_varint(&self, id: u64, version: Version) -> Result, DecodeError> { let Some(mut bytes) = self.0.get(&id).map(Vec::as_slice) else { return Ok(None); }; - let value = u64::decode(&mut bytes, Version::Lite05)?; + let value = u64::decode(&mut bytes, version)?; if !bytes.is_empty() { return Err(DecodeError::Long); } diff --git a/rs/moq-net/src/lite/setup.rs b/rs/moq-net/src/lite/setup.rs index c3868614e9..53810df047 100644 --- a/rs/moq-net/src/lite/setup.rs +++ b/rs/moq-net/src/lite/setup.rs @@ -194,7 +194,7 @@ impl Message for Setup { let params = Parameters::decode(r, version)?; let probe = params - .get_varint(PARAM_PROBE)? + .get_varint(PARAM_PROBE, version)? .map(ProbeLevel::from_code) .unwrap_or_default(); let path = match params.get_bytes(PARAM_PATH) { @@ -205,11 +205,13 @@ impl Message for Setup { ), None => None, }; - let role = params.get_varint(PARAM_ROLE)?.and_then(Role::from_code); - let cost = params.get_varint(PARAM_COST)?; + let role = params.get_varint(PARAM_ROLE, version)?.and_then(Role::from_code); + let cost = params.get_varint(PARAM_COST, version)?; // 0 is legal on the wire but carries no identity (it can't be excluded), // so it decodes as "not declared" rather than an error. - let hop = params.get_varint(PARAM_HOP)?.and_then(|id| crate::Hop::new(id).ok()); + let hop = params + .get_varint(PARAM_HOP, version)? + .and_then(|id| crate::Hop::new(id).ok()); Ok(Self { probe, @@ -228,7 +230,7 @@ impl Message for Setup { let mut params = Parameters::default(); // None is the wire default, so omit it to keep the message empty when nothing is set. if self.probe != ProbeLevel::None { - params.set_varint(PARAM_PROBE, self.probe.to_code()); + params.set_varint(PARAM_PROBE, self.probe.to_code(), version)?; } if let Some(path) = &self.path { params.set_bytes(PARAM_PATH, path.as_bytes().to_vec()); @@ -236,13 +238,13 @@ impl Message for Setup { // Bidirectional is the wire default (absence of the parameter), so only a // directional role is encoded. if let Some(role) = self.role { - params.set_varint(PARAM_ROLE, role.to_code()); + params.set_varint(PARAM_ROLE, role.to_code(), version)?; } if let Some(cost) = self.cost { - params.set_varint(PARAM_COST, cost); + params.set_varint(PARAM_COST, cost, version)?; } if let Some(hop) = self.hop { - params.set_varint(PARAM_HOP, hop.id()); + params.set_varint(PARAM_HOP, hop.id(), version)?; } params.encode(w, version) @@ -365,6 +367,25 @@ mod tests { } } + /// Parameter values follow the session's varint codec, not a fixed one: a value that + /// takes the two-byte QUIC form on lite-06 fits one leading-ones byte on lite-07. + #[test] + fn parameter_values_use_the_version_codec() { + let msg = Setup { + cost: Some(100), + ..Default::default() + }; + for (version, wire) in [ + (Version::Lite06, &[0x05, 0x01, 0x04, 0x02, 0x40, 0x64][..]), + (Version::Lite07, &[0x04, 0x01, 0x04, 0x01, 0x64][..]), + ] { + let mut buf = Vec::new(); + msg.encode(&mut buf, version).unwrap(); + assert_eq!(buf, wire, "{version}"); + assert_eq!(Setup::decode(&mut &buf[..], version).unwrap(), msg, "{version}"); + } + } + /// Never emit a SETUP our own receiver would refuse. #[test] fn encode_enforces_the_setup_limit() { @@ -415,7 +436,7 @@ mod tests { let version = Version::Lite05; let mut params = Parameters::default(); - params.set_varint(super::PARAM_HOP, 0); + params.set_varint(super::PARAM_HOP, 0, version).unwrap(); let mut body = bytes::BytesMut::new(); params.encode(&mut body, version).unwrap(); // Frame the body with the Message Length prefix `Setup::decode` expects. @@ -455,7 +476,7 @@ mod tests { // Frame a SETUP message carrying an unknown probe level (99) by hand: the // parameters body, prefixed with its length (the lite Message size prefix). let mut params = Parameters::default(); - params.set_varint(PARAM_PROBE, 99); + params.set_varint(PARAM_PROBE, 99, Version::Lite05).unwrap(); let mut body = Vec::new(); params.encode(&mut body, Version::Lite05).unwrap(); @@ -485,7 +506,7 @@ mod tests { // server. The draft mandates this fallback. for code in [0u64, 9, 250] { let mut params = Parameters::default(); - params.set_varint(PARAM_ROLE, code); + params.set_varint(PARAM_ROLE, code, Version::Lite05).unwrap(); let mut body = Vec::new(); params.encode(&mut body, Version::Lite05).unwrap(); diff --git a/rs/moq-net/src/model/origin.rs b/rs/moq-net/src/model/origin.rs index 1aef8cf8ba..2758b39347 100644 --- a/rs/moq-net/src/model/origin.rs +++ b/rs/moq-net/src/model/origin.rs @@ -36,7 +36,7 @@ use crate::{ /// on the wire, names nobody, and marks the chain anonymous for route selection. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub struct Hop { - /// 62-bit identifier. Encoded as a QUIC varint on the wire. + /// 62-bit identifier, so it fits a varint on every wire version. id: u64, } @@ -384,7 +384,8 @@ where /// /// The ceiling is the wire's, not the model's: lite-06 carries each cost as a QUIC /// varint, which tops out at 2^62-1, so a larger value could be selected on but -/// never forwarded. +/// never forwarded. It applies on every version, lite-07's 64-bit varints included, +/// so a cost stays forwardable to a QUIC-varint version. const MAX_COST: u64 = (1 << 62) - 1; /// What pulling content via a route costs, in two magnitudes that accumulate diff --git a/rs/moq-net/src/test_interop.rs b/rs/moq-net/src/test_interop.rs index 24d3d03abe..5b12c04544 100644 --- a/rs/moq-net/src/test_interop.rs +++ b/rs/moq-net/src/test_interop.rs @@ -1,8 +1,8 @@ //! Bytes cross the language boundary to a Bun script under `test/interop`: subscription -//! responses for subscribers on a mock transport, and varint encodings. +//! responses for subscribers on a mock transport, varint encodings, and moq-lite messages. use crate::coding::{Decode, Encode, VarInt}; -use crate::ietf; +use crate::{ietf, lite}; // Runs a script with a JSON argument, returning its stdout. fn bun(script: &str, input: serde_json::Value) -> Vec { @@ -68,3 +68,106 @@ fn varint_interop() { assert!(buf.is_empty()); } } + +/// Every leading-ones length boundary, plus the JS safe-integer edge and the 62-bit ceiling. +fn lite_values() -> Vec { + let mut values = vec![0]; + for bits in [6, 7, 14, 21, 28, 30, 35, 42, 49, 53, 56] { + values.extend([(1u64 << bits) - 1, 1 << bits]); + } + values.push(VarInt::MAX.into_inner()); + values +} + +/// One GROUP header and its frames, as a publisher writes them: a zigzag timestamp delta, +/// the payload size, then the payload. +fn lite_group(version: lite::Version) -> Vec { + let mut buf = Vec::new(); + let header = lite::Group { + subscribe: 5, + sequence: 1 << 20, + frame_start: 200, + }; + header.encode(&mut buf, version).unwrap(); + for (delta, size) in [(0i64, 10usize), (33_333, 300), (-1_000, 20_000), (1 << 40, 1)] { + VarInt::from_zigzag(delta).unwrap().encode(&mut buf, version).unwrap(); + size.encode(&mut buf, version).unwrap(); + buf.extend(std::iter::repeat_n(0xAB, size)); + } + buf +} + +#[derive(serde::Deserialize)] +struct Echo { + values: Vec, + varints: Vec, + setup: Vec, + datagram: Vec, + group: Vec, + beyond: Vec, +} + +/// JS decodes what Rust encodes back to the same values, and its own encoding of those +/// values is byte for byte Rust's, on lite-06 (QUIC) and lite-07 (leading-ones). +/// +/// Past 2^62-1 the range is per version. JS writes lite-07's 64-bit values, which Rust +/// must refuse with a decode error until its `VarInt` widens; on lite-06 JS refuses them. +#[test] +#[ignore = "requires Bun; run by just test interop"] +fn lite_varint_interop() { + let values = lite_values(); + for version in [lite::Version::Lite06, lite::Version::Lite07] { + let mut varints = Vec::new(); + for value in &values { + value.encode(&mut varints, version).unwrap(); + } + + let setup = lite::Setup { + hop: Some(crate::Hop::new(VarInt::MAX.into_inner()).unwrap()), + ..Default::default() + } + .encode_bytes(version) + .unwrap(); + + let datagram = lite::Datagram { + subscribe: (1 << 40) + 1, + sequence: 300, + timestamp: (1 << 50) + 5, + payload: bytes::Bytes::from_static(b"x"), + } + .encode_bytes(version) + .unwrap(); + + let group = lite_group(version); + + let input = serde_json::json!({ + "version": crate::Version::from(version).alpn(), + "values": values.len(), + "varints": varints, + "setup": setup.to_vec(), + "datagram": datagram.to_vec(), + "group": group, + }); + let echo: Echo = serde_json::from_slice(&bun("lite-varint.ts", input)).expect("JS returned its encodings"); + + let decoded: Vec = echo.values.iter().map(|v| v.parse().unwrap()).collect(); + assert_eq!(decoded, values, "{version}: JS decoded different values"); + assert_eq!(echo.varints, varints, "{version}: JS encoded the values differently"); + assert_eq!(echo.setup, setup, "{version}: SETUP"); + assert_eq!(echo.datagram, datagram, "{version}: datagram"); + assert_eq!(echo.group, group, "{version}: group"); + + match version { + lite::Version::Lite07 => { + let mut expected = vec![0xFF, 0x40, 0, 0, 0, 0, 0, 0, 0]; + expected.extend([0xFF; 9]); + assert_eq!(echo.beyond, expected, "{version}: JS's 64-bit encodings"); + for wire in echo.beyond.chunks(9) { + let err = VarInt::decode(&mut &wire[..], version).unwrap_err(); + assert!(matches!(err, crate::coding::DecodeError::BoundsExceeded), "{err:?}"); + } + } + _ => assert!(echo.beyond.is_empty(), "{version}: JS wrote a value past 2^62-1"), + } + } +} diff --git a/test/interop/README.md b/test/interop/README.md index 0a0a5c1d86..3d961b65ad 100644 --- a/test/interop/README.md +++ b/test/interop/README.md @@ -206,3 +206,10 @@ hands moq-net's QUIC and leading-ones encodings of each varint size boundary `varint.ts`. That script decodes them into js/net's `U64`, checks its `number` conversion, and returns js/net's own encodings, which Rust requires to match byte for byte and decode back to the same value. + +`lite_varint_interop` runs next to it and does the same through moq-lite's +version dispatch: `lite-varint.ts` decodes Rust's lite-06 (QUIC) and lite-07 +(leading-ones) varints, a SETUP carrying a 62-bit Hop ID, a datagram, and a +GROUP stream with frames, and re-encodes them byte for byte. Past 2^62-1 the +range is per version: JS writes lite-07's 64-bit values, which Rust must refuse +with a decode error until its `VarInt` widens, and JS refuses them on lite-06. diff --git a/test/interop/bare-fin.ts b/test/interop/bare-fin.ts index fc899157e1..96db7ddce2 100644 --- a/test/interop/bare-fin.ts +++ b/test/interop/bare-fin.ts @@ -21,6 +21,13 @@ import type { Subscriber as TrackSubscriber } from "../../js/net/src/track.ts"; console.debug = console.error; const input: { version: string; started: boolean; clean: boolean; responses: number[] } = JSON.parse(process.argv[2]); const pair = createMockTransportPair(input.version); +const path = Path.from("room"); +const lite: Record = { + "moq-lite-05": LiteVersion.DRAFT_05, + "moq-lite-06": LiteVersion.DRAFT_06, + "moq-lite-07-wip": LiteVersion.DRAFT_07, +}; +const version = lite[input.version]; const bytes: number[] = []; const output = new Writer( new WritableStream({ @@ -28,26 +35,20 @@ const output = new Writer( bytes.push(...chunk); }, }), + version ?? IetfVersion.DRAFT_19, ); -const path = Path.from("room"); -const lite: Record = { - "moq-lite-05": LiteVersion.DRAFT_05, - "moq-lite-06": LiteVersion.DRAFT_06, - "moq-lite-07-wip": LiteVersion.DRAFT_07, -}; -const version = lite[input.version]; let reader: TrackSubscriber; let peer: Stream | undefined; if (version !== undefined) { const subscriber = new LiteSubscriber(pair.client, version, randomHop()); reader = subscriber.consume(path).track("video").subscribe(); - const info = await Stream.accept(pair.server); + const info = await Stream.accept(pair.server, version); assert(info); assert.equal(await info.reader.u53(), StreamId.Track); await Track.decode(info.reader, version); await new TrackInfo({}).encode(info.writer, version); info.close(); - peer = await Stream.accept(pair.server); + peer = await Stream.accept(pair.server, version); assert(peer); assert.equal(await peer.reader.u53(), StreamId.Subscribe); await Subscribe.decode(peer.reader, version); diff --git a/test/interop/lite-varint.ts b/test/interop/lite-varint.ts new file mode 100644 index 0000000000..fcf72b9211 --- /dev/null +++ b/test/interop/lite-varint.ts @@ -0,0 +1,98 @@ +/** Decode Rust's lite varints and messages, then echo JS's own encoding of them back. */ +import assert from "node:assert/strict"; +import { Datagram } from "../../js/net/src/lite/datagram.ts"; +import { frameDecoder, Group } from "../../js/net/src/lite/group.ts"; +import { Setup } from "../../js/net/src/lite/setup.ts"; +import { Version } from "../../js/net/src/lite/version.ts"; +import { Reader, Writer } from "../../js/net/src/stream.ts"; + +const input: { + version: string; + values: number; + varints: number[]; + setup: number[]; + datagram: number[]; + group: number[]; +} = JSON.parse(process.argv[2]); + +const versions: Record = { + "moq-lite-06": Version.DRAFT_06, + "moq-lite-07-wip": Version.DRAFT_07, +}; +const version = versions[input.version]; +assert(version !== undefined, `unknown version ${input.version}`); + +const reader = (bytes: number[]) => new Reader(undefined, Uint8Array.from(bytes), version); + +// Collects whatever `f` writes with a Writer on this version. +async function write(f: (w: Writer) => Promise): Promise { + const out: number[] = []; + const w = new Writer(new WritableStream({ write: (chunk) => void out.push(...chunk) }), version); + await f(w); + w.close(); + await w.closed; + return out; +} + +// Varints: decode every value Rust wrote, then write them back. +const r = reader(input.varints); +const values: bigint[] = []; +for (let i = 0; i < input.values; i++) values.push(await r.u62()); +assert(await r.done(), "trailing varint bytes"); +const varints = await write(async (w) => { + for (const v of values) await w.u62(v); +}); + +// Past 2^62-1 the range is per version: lite-07 carries the full 64 bits, which JS writes and +// reads back and Rust (62-bit until its VarInt widens) must refuse loudly; lite-06's QUIC form +// cannot express it, so JS refuses to write it, as Rust does. +const huge = [1n << 62n, (1n << 64n) - 1n]; +let beyond: number[] = []; +if (version === Version.DRAFT_07) { + beyond = await write(async (w) => { + for (const v of huge) await w.u62(v); + }); + const back = reader(beyond); + for (const v of huge) assert.equal(await back.u62(), v); +} else { + for (const v of huge) await assert.rejects(write((w) => w.u62(v))); +} + +const setup = await Setup.decode(reader(input.setup), version); +const setupBytes = await write((w) => setup.encode(w, version)); + +const datagram = await Datagram.decode(Uint8Array.from(input.datagram), version); +const datagramBytes = [...datagram.encode(version)]; + +// A GROUP header, then frames until the stream ends. +const g = reader(input.group); +const group = await Group.decode(g, version); +const decode = frameDecoder(1_000_000); +const frames: { timestamp: bigint; payload: Uint8Array }[] = []; +for (;;) { + const frame = await g.decodeMaybe(decode); + if (!frame) break; + frames.push({ timestamp: BigInt(frame.timestamp.value), payload: frame.payload }); +} +const zigzag = (d: bigint) => (d << 1n) ^ (d >> 63n); +const groupBytes = await write(async (w) => { + await group.encode(w, version); + let prev = 0n; + for (const { timestamp, payload } of frames) { + await w.u62(zigzag(timestamp - prev)); + prev = timestamp; + await w.u53(payload.byteLength); + await w.write(payload); + } +}); + +console.log( + JSON.stringify({ + values: values.map((v) => v.toString()), + varints, + setup: setupBytes, + datagram: datagramBytes, + group: groupBytes, + beyond, + }), +); diff --git a/test/interop/varint.ts b/test/interop/varint.ts index 464a34b5f0..5a4689de12 100644 --- a/test/interop/varint.ts +++ b/test/interop/varint.ts @@ -7,7 +7,7 @@ import { U64 } from "../../js/net/src/util/u64.ts"; // Each value is a decimal string, since JSON numbers round past 2^53. const input: { values: string[]; quic: number[][]; leadingOnes: number[][] } = JSON.parse(process.argv[2]); -async function encode(v: U64, version?: Version): Promise { +async function encode(v: U64, version: Version): Promise { const bytes: number[] = []; const writer = new Writer(new WritableStream({ write: (chunk) => void bytes.push(...chunk) }), version); await writer.varint(v); @@ -16,7 +16,7 @@ async function encode(v: U64, version?: Version): Promise { return bytes; } -async function decode(bytes: number[], version?: Version): Promise { +async function decode(bytes: number[], version: Version): Promise { const reader = new Reader(undefined, new Uint8Array(bytes), version); const v = await reader.varint(); assert(await reader.done(), `trailing bytes after ${v}`); @@ -27,7 +27,7 @@ const output: { quic: number[][]; leadingOnes: number[][] } = { quic: [], leadin for (const [i, value] of input.values.entries()) { const expected = U64.fromBigInt(BigInt(value)); for (const [format, version] of [ - ["quic", undefined], + ["quic", Version.DRAFT_16], ["leadingOnes", Version.DRAFT_17], ] as const) { const decoded = await decode(input[format][i], version); diff --git a/test/justfile b/test/justfile index 2f3bc271f9..d98acad97d 100644 --- a/test/justfile +++ b/test/justfile @@ -38,7 +38,8 @@ harness: # client, and the GStreamer moqsrc plugin). --timeout 30 gives headless Chromium # cold-start headroom. Other flags pass through, e.g. # `just test interop --publishers rust,python --subscribers rust,c`. -# Both runs first check js/net's varints against moq-net's encoder at every size boundary. +# Both runs first check js/net's varints against moq-net's encoder at every size boundary, +# and on lite-06/07 messages (the filter also matches `lite_varint_interop`). interop *args: #!/usr/bin/env bash set -euo pipefail