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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .github/workflows/nightly.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 4 additions & 1 deletion doc/concept/moq-lite.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
8 changes: 8 additions & 0 deletions drafts/draft-lcurley-moq-lite.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ normative:
date: false
RFC3986:
RFC6455:
RFC9000:
RFC9002:

informative:
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
5 changes: 4 additions & 1 deletion js/net/bench/frames.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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: () => {} };
Expand Down
146 changes: 146 additions & 0 deletions js/net/bench/lite-varint.ts
Original file line number Diff line number Diff line change
@@ -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<void>;
decode(r: Reader, version: Version): Promise<unknown>;
}

// 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<Uint8Array> {
const chunks: Uint8Array[] = [];
const w = new Writer(new WritableStream<Uint8Array>({ 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<unknown>): Promise<number> {
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<Uint8Array>({
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");
3 changes: 2 additions & 1 deletion js/net/bench/reader.ts
Original file line number Diff line number Diff line change
@@ -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];
Expand Down Expand Up @@ -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);
Expand Down
2 changes: 1 addition & 1 deletion js/net/bench/varint.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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],
},
Expand Down
4 changes: 2 additions & 2 deletions js/net/src/connection/accept.ts
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ async function acceptSetup(
wiring: SessionProps,
): Promise<Established> {
// 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();
Expand Down Expand Up @@ -187,7 +187,7 @@ async function acceptNegotiated(
): Promise<Established> {
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();
Expand Down
2 changes: 1 addition & 1 deletion js/net/src/connection/connect.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
8 changes: 2 additions & 6 deletions js/net/src/ietf/adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand Down Expand Up @@ -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 });

Expand Down
6 changes: 3 additions & 3 deletions js/net/src/ietf/ietf.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Uint8Array> {
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;
Expand All @@ -1077,14 +1077,14 @@ async function encodeNamespace(namespace: Path.Valid): Promise<Uint8Array> {

// Helper to decode a namespace from raw bytes
async function decodeNamespace(bytes: Uint8Array): Promise<Path.Valid> {
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<Uint8Array> {
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();
Expand Down
19 changes: 7 additions & 12 deletions js/net/src/ietf/object.ts
Original file line number Diff line number Diff line change
@@ -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";

Expand Down Expand Up @@ -47,12 +47,7 @@ function hasDeltaObjectPropertyTypes(version: IetfVersion | undefined): boolean
}
}

async function encodeObjectPropertyType(
w: Writer,
id: bigint,
prev: bigint,
version: IetfVersion | undefined,
): Promise<void> {
async function encodeObjectPropertyType(w: Writer, id: bigint, prev: bigint, version: IetfVersion): Promise<void> {
const encoded = hasDeltaObjectPropertyTypes(version) ? id - prev : id;
await w.u62(encoded);
}
Expand All @@ -65,7 +60,7 @@ async function encodeObjectTime(
w: Writer,
timestamp: Timestamp,
timescale: Timescale,
version: IetfVersion | undefined,
version: IetfVersion,
): Promise<void> {
const value = Math.round((timestamp.value * timescale) / timestamp.scale);
await encodeObjectPropertyType(w, PROP_TIMESTAMP, 0n, version);
Expand All @@ -75,7 +70,7 @@ async function encodeObjectTime(
async function encodeObjectExtensions(
timestamp: Timestamp | undefined,
timescale: Timescale,
version: IetfVersion | undefined,
version: IetfVersion,
): Promise<Uint8Array> {
if (timestamp === undefined) {
return new Uint8Array();
Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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<void> {
async encode(w: Writer, flags: GroupFlags, timescale: Timescale, version: IetfVersion, idDelta = 0): Promise<void> {
await w.u53(idDelta);

if (flags.hasExtensions) {
Expand Down Expand Up @@ -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<void> {
async encode(w: Writer, position: FetchPosition, timescale: Timescale, version: IetfVersion): Promise<void> {
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;
Expand Down
Loading
Loading