diff --git a/dart/moq_ffi/lib/src/moq.dart b/dart/moq_ffi/lib/src/moq.dart index e4d5f7532b..b91803e2d7 100644 --- a/dart/moq_ffi/lib/src/moq.dart +++ b/dart/moq_ffi/lib/src/moq.dart @@ -1775,7 +1775,7 @@ class MoqTrackInfo { final int? maxAgeUs; final int? timescale; MoqTrackInfo({ - this.priority = 0, + this.priority = 127, this.maxAgeUs = null, this.timescale = null, }); @@ -12573,7 +12573,7 @@ void _checkApiChecksums() { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqtrackconsumer_recv_datagram() != - 29049) { + 17412) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqtrackconsumer_recv_group() != 60887) { diff --git a/doc/concept/standard.md b/doc/concept/standard.md index 54801adec8..a5d59d0daa 100644 --- a/doc/concept/standard.md +++ b/doc/concept/standard.md @@ -34,18 +34,35 @@ the block; other groups and the session stay open. An IETF publisher declares the track's default priority in `SUBSCRIBE_OK` or `PUBLISH` when that draft carries track properties. Groups without a priority flag inherit it. If the property is absent, the IETF wire default of 128 maps -to model priority 127, where higher values are served first. - -On drafts 14–19, the Rust publisher serves relative joining `FETCH` requests -with offset zero for `NextObject` subscriptions. The fetch delivers the saved -current-group prefix, and the subscription delivers later objects. Standalone, -absolute joining, and nonzero-offset fetches are refused. Draft-20 uses -subscription fills instead. JavaScript publishing does not yet serve `FETCH`; +to model priority 127, where higher values are served first. A track that +never sets a priority is 127 as well, so it goes out as 128 on IETF and 127 +on moq-lite. + +On drafts 14–19, the Rust publisher answers a standalone `FETCH` within one +group from the cache. A relay fetches a missing group upstream with a `FETCH` +of that one whole group, and an upstream refusal is the refusal the fetcher +sees. A range touching several groups is refused with `NOT_SUPPORTED`, as is +any `FETCH` on draft-20 and later, which moved the range into +`LOCATION_FILTER`. A standalone `FETCH` +carries no timestamps, since no `SUBSCRIBE_OK` declared a timescale for it. + +On drafts 14–19, the Rust publisher also serves relative and absolute joining +`FETCH` requests for `NextObject` subscriptions, for the subscription group's +saved prefix only, while the subscription delivers later objects. One reaching +back to earlier groups is refused with `NOT_SUPPORTED`. Draft-20 uses +subscription fills instead. JavaScript +publishing does not yet serve `FETCH`; Rust and JavaScript subscribers request unfiltered delivery on older drafts because they do not issue joining fetches. Other publishers may replay a cached backlog for that filter; selecting the next group instead would leave static tracks waiting for a group that never arrives. +A moq-lite datagram is a single-frame group, so on moq-transport it travels +as an `OBJECT_DATAGRAM` at object 0 whose Group ID is the sequence, and a relay +forwards it without renumbering. A datagram carrying any other Object ID, or a +status other than Normal, is dropped. JavaScript does not yet carry datagrams +on moq-transport. + A client may present one credential in its `SETUP` with the `AUTHORIZATION TOKEN` option. The server reads a value (`USE_VALUE`, or `REGISTER`, which it treats as a value since it advertises no token cache) and hands its Token Type diff --git a/doc/lib/rs/moq-net.md b/doc/lib/rs/moq-net.md index 910211481f..322894a4c3 100644 --- a/doc/lib/rs/moq-net.md +++ b/doc/lib/rs/moq-net.md @@ -21,7 +21,7 @@ above ([hang](/lib/rs/hang)); relays and CDNs implement only this. - **Tracks** carry groups with a priority, a retention window, and a timescale. Subscribers set their own priority and max age and can change them live. - **Groups** are written frame by frame and delivered on independent streams. Old groups are cached for fetch-by-sequence; stale groups are skipped per the subscriber's budget. - **Track ends**: `finish()` ends a track at its live edge, while `finish_at(n)` declares the exclusive end ahead of it and still accepts the groups below. A subscriber awaits it with `finished()`. A remote track ends only once every group below its end has arrived or was dropped; one reset before its header arrived is skipped after the subscription's max age on moq-lite (one second without one), or after one second on IETF. -- **Datagrams** send a single small frame unreliably on moq-lite 05+. +- **Datagrams** send a single small frame unreliably on moq-lite 05+ and moq-transport. - **Routes** record the relay hops and a cost, which is what the relay [cluster](/bin/relay/cluster) routes on. A hop of 0 marks the chain anonymous: `Route::is_anonymous()` is true, and that route ranks below every fully identified one. `Route::source()` says where a delivered route entered: `Source::Local`, or `Source::Peer(hop)` when a handle marked `origin::Producer::peer()` announced it. `origin::Consumer::local()` sees only the local ones. - **Stats** counters per broadcast and session, drained by [`moq-stats`](https://docs.rs/moq-stats). diff --git a/drafts/draft-lcurley-moq-e2ee.md b/drafts/draft-lcurley-moq-e2ee.md index ec68f21aa9..25ca49c8ea 100644 --- a/drafts/draft-lcurley-moq-e2ee.md +++ b/drafts/draft-lcurley-moq-e2ee.md @@ -227,7 +227,8 @@ An Object ID above `2^32-1` MUST be refused as `identity` before encryption or d ## Datagrams A datagram uses `domain = 0x01`, its 64-bit sequence as `group`, and `frame = 0`. -MoQ Transport has no datagram mapping in this profile; shared vectors cover grouped tracks on both transports and datagrams on moq-lite only. +On MoQ Transport, a datagram is an Object at ID 0 in an OBJECT_DATAGRAM whose Group ID is the sequence. +Shared vectors cover grouped tracks on both transports and datagrams on moq-lite only. ## Nonce The 96-bit AES-GCM nonce is: diff --git a/go/wrapper/README.md b/go/wrapper/README.md index 781c8d055e..1f9d375030 100644 --- a/go/wrapper/README.md +++ b/go/wrapper/README.md @@ -110,7 +110,7 @@ Raw tracks support best-effort datagrams alongside groups: `TrackProducer.Append sends one `Frame` and returns its sequence number, while `TrackConsumer.RecvDatagram` and `TrackConsumer.Datagrams` receive them in arrival order. Payloads are capped at 1200 bytes. Datagram delivery requires a datagram-capable transport and lite-05 or -newer moq-lite; IETF moq-transport, pre-lite-05, WebSocket, and TCP paths do not +newer moq-lite, or moq-transport; pre-lite-05, WebSocket, and TCP paths do not deliver them, and there is no stream fallback. ## Versioning diff --git a/go/wrapper/types.go b/go/wrapper/types.go index 40bf4f53ea..5c21af7a4d 100644 --- a/go/wrapper/types.go +++ b/go/wrapper/types.go @@ -48,6 +48,7 @@ type ( // Subscription holds subscriber-side delivery preferences: priority, ordering, max age, and group range. Subscription = ffi.MoqSubscription // TrackInfo holds publisher-side track properties: priority, ordering, max age, and timescale. + // A zero Priority is the least urgent, not the default; set 127 for the midpoint a nil TrackInfo uses. TrackInfo = ffi.MoqTrackInfo // Video describes one catalog rendition, including whether the publisher recommends temporarily avoiding it. Video = ffi.MoqVideo diff --git a/js/net/src/ietf/priority.test.ts b/js/net/src/ietf/priority.test.ts index a05f8b2650..b04da83354 100644 --- a/js/net/src/ietf/priority.test.ts +++ b/js/net/src/ietf/priority.test.ts @@ -1,4 +1,5 @@ import { expect, test } from "bun:test"; +import { infoDefaults } from "../track.ts"; import { fromWire, toWire } from "./priority.ts"; test("IETF subscriber priority is lower first", () => { @@ -17,3 +18,7 @@ test("subscriber priority round trips", () => { expect(fromWire(toWire(priority))).toBe(priority); } }); + +test("an unset track priority is the draft's usual publisher priority", () => { + expect(toWire(infoDefaults().priority)).toBe(128); +}); diff --git a/js/net/src/lite/track.test.ts b/js/net/src/lite/track.test.ts index ca4d25d44f..34217581b9 100644 --- a/js/net/src/lite/track.test.ts +++ b/js/net/src/lite/track.test.ts @@ -48,7 +48,7 @@ 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), Version.DRAFT_05)).toEqual( - new Uint8Array([0x06, 0x00, 0x00, 0x53, 0x88, 0x43, 0xe8]), + new Uint8Array([0x06, 0x7f, 0x00, 0x53, 0x88, 0x43, 0xe8]), ); }); diff --git a/js/net/src/track.test.ts b/js/net/src/track.test.ts index c8c516b67e..04ac64ae9b 100644 --- a/js/net/src/track.test.ts +++ b/js/net/src/track.test.ts @@ -27,9 +27,9 @@ function mockMonotonicTime(initial: number) { }; } -test("priority reads the committed info and is 0 before accept", () => { +test("priority reads the committed info and is the midpoint before accept", () => { const producer = new TrackProducer("video"); - expect(producer.priority).toBe(0); + expect(producer.priority).toBe(127); producer.accept({ priority: 60 }); expect(producer.priority).toBe(60); }); diff --git a/js/net/src/track.ts b/js/net/src/track.ts index 52dc4034e0..45df6e4320 100644 --- a/js/net/src/track.ts +++ b/js/net/src/track.ts @@ -30,6 +30,10 @@ const PRUNE_SLICES = 8; /** Default {@link Info.maxAge} window (milliseconds) when the publisher does not set one. */ export const DEFAULT_MAX_AGE_MS = Milli(5000); +// The higher-first midpoint. IETF flips priority (lower first), so this goes out as 128, the +// draft's usual publisher priority, while moq-lite carries 127 as written: one urgency on both. +const DEFAULT_PRIORITY = 127; + /** Maximum buffered datagrams per subscriber; mirrors Rust's bounded send buffer. */ const MAX_DATAGRAMS = 64; @@ -70,7 +74,7 @@ export interface Info { * or non-finite value and a result past `Number.MAX_SAFE_INTEGER`. */ maxAge: Milli; - /** Tie-break priority between subscriptions of equal subscriber priority (`0..=255`). */ + /** Tie-break priority between subscriptions of equal subscriber priority (`0..=255`, higher first). Defaults to `127`. */ priority: number; } @@ -103,7 +107,7 @@ export function infoDefaults(info: Partial = {}): Info { return { timescale: Timescale(info.timescale ?? Timescale.MILLI), maxAge: maxAgeMillis(info.maxAge ?? DEFAULT_MAX_AGE_MS), - priority: priorityByte(info.priority ?? 0), + priority: priorityByte(info.priority ?? DEFAULT_PRIORITY), }; } @@ -481,13 +485,13 @@ export class Producer { } /** - * Publisher priority from the committed {@link Info}, or 0 before {@link accept}. + * Publisher priority from the committed {@link Info}, or the default before {@link accept}. * * Higher is served first. Hang publishers set this from `Catalog.PRIORITY` so * audio outranks video on the wire and in the bandwidth allocator. */ get priority(): number { - return this.#state.info.peek()?.priority ?? 0; + return this.#state.info.peek()?.priority ?? DEFAULT_PRIORITY; } /** diff --git a/kt/README.md b/kt/README.md index 2afe6524f6..7e0958fd69 100644 --- a/kt/README.md +++ b/kt/README.md @@ -52,7 +52,7 @@ The `dev.moq` package is intentionally thin: Kotlin has extension functions, so - **Fetched media**: `fetchMediaGroup(...).frames()` streams the decoded frames of one retained group, then completes. - **Duration extensions** (`Durations.kt`): the FFI carries microseconds as integers, so `stats.rtt`, `backoff.initial`, `frame.timestamp`, and their siblings read back as a `kotlin.time.Duration`. - **`logLevel(...)`**: configures native Rust tracing without importing the raw bindings package. -- **Raw datagrams**: `TrackProducer.appendDatagram(Frame(payload, timestampUs))` sends one best-effort frame and returns its sequence; `TrackConsumer.recvDatagram()` and `datagrams()` receive them. Payloads are capped at 1200 bytes, require a datagram-capable transport plus lite-05 or newer moq-lite, and have no stream fallback. +- **Raw datagrams**: `TrackProducer.appendDatagram(Frame(payload, timestampUs))` sends one best-effort frame and returns its sequence; `TrackConsumer.recvDatagram()` and `datagrams()` receive them. Payloads are capped at 1200 bytes, require a datagram-capable transport plus lite-05 or newer moq-lite or moq-transport, and have no stream fallback. - **`MoqException.isShutdown`** (`Errors.kt`): true for the graceful `Cancelled`/`Closed` cases. ## Versioning diff --git a/py/moq-rs/README.md b/py/moq-rs/README.md index 0902031a5a..90281ccc6d 100644 --- a/py/moq-rs/README.md +++ b/py/moq-rs/README.md @@ -212,7 +212,7 @@ Every handle whose cleanup is `cancel()` is an async context manager, so exiting - **`Catalog`**. `.audio: dict[str, Audio]`, `.video: dict[str, Video]`, `.display`, `.rotation`, `.flip`. - **`Frame`**. `.payload: bytes`, `.timestamp_us: int`. The unit of every write and every raw read. - **`MediaFrame`**. `.payload: bytes`, `.timestamp_us: int`, `.keyframe: bool`. Returned by media subscriptions. `keyframe` marks a group start or video keyframe; for audio it is true only at a group start. -- **`Datagram`**. `.sequence: int`, `.timestamp_us: int`, `.payload: bytes`. Delivered only on datagram-capable transports and lite-05 or newer moq-lite. +- **`Datagram`**. `.sequence: int`, `.timestamp_us: int`, `.payload: bytes`. Delivered only on datagram-capable transports with lite-05 or newer moq-lite, or moq-transport. - **`Audio`**. `.codec`, `.sample_rate`, `.channel_count`, `.bitrate`, `.description`. - **`Video`**. `.codec`, `.coded: Dimensions`, `.display_aspect`, `.bitrate`, `.stalled`, `.framerate`, `.description`. A true `.stalled` recommends temporarily avoiding the rendition without making it unavailable. - **`Subscription`**. Subscriber delivery preferences: priority, staleness, and optional group range. diff --git a/quest/m0/ietf-legal-input.md b/quest/m0/ietf-legal-input.md index c90b5f547b..32ebf87922 100644 --- a/quest/m0/ietf-legal-input.md +++ b/quest/m0/ietf-legal-input.md @@ -14,7 +14,7 @@ Each of these closes the session today with `PROTOCOL_VIOLATION`: - FETCH on drafts 20+ still decodes the removed Fetch Type field (`ietf/fetch.rs`). Decode the draft-20 layout: namespace, name and params, with the range in LOCATION_FILTER. Refuse it `NOT_SUPPORTED` until - [moxygen FETCH](/quest/m1/moxygen/fetch.md) serves it. Fix the encoder and + [Draft-20 FETCH](/quest/m1/ietf-fetch-location.md) serves it. Fix the encoder and the pinned test in `version.rs`. - Request parameters the draft allows on a message fail `decode_params!`: AUTHORIZATION TOKEN (0x03) anywhere, NEW_GROUP_REQUEST (0x32) and the diff --git a/quest/m0/ietf-subgroup-refusal.md b/quest/m0/ietf-subgroup-refusal.md index 969ea14338..872ad8504c 100644 --- a/quest/m0/ietf-subgroup-refusal.md +++ b/quest/m0/ietf-subgroup-refusal.md @@ -20,7 +20,3 @@ is the check. Moved from the moxygen line to m0 in the 2026-09-30 audit: a non-zero subgroup ending the upstream session is exactly m0's "legal input never fails a session", and moxygen will send it at Seattle interop on 2026-10-12. - -## Related - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - subgroups stay out of scope; only the blast radius is in diff --git a/quest/m1/README.md b/quest/m1/README.md index 027063f5ee..72c49c8bda 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -61,7 +61,8 @@ QUIC studies there on that rule. - [Data track clock](/quest/m1/data-track-clock.md) - JSON and binary data tracks stamp on the catalog's clock at write time, matching the media's anchored clock - [Go and Dart doc samples](/quest/m1/doc-samples-go-dart.md) - Go and Dart doc samples compile against their wrappers - [Data capture in bindings](/quest/m1/data-capture-bindings.md) - moq-ffi and every wrapper pass a data frame's capture time, and the JSON window producer takes one -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - one subgroup per group, whole-group FETCH, and one datagram per group, never a full moxygen pass +- [Draft-20 FETCH](/quest/m1/ietf-fetch-location.md) - a draft-20+ FETCH within one group is served from its LOCATION_FILTER, as older drafts are +- [Fetch without SUBSCRIBE](/quest/m1/ietf-fetch-only.md) - a relay fetches from an IETF upstream without subscribing, finished tracks included, with End of Track always reported - [JS IETF datagrams](/quest/m1/js-ietf-datagram.md) - `@moq/net` sends and receives datagram groups over moq-transport, like Rust - [Datagrams are live-only](/quest/m1/datagram-unfetchable.md) - no FETCH, replay to a new subscriber, or cache fill ever returns a datagram group - [#2991](/quest/m1/2991-net-coalesce-dynamic-tracks-and-preserve-sequences-across.md) - one dynamic producer per track name in both languages, with the sequence namespace surviving a replacement diff --git a/quest/m1/datagram-unfetchable.md b/quest/m1/datagram-unfetchable.md index 2de867c657..76fa052c69 100644 --- a/quest/m1/datagram-unfetchable.md +++ b/quest/m1/datagram-unfetchable.md @@ -24,7 +24,7 @@ datagram flag, and reports it as "unsupported": - A received fetch object with the datagram flag is refused as not fetchable. It fails only that fetch or fill, not the session, because draft-16+ allows the flag (maintainer, 09-29, on Codex's review). That covers the relay's group - fill (`recv_group_fetch_objects`, on the moxygen line) and the joining-fetch + fill (`recv_group_fetch_objects`) and the joining-fetch fill (`run_fill_objects`). - Update `drafts/draft-lcurley-moq-lite.md` to say datagrams are neither cached nor fetchable, and `doc/concept/`. Run `just drafts check` and @@ -51,6 +51,5 @@ range apart from one dropped for any other reason. ## Related -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - says it carries one datagram per group live; this keeps them out of FETCH - [JS IETF datagrams](/quest/m1/js-ietf-datagram.md) - JS datagram send and receive over moq-transport - [SUBSCRIBE_DROP](/quest/m1/subscribe-drop.md) - datagram groups stay best effort there too diff --git a/quest/m1/e2ee/README.md b/quest/m1/e2ee/README.md index 07e41ce3c9..0dcb75b674 100644 --- a/quest/m1/e2ee/README.md +++ b/quest/m1/e2ee/README.md @@ -40,7 +40,7 @@ The Rust and TypeScript cores expose the same surface, and nothing else: - Deterministic secret-derived physical names hide catalog, codec, role, quality, timeline, and custom-track semantics. Authorized clients derive the encrypted catalog track name, then learn the remaining opaque names from its decrypted contents. Every catalog representation is encrypted; Rust publishers must not emit a plaintext MSF catalog. - A platform that forwards and meters protected bytes must never preview, record, archive, transmux, transcode, transcribe, compose, or inspect them, rejecting those paths before opening a processing session or writing product state. Applications needing those operations terminate E2EE outside the platform. A platform classifies protected broadcasts by its own credential or product state, never by name; the moq.pro (downstream) exclusion classifier and dashboard work stay downstream. -- The first proof covers browser TypeScript and native Rust publication and playback in both directions, with grouped audio and video over both moq-lite and MoQ Transport. Shared vectors cover groups and moq-lite datagrams; MoQ Transport has no datagram delivery. +- The first proof covers browser TypeScript and native Rust publication and playback in both directions, with grouped audio and video over both moq-lite and MoQ Transport. Shared vectors cover groups and moq-lite datagrams; JavaScript has no MoQ Transport datagram delivery yet. ## Required diff --git a/quest/m1/ietf-fetch-location.md b/quest/m1/ietf-fetch-location.md new file mode 100644 index 0000000000..54daf01a6a --- /dev/null +++ b/quest/m1/ietf-fetch-location.md @@ -0,0 +1,23 @@ +# [S] Draft-20 FETCH + +## Goal + +On drafts 20 and later, a FETCH whose LOCATION_FILTER stays within one group +is answered the way drafts 14 to 19 answer a standalone FETCH: from cache, with +a miss fetched upstream, and an upstream refusal passed through. A range +touching several groups is refused `NOT_SUPPORTED`, as on older drafts. + +## Plan + +Today both sides refuse FETCH on draft-20+ `NOT_SUPPORTED`: the publisher in +`run_fetch_stream`, and the relay's upstream group fill in `run_group_fetch`, +because the codec still carries the removed Fetch Type field. +[Legal IETF input](/quest/m0/ietf-legal-input.md) fixes the codec. Then drop +both refusals, and flip `fetch_moq_transport_20` in `rs/moq-tokio/tests` back +to a served fetch. + +Update the draft-20 note in `doc/concept/standard.md`. + +## Required + +- [Legal IETF input](/quest/m0/ietf-legal-input.md) - decodes the draft-20 FETCH this serves diff --git a/quest/m1/ietf-fetch-only.md b/quest/m1/ietf-fetch-only.md new file mode 100644 index 0000000000..a7bdc908e6 --- /dev/null +++ b/quest/m1/ietf-fetch-only.md @@ -0,0 +1,29 @@ +# [M] Fetch without SUBSCRIBE + +## Goal + +A relay serves a FETCH-only demand for an IETF upstream track without +SUBSCRIBE upstream. A finished upstream track can still be fetched, and a +FETCH_OK that reaches the end of the track says so every time. + +## Plan + +The origin splices a route only once the track's info is known. moq-lite gets +it from TRACK_INFO, so its fetch-only demand never subscribes. IETF has no +TRACK_INFO, so today fetch-only demand still makes the relay SUBSCRIBE upstream. +Two symptoms follow: + +- A finished upstream track refuses the SUBSCRIBE, so its groups cannot be + fetched either. +- The live subscription races the group FETCHes. End of Track is known only + once an upstream FETCH_OK reports it, so a downstream FETCH that runs to the + end can answer before that and leave End of Track unset. moxygen's "FETCH + with large objects" case flakes on this. + +TRACK_STATUS is the likely source of the info. The publisher refuses it today, +so both sides are in scope. Keep the SUBSCRIBE path for real subscription +demand. + +## Related + +- [JavaScript FETCH](/quest/m1/js-fetch.md) - the browser publisher answers these fetches diff --git a/quest/m1/js-ietf-datagram.md b/quest/m1/js-ietf-datagram.md index db210e7b25..7ecfb8faab 100644 --- a/quest/m1/js-ietf-datagram.md +++ b/quest/m1/js-ietf-datagram.md @@ -1,29 +1,21 @@ -# [M] @moq/net datagrams over moq-transport +# [M] JavaScript moq-transport datagrams ## Goal -A `@moq/net` session on moq-transport sends and receives datagram groups as -`OBJECT_DATAGRAM`, one object per group with its sequence kept, as Rust does. -A datagram track crosses IETF between Rust and JS in both directions in -`just test interop`. +`@moq/net` sends and receives datagrams over moq-transport as +`OBJECT_DATAGRAM`, matching Rust: one Object at ID 0 is a single-frame group +whose Group ID is the sequence. A JavaScript publisher's datagrams reach a Rust +subscriber, and a Rust publisher's reach a JavaScript one. ## Plan -- JS routes datagrams on moq-lite only (`js/net/src/lite/datagram.ts`, - `runDatagrams` from lite-05); `js/net/src/ietf/` reads and writes none. - Mirror the lite path on the IETF session and the Rust mapping from #4274: - receive decodes `OBJECT_DATAGRAM` and inserts on the aliased subscription's - track; send writes object 0 with END_OF_GROUP, the explicit publisher - priority, and the timestamp property when the track has a timescale. -- Keep Rust's edges: an Object ID other than 0, a non-Normal status, or an - unbound alias is dropped; a malformed Type closes the session. Rust covers - drafts 14 and later, whose Type flags differ between 14 and 15+; decide - what draft 07, which JS also speaks, does. -- Add IETF datagram cases beside the lite ones in the interop harness. +Port `rs/moq-net/src/ietf/datagram.rs` and the session's send and receive +loops. Decode every draft's Type flags, drop what the model cannot carry the +same way Rust does, and close the session on a malformed Type. -Public API: none expected. Wire: `@moq/net` moq-transport sessions send and -accept `OBJECT_DATAGRAM` as the drafts define; no project draft changes. +The integration test `ietf does not deliver datagrams` flips to delivery on +every supported draft, and `just test interop --all` covers both directions. -## Required +## Related -- [moxygen interop](/quest/m1/moxygen/README.md) - its datagram groups quest (done on the line) is the Rust side this mirrors and interops with +- [Datagrams are live-only](/quest/m1/datagram-unfetchable.md) - the subscribe range for datagrams, settled on both protocols diff --git a/quest/m1/moxygen/README.md b/quest/m1/moxygen/README.md deleted file mode 100644 index e478e9d030..0000000000 --- a/quest/m1/moxygen/README.md +++ /dev/null @@ -1,40 +0,0 @@ -# Moxygen compatibility - -## Goal - -More compatibility with moxygen's moq-test suite, never a full pass. A peer -that speaks one subgroup per group, FETCH of whole groups, and one datagram -per group gets those through the relay. - -## Plan - -`conformance_test.sh` on draft-16 scored 0/76 against moq-relay 0.15.1. -Subgroup subscribes did return objects. The relay decodes IETF into the -moq-lite track and writes a new subgroup. A green run of all 76 cases is not -the exit test. - -Out of scope, so they are not reopened as bugs: - -- PUBLISH. The relay refuses it. Routing stays SUBSCRIBE per namespace. -- Subgroups. A non-zero subgroup id is dropped. One stream per group. -- Per-group priority. `track::Info` has one priority. A 200/201 split by - group parity stays a gap. -- Foreign object extensions. Dropped. A timestamp we add is ours. -- The end-of-group bit. Optional in draft-16. The group's one stream FINs - at the end, so the bit adds nothing this line will do. -- Several datagram objects in one group. -- A peer that answers SUBSCRIBE_NAMESPACE with unimplemented. The session - already continues. - -Docs stay inline in the change that makes them stale. No new guide. - -## Required - -- [Default track priority](/quest/m1/moxygen/priority.md) - an unset track priority is the midpoint on moq-lite and on IETF, not the least urgent value -- [Group FETCH](/quest/m1/moxygen/fetch.md) - an IETF FETCH of whole groups is served from cache or fetched upstream, one group at a time -- [Datagram groups](/quest/m1/moxygen/datagram.md) - an IETF datagram that is one object in a group arrives as a moq-lite datagram group - -## Related - -- [JavaScript FETCH](/quest/m1/js-fetch.md) - the browser publisher answers a FETCH this relay forwards -- [Track priority scope](/quest/m1/track-priority-scope.md) - send-order fairness, not the default value diff --git a/quest/m1/moxygen/datagram.md b/quest/m1/moxygen/datagram.md deleted file mode 100644 index c794b5e57d..0000000000 --- a/quest/m1/moxygen/datagram.md +++ /dev/null @@ -1,23 +0,0 @@ -# [M] Datagram groups - -## Goal - -An IETF datagram that carries one object for a group is delivered as a -moq-lite datagram: one single-frame group, sequence preserved. Several -objects in one datagram group stay unsupported. - -## Plan - -moq-lite already does this. A datagram is a subscribe id, a group sequence, a -timestamp, and a payload. `insert_datagram` keeps the sequence so a relay -does not renumber it. - -The IETF session does not read or write QUIC datagrams, so a datagram -subscribe delivers nothing. Copy the lite path onto that session. Do not -invent a second object list inside the group. - -The moxygen cases with one object per group are the ones this can pass. - -## Related - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line this belongs to diff --git a/quest/m1/moxygen/fetch.md b/quest/m1/moxygen/fetch.md deleted file mode 100644 index e38d1e533b..0000000000 --- a/quest/m1/moxygen/fetch.md +++ /dev/null @@ -1,31 +0,0 @@ -# [L] Group FETCH - -## Goal - -An IETF FETCH that asks for whole groups is answered. A group already in -cache is served from there. A miss fetches that group upstream. Subscribers -still FETCH one group at a time. A publisher may be asked for a range, which -the relay walks one group at a time. A publisher refusal is what the -subscriber sees. - -## Plan - -Standalone FETCH and a non-zero joining FETCH are refused today with -"not supported". Those forms exist only before draft-20; on draft-20+ the -range comes from LOCATION_FILTER. Walk `track::Consumer::fetch_group`, one group, then the -next. Do not add an archive. - -A joining FETCH is the same walk for the groups it names. A form the walk -cannot express is still an explicit refusal, not a hang. - -The moxygen FETCH cases that ask for whole groups are the check. The rest of -that suite is not. - -## Required - -- [Legal IETF input](/quest/m0/ietf-legal-input.md) - decodes the draft-20+ FETCH layout this serves - -## Related - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line this belongs to -- [JavaScript FETCH](/quest/m1/js-fetch.md) - the browser publisher that fills an upstream miss diff --git a/quest/m1/moxygen/priority.md b/quest/m1/moxygen/priority.md deleted file mode 100644 index 6efd6aaa5a..0000000000 --- a/quest/m1/moxygen/priority.md +++ /dev/null @@ -1,28 +0,0 @@ -# [S] Default track priority - -## Goal - -A track that never set a priority is the midpoint on both wires. moq-lite -TRACK_INFO carries that midpoint as written. An IETF subgroup header carries -128. A track that set `track::Info.priority` still uses that value. Per-group -priority stays unsupported. - -## Plan - -`track::Info.priority` defaults to 0, higher-first. moq-lite writes it -verbatim, so the wire byte is 0. IETF flips it, so 0 goes out as 255. The -draft's usual publisher priority is 128. - -One model default serves both. 127 is the higher-first midpoint: lite writes -127, IETF writes 128. Do not special-case each wire to the byte 128. That -would make the two encodings different urgencies. - -Rust and JavaScript both pin the old default TRACK_INFO bytes -(`0x06, 0x00, ...`). Update those fixtures with the new default. - -The moxygen 200 versus 201 split by group parity is out of scope. - -## Related - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line this belongs to -- [Track priority scope](/quest/m1/track-priority-scope.md) - fairness across owners, not this default diff --git a/quest/m2/signed-priority.md b/quest/m2/signed-priority.md index 5f2136fb72..bd41dd592a 100644 --- a/quest/m2/signed-priority.md +++ b/quest/m2/signed-priority.md @@ -19,11 +19,10 @@ Decided with the maintainer: range, with no saturation. An unset priority goes out as IETF 127, not the draft's usual 128, and an absent IETF priority decodes to 0, the unset default. Pin both ends in tests. -- Invariant, kept from the [moxygen line](/quest/m1/moxygen/README.md)'s - default-priority quest (#4273): one urgency on both wires, so the IETF byte - is always `255 -` the lite byte. Mapping IETF as `128 - p` to hit the - draft's 128 would break it, which is why the default is one step off the - draft there. +- Invariant, kept from the default-priority change (#4273): one urgency on + both wires, so the IETF byte is always `255 -` the lite byte. Mapping IETF + as `128 - p` to hit the draft's 128 would break it, which is why the + default is one step off the draft there. - hang's built-in priorities move above 0, so hang media outranks a track that never set one. Something like catalog 40, text 30, audio 20, video 10; the spacing is the implementer's call. Rust and JS keep matching values. @@ -48,10 +47,6 @@ Decided in the 2026-09-30 audit: moved to m2. It stays deferred unless it ships in the same `dev` release as the moxygen default change, so the default byte moves once instead of twice. -## Required - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - ships the 127 default and the one-urgency invariant this re-maps - ## Related - [Scope track priority](/quest/m1/track-priority-scope.md) - which streams a priority competes with, not its type diff --git a/rs/libmoq/src/api.rs b/rs/libmoq/src/api.rs index d8c6dbcfda..5d429aaede 100644 --- a/rs/libmoq/src/api.rs +++ b/rs/libmoq/src/api.rs @@ -491,8 +491,8 @@ pub struct moq_datagram { /// Publisher-side raw track properties. /// /// A null [moq_publish_track] `info` pointer uses the moq-net defaults. -/// A zero-initialized struct also uses those defaults, except `priority` where -/// zero is the default itself. +/// A zero-initialized struct also uses those defaults, except `priority`, which +/// has no presence flag: zero is the least urgent, and 127 is the moq-net default. #[repr(C)] #[allow(non_camel_case_types)] pub struct moq_track_info { @@ -3619,7 +3619,7 @@ pub extern "C" fn moq_consume_track_cancel(track: u32) -> i32 { /// touched again, so release `user_data` there. The terminal callback fires even after /// [moq_consume_datagrams_cancel]. Read each datagram with [moq_consume_datagram] and release /// it with [moq_consume_datagram_free]. Datagrams arrive only over datagram-capable -/// transports and lite-05 or newer moq-lite; there is no stream fallback. +/// transports on moq-transport or lite-05 and newer moq-lite; there is no stream fallback. /// /// Returns a non-zero handle to the subscription on success, or a negative code on failure. /// diff --git a/rs/moq-e2ee/tests/transport.rs b/rs/moq-e2ee/tests/transport.rs index bb114105e2..72dde42bc7 100644 --- a/rs/moq-e2ee/tests/transport.rs +++ b/rs/moq-e2ee/tests/transport.rs @@ -1,4 +1,4 @@ -//! Grouped tracks on moq-lite and MoQ Transport; datagrams on moq-lite. +//! Grouped tracks and datagrams on moq-lite and MoQ Transport. mod support; @@ -104,10 +104,9 @@ async fn grouped_over_ietf() { .expect("timed out"); } -#[tokio::test] -async fn datagrams_over_lite() { +async fn datagram_roundtrip(version: &str) { tokio::time::timeout(TEST_TIMEOUT, async { - let mut fixture = connect_protected("moq-lite-05".parse().unwrap(), "audio").await; + let mut fixture = connect_protected(version.parse().unwrap(), "audio").await; fixture .producer .append_datagram(Timestamp::from_millis(9).unwrap(), b"opus") @@ -123,3 +122,13 @@ async fn datagrams_over_lite() { .await .expect("timed out"); } + +#[tokio::test] +async fn datagrams_over_lite() { + datagram_roundtrip("moq-lite-05").await; +} + +#[tokio::test] +async fn datagrams_over_ietf() { + datagram_roundtrip("moq-transport-21").await; +} diff --git a/rs/moq-ffi/src/consumer.rs b/rs/moq-ffi/src/consumer.rs index 84da7e0076..14e441e320 100644 --- a/rs/moq-ffi/src/consumer.rs +++ b/rs/moq-ffi/src/consumer.rs @@ -633,7 +633,7 @@ impl MoqTrackConsumer { /// Receive the next best-effort datagram in arrival order. /// /// Returns `None` when the track ends. Datagram delivery is unavailable over - /// IETF moq-transport, pre-lite-05 moq-lite, and stream-only transports. + /// pre-lite-05 moq-lite and stream-only transports. /// Datagrams are a separate cursor from groups, so this works alongside either /// group order, never commits the track to one, and progresses while a group /// read is pending. diff --git a/rs/moq-ffi/src/producer.rs b/rs/moq-ffi/src/producer.rs index 377ec3c7a4..dfee87f50d 100644 --- a/rs/moq-ffi/src/producer.rs +++ b/rs/moq-ffi/src/producer.rs @@ -11,11 +11,12 @@ use crate::media::{MoqAudioInit, MoqContainerFormat, MoqContainerInit, MoqFrame, /// Publisher-side track properties, mirroring [`moq_net::track::Info`]. /// /// Construct with the fields you care about; the rest use raw-track defaults -/// (priority 0, the publisher's default max age, microsecond timescale). +/// (priority 127, the publisher's default max age, microsecond timescale). #[derive(Clone, uniffi::Record)] pub struct MoqTrackInfo { /// Priority, used only to break ties between subscriptions of equal subscriber priority. - #[uniffi(default = 0)] + /// Higher is more urgent; the default 127 is the midpoint. + #[uniffi(default = 127)] pub priority: u8, /// Maximum age of a non-latest group before the publisher evicts it, in /// microseconds. Null uses the default. This is the publisher-side half of diff --git a/rs/moq-net/src/fuzz.rs b/rs/moq-net/src/fuzz.rs index a7dd026195..36c0a4be23 100644 --- a/rs/moq-net/src/fuzz.rs +++ b/rs/moq-net/src/fuzz.rs @@ -61,7 +61,7 @@ const IETF_VERSIONS: &[ietf::Version] = &[ const LITE_KINDS: u8 = 21; /// How many types [`ietf_wire`] dispatches over. -const IETF_KINDS: u8 = 39; +const IETF_KINDS: u8 = 40; /// Split the two selector bytes off the input: a version and a type. fn select(data: &[u8], versions: usize) -> Option<(usize, u8, &[u8])> { @@ -436,6 +436,7 @@ pub fn ietf_wire(data: &[u8]) -> bool { 36 => roundtrip::(rest, version, stable), 37 => roundtrip::(rest, version, stable), 38 => roundtrip::(rest, version, stable), + 39 => roundtrip::(rest, version, stable), _ => unreachable!("kind is taken modulo IETF_KINDS"), } } diff --git a/rs/moq-net/src/ietf/datagram.rs b/rs/moq-net/src/ietf/datagram.rs new file mode 100644 index 0000000000..fae9b48bdf --- /dev/null +++ b/rs/moq-net/src/ietf/datagram.rs @@ -0,0 +1,301 @@ +//! OBJECT_DATAGRAM: one Object carried in a QUIC datagram (draft-14 section 10.3.1 through +//! draft-20 section 11.3.1). +//! +//! The model counterpart is [`crate::Datagram`], a single-frame group, so only an Object at +//! ID 0 maps onto it. The Type is a set of flags on every draft; draft-14 lacks the +//! DEFAULT_PRIORITY bit and a status with an omitted Object ID. + +use bytes::{Buf, BufMut, Bytes}; + +use crate::coding::{Decode, DecodeError, Encode, EncodeError}; + +use super::Version; + +/// The bits of an OBJECT_DATAGRAM Type. +mod flag { + pub const PROPERTIES: u64 = 0x01; + pub const END_OF_GROUP: u64 = 0x02; + pub const ZERO_OBJECT_ID: u64 = 0x04; + pub const DEFAULT_PRIORITY: u64 = 0x08; + pub const STATUS: u64 = 0x20; + /// Every defined bit. Anything else, including the reserved 0x10, is invalid. + pub const ALL: u64 = PROPERTIES | END_OF_GROUP | ZERO_OBJECT_ID | DEFAULT_PRIORITY | STATUS; +} + +/// What follows an OBJECT_DATAGRAM's header. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum DatagramBody { + /// The Object Payload, delimited by the datagram boundary. + Payload(Bytes), + /// The Object Status of an Object without a payload. + Status(u64), +} + +/// A decoded OBJECT_DATAGRAM. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ObjectDatagram { + pub track_alias: u64, + pub group_id: u64, + /// The Object ID, or `None` when the ZERO_OBJECT_ID bit omits it (Object 0). + pub object_id: Option, + /// The Publisher Priority, or `None` to inherit the subscription's (draft-15+). + pub publisher_priority: Option, + /// No Object past this one exists in the group. + pub end_of_group: bool, + /// The Object Properties block without its length prefix, which carries the Timestamp. + pub properties: Option>, + pub body: DatagramBody, +} + +impl ObjectDatagram { + /// Whether `kind` is a Type this draft defines. + fn valid(kind: u64, version: Version) -> bool { + match version { + Version::Draft14 => kind <= 0x07 || kind == 0x20 || kind == 0x21, + // A status cannot also end the group. + _ => kind & !flag::ALL == 0 && !(kind & flag::STATUS != 0 && kind & flag::END_OF_GROUP != 0), + } + } +} + +impl Encode for ObjectDatagram { + fn encode(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { + let mut kind = 0; + if self.properties.is_some() { + kind |= flag::PROPERTIES; + } + if self.end_of_group { + kind |= flag::END_OF_GROUP; + } + if self.object_id.is_none() { + kind |= flag::ZERO_OBJECT_ID; + } + if self.publisher_priority.is_none() { + kind |= flag::DEFAULT_PRIORITY; + } + if matches!(self.body, DatagramBody::Status(_)) { + kind |= flag::STATUS; + } + if !Self::valid(kind, version) { + return Err(EncodeError::InvalidState); + } + + kind.encode(w, version)?; + self.track_alias.encode(w, version)?; + self.group_id.encode(w, version)?; + if let Some(object_id) = self.object_id { + object_id.encode(w, version)?; + } + if let Some(priority) = self.publisher_priority { + priority.encode(w, version)?; + } + if let Some(properties) = &self.properties { + // A present but empty block is a protocol violation for the peer. + if properties.is_empty() { + return Err(EncodeError::InvalidState); + } + properties.encode(w, version)?; + } + + match &self.body { + DatagramBody::Status(status) => { + // Draft-17 on: only a Normal Object may carry Properties. + let legacy = matches!(version, Version::Draft14 | Version::Draft15 | Version::Draft16); + if !legacy && *status != 0 && self.properties.is_some() { + return Err(EncodeError::InvalidState); + } + status.encode(w, version)? + } + DatagramBody::Payload(payload) => { + // Runs to the datagram boundary: written raw, no length prefix. + if w.remaining_mut() < payload.len() { + return Err(EncodeError::Short); + } + w.put_slice(payload); + } + } + Ok(()) + } +} + +impl Decode for ObjectDatagram { + fn decode(r: &mut R, version: Version) -> Result { + let kind = u64::decode(r, version)?; + if !Self::valid(kind, version) { + return Err(DecodeError::InvalidValue); + } + + let track_alias = u64::decode(r, version)?; + let group_id = u64::decode(r, version)?; + let object_id = match kind & flag::ZERO_OBJECT_ID != 0 { + true => None, + false => Some(u64::decode(r, version)?), + }; + let publisher_priority = match kind & flag::DEFAULT_PRIORITY != 0 { + true => None, + false => Some(u8::decode(r, version)?), + }; + let properties = match kind & flag::PROPERTIES != 0 { + true => { + let properties = Vec::::decode(r, version)?; + if properties.is_empty() { + return Err(DecodeError::InvalidValue); + } + Some(properties) + } + false => None, + }; + + let body = match kind & flag::STATUS != 0 { + true => { + let status = u64::decode(r, version)?; + // Draft-17 on: only a Normal Object may carry Properties. + let legacy = matches!(version, Version::Draft14 | Version::Draft15 | Version::Draft16); + if !legacy && status != 0 && properties.is_some() { + return Err(DecodeError::InvalidValue); + } + if r.has_remaining() { + return Err(DecodeError::TrailingBytes); + } + DatagramBody::Status(status) + } + false => DatagramBody::Payload(r.copy_to_bytes(r.remaining())), + }; + + Ok(Self { + track_alias, + group_id, + object_id, + publisher_priority, + end_of_group: kind & flag::END_OF_GROUP != 0, + properties, + body, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const ALL: [Version; 9] = [ + Version::Draft14, + Version::Draft15, + Version::Draft16, + Version::Draft17, + Version::Draft18, + Version::Draft19, + Version::Draft20, + Version::Draft21, + Version::Draft22, + ]; + + fn single(priority: Option) -> ObjectDatagram { + ObjectDatagram { + track_alias: 3, + group_id: 42, + object_id: None, + publisher_priority: priority, + end_of_group: true, + properties: Some(vec![0x10, 0x05]), + body: DatagramBody::Payload(Bytes::from_static(b"hello")), + } + } + + #[test] + fn roundtrip_every_draft() { + for version in ALL { + let datagram = single(Some(7)); + let mut buf = datagram.encode_bytes(version).unwrap(); + assert_eq!(buf[0], 0x07, "{version}: properties, end of group, object 0"); + let decoded = ObjectDatagram::decode(&mut buf, version).unwrap(); + assert_eq!(decoded, datagram, "{version}"); + assert!(!buf.has_remaining()); + } + } + + #[test] + fn default_priority_needs_draft15() { + assert!(matches!( + single(None).encode_bytes(Version::Draft14), + Err(EncodeError::InvalidState) + )); + assert!(ObjectDatagram::decode(&mut &[0x0C, 0x01, 0x02][..], Version::Draft14).is_err()); + + let mut buf = single(None).encode_bytes(Version::Draft16).unwrap(); + assert_eq!(buf[0], 0x0F); + assert_eq!( + ObjectDatagram::decode(&mut buf, Version::Draft16).unwrap(), + single(None) + ); + } + + #[test] + fn explicit_object_id_and_status() { + let datagram = ObjectDatagram { + track_alias: 1, + group_id: 2, + object_id: Some(0), + publisher_priority: Some(128), + end_of_group: false, + properties: None, + body: DatagramBody::Status(0), + }; + for version in ALL { + let mut buf = datagram.encode_bytes(version).unwrap(); + assert_eq!(buf[0], 0x20, "{version}"); + assert_eq!(ObjectDatagram::decode(&mut buf, version).unwrap(), datagram); + } + } + + #[test] + fn rejects_invalid_types() { + for version in ALL { + // A status that ends the group, the reserved bit, and a bit past the defined ones. + for kind in [0x22u64, 0x10, 0x40] { + let mut bytes = kind.encode_bytes(version).unwrap().to_vec(); + bytes.extend_from_slice(&[0x01, 0x02, 0x03, 0x04]); + assert!( + ObjectDatagram::decode(&mut &bytes[..], version).is_err(), + "{version}: {kind:#x}" + ); + } + } + } + + #[test] + fn status_with_properties_needs_normal() { + let datagram = ObjectDatagram { + track_alias: 1, + group_id: 2, + object_id: Some(0), + publisher_priority: Some(0), + end_of_group: false, + properties: Some(vec![0x10, 0x05]), + body: DatagramBody::Status(3), + }; + for version in ALL { + let legacy = matches!(version, Version::Draft14 | Version::Draft15 | Version::Draft16); + match datagram.encode_bytes(version) { + Ok(mut buf) => { + assert!(legacy, "{version}"); + assert_eq!(ObjectDatagram::decode(&mut buf, version).unwrap(), datagram); + } + Err(err) => { + assert!(!legacy, "{version}"); + assert!(matches!(err, EncodeError::InvalidState)); + } + } + } + } + + #[test] + fn rejects_empty_properties() { + // PROPERTIES and ZERO_OBJECT_ID, alias 1, group 2, priority 0, empty block. + let bytes = [0x05, 0x01, 0x02, 0x00, 0x00, b'x']; + assert!(matches!( + ObjectDatagram::decode(&mut &bytes[..], Version::Draft16), + Err(DecodeError::InvalidValue) + )); + } +} diff --git a/rs/moq-net/src/ietf/fetch.rs b/rs/moq-net/src/ietf/fetch.rs index 5987d8efbc..57e7d07fe3 100644 --- a/rs/moq-net/src/ietf/fetch.rs +++ b/rs/moq-net/src/ietf/fetch.rs @@ -167,7 +167,8 @@ impl Message for Fetch<'_> { ); let subscriber_priority = subscriber_priority.unwrap_or(128); - let group_order = group_order.unwrap_or(GroupOrder::Descending); + // No preference: the publisher picks the order. + let group_order = group_order.unwrap_or(GroupOrder::Any); Ok(Self { request_id, diff --git a/rs/moq-net/src/ietf/mod.rs b/rs/moq-net/src/ietf/mod.rs index 9c8577b340..944c277893 100644 --- a/rs/moq-net/src/ietf/mod.rs +++ b/rs/moq-net/src/ietf/mod.rs @@ -9,6 +9,7 @@ mod parameters; mod adapter; pub mod cluster; mod control; +mod datagram; pub(crate) mod error; mod fetch; mod filter; @@ -35,6 +36,7 @@ mod track; mod version; use control::Control; +pub use datagram::*; pub use fetch::*; pub use filter::*; pub use goaway::*; diff --git a/rs/moq-net/src/ietf/publisher.rs b/rs/moq-net/src/ietf/publisher.rs index b6bbe92179..5de4a889e8 100644 --- a/rs/moq-net/src/ietf/publisher.rs +++ b/rs/moq-net/src/ietf/publisher.rs @@ -15,7 +15,7 @@ use web_transport_trait::poll::SendStream as _; use crate::{ AsPath, Error, Timescale, Timestamp, - coding::{Stream, Writer}, + coding::{Encode as _, Stream, Writer}, ietf::{self, Control, EndLocation, FetchHeader, FetchType, Filter, GroupOrder, Location, RequestId}, track::Subscription, util::{MaybeBoxedExt, MaybeSendBox}, @@ -44,6 +44,83 @@ enum FillStep { Done, } +/// Where a fetch object sits relative to the one written before it on the same stream. +#[derive(Clone, Copy)] +enum FetchPrior { + /// The first object on the stream. + None, + /// The next object of the same group. + Same, +} + +/// The group a FETCH answers with, read out of the cache so a later eviction cannot +/// truncate what FETCH_OK already promised. +struct FetchedGroup { + sequence: u64, + /// The Object ID of the first frame. + first: u64, + frames: Vec, + /// The group ended within the range, so every frame it will ever hold was read. + complete: bool, +} + +impl FetchPrior { + /// The prior for the next object of a single-group stream, clearing `first`. + fn next(first: &mut bool) -> Self { + match std::mem::take(first) { + true => Self::None, + false => Self::Same, + } + } +} + +impl FetchedGroup { + /// One past the last Object ID read. + fn end(&self) -> u64 { + self.first + self.frames.len() as u64 + } +} + +/// Read group `sequence` of `track` for a FETCH, from object `skip` up to the exclusive +/// `until`, or through the end of the group without one. +/// +/// This is a [`track::Consumer::fetch_group`], so a relay fetches a miss upstream. +async fn read_fetch( + track: &track::Consumer, + sequence: u64, + skip: u64, + until: Option, + priority: u8, +) -> Result { + let fetch = group::Fetch { + priority, + frame_start: skip, + }; + let mut group = track.fetch_group(sequence, fetch).await?; + + // `fetch_group` positions the consumer at `skip`, or refuses a group that no + // longer holds it. + let first = group.index(); + let mut frames = Vec::new(); + let mut complete = false; + while until.is_none_or(|until| first + (frames.len() as u64) < until) { + match group.read_frame().await? { + Some(frame) => frames.push(frame), + None => { + complete = true; + break; + } + } + } + + Ok(FetchedGroup { + sequence, + first, + frames, + complete, + }) +} + /// A broadcast whose route table is watched for changes in what we advertise: the /// namespace becoming (un)advertisable, or its path or cost moving. struct Watched { @@ -877,20 +954,13 @@ where stream, sequence, index, - std::mem::take(&mut first), + FetchPrior::next(&mut first), frame.timestamp, timescale, version, ) .await?; - stream.encode(&(frame.payload.len() as u64)).await?; - if frame.payload.is_empty() && matches!(version, Version::Draft14 | Version::Draft15) { - stream.encode(&0u64).await?; - } - if !frame.payload.is_empty() { - let mut payload = frame.payload; - stream.write_all(&mut payload).await?; - } + Self::write_fetch_payload(stream, frame.payload, version).await?; } index += 1; group.keep_alive(); @@ -922,7 +992,7 @@ where stream, sequence, index, - std::mem::take(&mut first), + FetchPrior::next(&mut first), frame.timestamp, timescale, version, @@ -968,7 +1038,7 @@ where stream: &mut Writer, sequence: u64, object: u64, - first: bool, + prior: FetchPrior, timestamp: Timestamp, timescale: Option, version: Version, @@ -993,10 +1063,10 @@ where return Ok(()); } - let header = match first { + let header = match prior { // The first object must carry its absolute Group and Object IDs. Include the // priority too: "same as the prior object" has no prior to refer to. - true => ietf::FetchObject::Object { + FetchPrior::None => ietf::FetchObject::Object { subgroup: ietf::FetchSubgroup::Zero, group: Some(sequence), object: Some(object), @@ -1004,7 +1074,7 @@ where properties, }, // Same group and priority; the Object ID is the prior one plus one. - false => ietf::FetchObject::Object { + FetchPrior::Same => ietf::FetchObject::Object { subgroup: ietf::FetchSubgroup::Zero, group: None, object: None, @@ -1018,6 +1088,24 @@ where Ok(()) } + /// Write a whole fetch object's length and payload, after its header. + async fn write_fetch_payload( + stream: &mut Writer, + payload: bytes::Bytes, + version: Version, + ) -> Result<(), Error> { + stream.encode(&(payload.len() as u64)).await?; + // Draft-14 and 15 still carry a Normal status after an empty object. + if payload.is_empty() && matches!(version, Version::Draft14 | Version::Draft15) { + stream.encode(&0u64).await?; + } + if !payload.is_empty() { + let mut payload = payload; + stream.write_all(&mut payload).await?; + } + Ok(()) + } + /// Register a pending subscription until the returned guard drops. fn register_join(&self, request_id: RequestId) -> Join { self.joins.lock().insert(request_id, None); @@ -1027,166 +1115,160 @@ where } } - /// Serve the current-group prefix for a relative joining FETCH with offset zero. + /// Answer a FETCH within one group: a standalone range of the named track, or a + /// joining FETCH's prefix of its subscription's group. + /// + /// The answer is buffered before replying: FETCH_OK names where the response ends, + /// which a range running past the track only learns by reading it, and a refusal can + /// still replace it until then. async fn run_fetch_stream(mut self, mut stream: Stream, msg: ietf::Fetch<'_>) -> Result<(), Error> { + let priority = super::priority::from_wire(msg.subscriber_priority); + + // Draft-20 moved the range into LOCATION_FILTER and replaced joining FETCH with + // subscription fills; this still reads the older layout. if Filter::is_draft20(self.version) { return self .reject_fetch( stream, msg.request_id, &Error::Unsupported, - "joining FETCH removed in draft-20", + "FETCH not supported on draft-20", ) .await; } - let subscribe_id = match msg.fetch_type { - FetchType::Standalone { .. } => { - return self - .reject_fetch(stream, msg.request_id, &Error::Unsupported, "not supported") - .await; - } - FetchType::RelativeJoining { - subscriber_request_id, - group_offset, + let (track, start, end, timescale, joined) = match msg.fetch_type { + FetchType::Standalone { + namespace, + track, + start, + end, } => { - if group_offset != 0 { + // An End Object of 0 asks for the whole End Group. + let end = match end.object { + 0 => end.group.checked_add(1).map(|group| Location { group, object: 0 }), + _ => Some(end), + }; + let Some(end) = end.filter(|end| (start.group, start.object) < (end.group, end.object)) else { return self - .reject_fetch(stream, msg.request_id, &Error::Unsupported, "not supported") + .reject_fetch(stream, msg.request_id, &Error::InvalidRange, "empty range") .await; - } - subscriber_request_id + }; + + // The peer must have seen the announcement to name this namespace, so this + // resolves like a SUBSCRIBE does. + let broadcast = match self.serving_origin().await.request_broadcast(&namespace).await { + Ok(broadcast) => broadcast, + Err(err) => return self.reject_fetch(stream, msg.request_id, &err, &err.to_string()).await, + }; + let track = match broadcast.track(&track) { + Ok(track) => track, + Err(err) => return self.reject_fetch(stream, msg.request_id, &err, &err.to_string()).await, + }; + + // No SUBSCRIBE declared a timescale for this request, so its objects go + // out unstamped. + (track, start, end, None, false) } - FetchType::AbsoluteJoining { .. } => { - return self - .reject_fetch(stream, msg.request_id, &Error::Unsupported, "not supported") - .await; + FetchType::RelativeJoining { + subscriber_request_id, .. + } + | FetchType::AbsoluteJoining { + subscriber_request_id, .. + } => { + let (end, cache, timescale) = match self.joined(&mut stream, subscriber_request_id).await? { + Ok(joined) => joined, + Err((err, reason)) => return self.reject_fetch(stream, msg.request_id, &err, reason).await, + }; + let start = match msg.fetch_type { + FetchType::RelativeJoining { group_offset, .. } => end.group.saturating_sub(group_offset), + FetchType::AbsoluteJoining { group_id, .. } if group_id <= end.group => group_id, + _ => { + return self + .reject_fetch( + stream, + msg.request_id, + &Error::InvalidRange, + "joining group past the subscription", + ) + .await; + } + }; + ( + cache, + Location { + group: start, + object: 0, + }, + end, + timescale, + true, + ) } }; - // Request streams can arrive out of order. Wait on registration, while bounding - // the lifetime of a request whose subscription never arrives or resolves. - let joined = { - let mut pending = false; - let mut deadline = crate::runtime::Deadline::after(&self.runtime, Duration::from_secs(10)); + // One group per FETCH: on a relay, a range costs a serial upstream fetch per missing + // group, all buffered until FETCH_OK. Ranges wait for upstream fills by range. + let last = match end.object { + 0 => end.group - 1, + _ => end.group, + }; + if start.group != last { + return self + .reject_fetch( + stream, + msg.request_id, + &Error::Unsupported, + "FETCH spanning several groups not supported", + ) + .await; + } + let until = (end.object > 0).then_some(end.object); + + // The subscriber cancelling is the only other way this ends early, and is owed + // nothing. + let group = { + let mut read = std::pin::pin!(read_fetch(&track, start.group, start.object, until, priority)); kio::wait(|waiter| { let mut cx = std::task::Context::from_waker(waiter.waker()); - // The request reader is what the subscriber FINs or resets. The writer - // on a draft-14-16 virtual stream reports closed immediately, which is - // not a cancellation. if stream.reader.poll_closed(&mut cx).is_ready() { - return Poll::Ready(Err(Error::Cancel)); + return Poll::Ready(None); } - let joins = self.joins.poll(waiter, |joins| match joins.get(&subscribe_id) { - Some(Some(_)) => Poll::Ready(()), - Some(None) => { - pending = true; - Poll::Pending - } - None if pending => Poll::Ready(()), - None => Poll::Pending, - }); - if let Poll::Ready(joins) = joins { - return Poll::Ready(Ok(joins.get(&subscribe_id).cloned().flatten())); - } - if deadline.poll(waiter).is_ready() { - return Poll::Ready(if pending { Err(Error::Timeout) } else { Ok(None) }); - } - Poll::Pending + waiter.poll_future(read.as_mut()).map(Some) }) .await }; - let joined = match joined { - Err(Error::Timeout) => { - return self - .reject_fetch(stream, msg.request_id, &Error::Timeout, "subscription not ready") - .await; - } - result => result?, + let group = match group { + Some(Ok(group)) => group, + Some(Err(err)) => return self.reject_fetch(stream, msg.request_id, &err, &err.to_string()).await, + None => return Ok(()), }; - let (end, cache, timescale) = match joined { - None => { - return self - .reject_fetch( - stream, - msg.request_id, - &if matches!( - self.version, - Version::Draft14 - | Version::Draft15 | Version::Draft16 - | Version::Draft17 | Version::Draft18 - | Version::Draft19 - ) { - Error::InvalidJoiningRequestId - } else { - Error::NotFound - }, - "no such subscription", - ) - .await; - } - Some(Joined::Unsupported) => { - if matches!(self.version, Version::Draft14 | Version::Draft15 | Version::Draft16) { - self.session.close( - crate::SessionError::ProtocolViolation.to_code(), - "joining FETCH requires Largest Object filter", - ); - return Err(Error::ProtocolViolation); - } - return self - .reject_fetch( - stream, - msg.request_id, - &Error::Unsupported, - "joining filter not supported", - ) - .await; - } - Some(Joined::Empty) => { + + let (end_location, end_of_track) = if joined { + // The subscription starts at the saved Largest Object, so the prefix of that + // group has to be here in full or the two leave a gap. + if group.end() != end.object { return self - .reject_fetch( - stream, - msg.request_id, - &Error::InvalidRange, - "no objects at subscription start", - ) + .reject_fetch(stream, msg.request_id, &Error::Evicted, "joining prefix unavailable") .await; } - Some(Joined::Group { end, cache, timescale }) => (end, cache, timescale), - }; - let priority = super::priority::from_wire(msg.subscriber_priority); - let mut group = match cache - .fetch_group( - end.group, - group::Fetch { - priority, - ..Default::default() - }, - ) - .await - { - Ok(group) => group, - Err(err) => { - return self - .reject_fetch(stream, msg.request_id, &err, "joining group unavailable") - .await; + (end, false) + } else { + let end_of_track = group.complete && track.final_sequence() == group.sequence.checked_add(1); + // A group that ends the track ends the response at its last object; otherwise + // the response ends where it was asked to. + match end_of_track { + true => ( + Location { + group: group.sequence, + object: group.end(), + }, + true, + ), + false => (end, false), } }; - // Retain the promised prefix before success: a cached group may already have - // evicted its front, and later cache eviction must not truncate this FETCH. - let mut prefix = Vec::new(); - for _ in 0..end.object { - match group.read_frame().await { - Ok(Some(frame)) => prefix.push(frame), - _ => { - return self - .reject_fetch(stream, msg.request_id, &Error::Evicted, "joining prefix unavailable") - .await; - } - } - } - // FETCH_OK on every draft, never REQUEST_OK: section 5.2 allows exactly one FETCH_OK or // REQUEST_ERROR in answer to a FETCH, and REQUEST_OK's own definition lists the other // requests it answers without ever naming this one. @@ -1199,13 +1281,15 @@ where _ => None, }, // Only draft-14 encodes it, and only as the publisher restating the order. - group_order: msg.group_order.any_to_descending(), - end_of_track: false, - end_location: end, + group_order: match msg.group_order { + GroupOrder::Descending => GroupOrder::Descending, + _ => GroupOrder::Ascending, + }, + end_of_track, + end_location, }) .await?; - // The FETCH owns the retained prefix, capped at the subscription snapshot. let uni = self.session.open_uni().await.map_err(Error::from_transport)?; let mut writer = Writer::new(uni, self.version); writer.set_priority(priority); @@ -1215,25 +1299,19 @@ where request_id: msg.request_id, }) .await?; - for (index, frame) in prefix.into_iter().enumerate() { + let mut first = true; + for (object, frame) in (group.first..).zip(group.frames) { Self::write_fetch_object( &mut writer, - end.group, - index as u64, - index == 0, + group.sequence, + object, + FetchPrior::next(&mut first), frame.timestamp, timescale, self.version, ) .await?; - writer.encode(&(frame.payload.len() as u64)).await?; - if frame.payload.is_empty() && matches!(self.version, Version::Draft14 | Version::Draft15) { - writer.encode(&0u64).await?; - } - if !frame.payload.is_empty() { - let mut payload = frame.payload; - writer.write_all(&mut payload).await?; - } + Self::write_fetch_payload(&mut writer, frame.payload, self.version).await?; } writer.close().await?; @@ -1245,6 +1323,76 @@ where Ok(()) } + /// Resolve the subscription a joining FETCH names to its saved start, or the refusal + /// to answer the FETCH with. + async fn joined( + &mut self, + stream: &mut Stream, + subscribe_id: RequestId, + ) -> Result), (Error, &'static str)>, Error> { + // Request streams can arrive out of order. Wait on registration, while bounding + // the lifetime of a request whose subscription never arrives or resolves. + let joined = { + let mut pending = false; + let mut deadline = crate::runtime::Deadline::after(&self.runtime, Duration::from_secs(10)); + kio::wait(|waiter| { + let mut cx = std::task::Context::from_waker(waiter.waker()); + // The request reader is what the subscriber FINs or resets. The writer + // on a draft-14-16 virtual stream reports closed immediately, which is + // not a cancellation. + if stream.reader.poll_closed(&mut cx).is_ready() { + return Poll::Ready(Err(Error::Cancel)); + } + let joins = self.joins.poll(waiter, |joins| match joins.get(&subscribe_id) { + Some(Some(_)) => Poll::Ready(()), + Some(None) => { + pending = true; + Poll::Pending + } + None if pending => Poll::Ready(()), + None => Poll::Pending, + }); + if let Poll::Ready(joins) = joins { + return Poll::Ready(Ok(joins.get(&subscribe_id).cloned().flatten())); + } + if deadline.poll(waiter).is_ready() { + return Poll::Ready(if pending { Err(Error::Timeout) } else { Ok(None) }); + } + Poll::Pending + }) + .await + }; + let refusal = match joined { + Err(Error::Timeout) => (Error::Timeout, "subscription not ready"), + Err(err) => return Err(err), + Ok(Some(Joined::Group { end, cache, timescale })) => return Ok(Ok((end, cache, timescale))), + Ok(None) => ( + match self.version { + Version::Draft14 + | Version::Draft15 + | Version::Draft16 + | Version::Draft17 + | Version::Draft18 + | Version::Draft19 => Error::InvalidJoiningRequestId, + _ => Error::NotFound, + }, + "no such subscription", + ), + Ok(Some(Joined::Unsupported)) => { + if matches!(self.version, Version::Draft14 | Version::Draft15 | Version::Draft16) { + self.session.close( + crate::SessionError::ProtocolViolation.to_code(), + "joining FETCH requires Largest Object filter", + ); + return Err(Error::ProtocolViolation); + } + (Error::Unsupported, "joining filter not supported") + } + Ok(Some(Joined::Empty)) => (Error::InvalidRange, "no objects at subscription start"), + }; + Ok(Err(refusal)) + } + async fn reject_track_status(&self, mut stream: Stream, request_id: RequestId) -> Result<(), Error> { let error_code = request::to_code(&Error::Unsupported, request::Kind::TrackStatus, self.version); if self.version == Version::Draft14 { @@ -1968,6 +2116,9 @@ struct TrackServe { opened: Arc, /// The track's exclusive end, once its groups ran out because it finished. end: Option, + /// Serve the track's datagrams too, as OBJECT_DATAGRAMs. Off when the transport has no + /// datagrams; there is no stream fallback. + datagrams: bool, } impl TrackServe { @@ -1988,8 +2139,10 @@ impl TrackServe { } } track.end_at(range.end.map_or(Bound::Unbounded, |end| Bound::Included(end.group))); + let datagrams = session.max_datagram_size() > 0; Self { + datagrams, session, track, request_id, @@ -2103,6 +2256,8 @@ impl TrackServe { ); } Poll::Ready(Ok(None)) => { + // Datagrams written before the track finished still go out. + self.poll_datagrams(waiter); self.draining = true; if let Poll::Ready(Ok(end)) = self.track.poll_finished(waiter) { self.end = Some(end); @@ -2113,10 +2268,59 @@ impl TrackServe { Poll::Pending => break, } } + // Groups first, so a burst of datagrams cannot starve them. + self.poll_datagrams(waiter); // Newly created group machines start now rather than on the next wake. let _ = self.children.poll(waiter); Poll::Pending } + + /// Send every buffered datagram as an OBJECT_DATAGRAM, best-effort, like moq-lite. + /// + /// A datagram is Object 0 of a group that has no other, so the Object ends its group. + /// The track's end or failure surfaces through its groups, so this only stops. + fn poll_datagrams(&mut self, waiter: &kio::Waiter) { + if !self.datagrams { + return; + } + while let Poll::Ready(Ok(Some(datagram))) = self.track.poll_recv_datagram(waiter) { + let sequence = datagram.sequence; + let properties = match self.timescale { + Some(timescale) => { + let mut properties = Vec::new(); + if ietf::encode_object_time(&mut properties, datagram.timestamp, timescale, self.version).is_err() { + continue; + } + Some(properties) + } + None => None, + }; + let body = ietf::ObjectDatagram { + track_alias: self.request_id.0, + group_id: sequence, + object_id: None, + publisher_priority: Some(super::priority::to_wire(self.track.info().priority)), + end_of_group: true, + properties, + body: ietf::DatagramBody::Payload(datagram.payload), + }; + let Ok(body) = body.encode_bytes(self.version) else { + continue; + }; + + let max = self.session.max_datagram_size(); + if body.len() > max { + tracing::debug!( + sequence, + size = body.len(), + max, + "dropping datagram larger than the transport limit" + ); + continue; + } + let _ = self.session.send_datagram(&body); + } + } } /// Serves one group on its own unidirectional stream in the moq-transport @@ -2496,10 +2700,31 @@ mod group_priority_test { /// every moq-transport peer. #[tokio::test] async fn group_header_carries_the_publisher_priority() { + let header = serve_group_header(track::Info::default().with_priority(hang_audio_priority())).await; + assert_eq!( + header.publisher_priority, + priority::to_wire(hang_audio_priority()), + "the wire is lower-first, so audio must encode below video" + ); + assert!( + priority::to_wire(hang_audio_priority()) < priority::to_wire(hang_video_priority()), + "audio outranks video on the wire" + ); + } + + /// A track that never set a priority is the draft's usual publisher priority, 128, + /// not 255, the least urgent value a peer like moxygen would deprioritize. + #[tokio::test] + async fn group_header_defaults_to_the_midpoint() { + let header = serve_group_header(track::Info::default()).await; + assert_eq!(header.publisher_priority, 128); + } + + /// Serve one group of a track with `info` and decode the subgroup header it opens with. + async fn serve_group_header(info: track::Info) -> ietf::GroupHeader { let log = crate::lite::test_transport::Log::default(); let session = SinkSession::new(log.clone()); - let info = track::Info::default().with_priority(hang_audio_priority()); let track = track::Producer::new(std::sync::Arc::new(crate::broadcast::Info::default()), "test", info); let subscriber = track.subscribe(None); @@ -2520,16 +2745,7 @@ mod group_priority_test { let written = log.writes.lock().unwrap().clone(); let mut buf = bytes::Bytes::from(written); - let header = ietf::GroupHeader::decode(&mut buf, Version::Draft14).expect("a group header"); - assert_eq!( - header.publisher_priority, - priority::to_wire(hang_audio_priority()), - "the wire is lower-first, so audio must encode below video" - ); - assert!( - priority::to_wire(hang_audio_priority()) < priority::to_wire(hang_video_priority()), - "audio outranks video on the wire" - ); + ietf::GroupHeader::decode(&mut buf, Version::Draft14).expect("a group header") } /// `hang::catalog::PRIORITY` isn't reachable from `moq-net` (hang depends on it, not the @@ -3512,6 +3728,337 @@ mod serve_tests { assert!(buf.is_empty()); } + /// Every draft that carries a standalone FETCH. + const FETCH_DRAFTS: [Version; 6] = [ + Version::Draft14, + Version::Draft15, + Version::Draft16, + Version::Draft17, + Version::Draft18, + Version::Draft19, + ]; + + /// Groups `0..count`, each holding `g-0` and `g-1`, skipping `hole`. + fn publish_pairs(h: &mut Serve, count: u64, hole: Option) { + for sequence in (0..count).filter(|sequence| Some(*sequence) != hole) { + let mut group = h.track.create_group(group::Info { sequence }).unwrap(); + for object in 0..2 { + group + .write_frame(timestamp(), format!("{sequence}-{object}").into_bytes()) + .unwrap(); + } + group.finish().unwrap(); + } + } + + /// Run a standalone FETCH of `room/video`, returning what the peer reads back. + async fn standalone_fetch(h: &Serve, start: Location, end: Location, group_order: GroupOrder) -> bytes::Bytes { + let version = h.publisher.version; + let mark = h.log.writes.lock().unwrap().len(); + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + h.publisher + .clone() + .run_fetch_stream( + stream, + ietf::Fetch { + request_id: FETCH_ID, + subscriber_priority: 128, + group_order, + fetch_type: FetchType::Standalone { + namespace: crate::Path::new("room"), + track: "video".into(), + start, + end, + }, + }, + ) + .await + .unwrap(); + bytes::Bytes::from(h.log.writes.lock().unwrap()[mark..].to_vec()) + } + + /// Decode a FETCH_OK and the fetch stream after it, as `(group, object, payload)`. + fn fetch_answer(mut buf: bytes::Bytes, version: Version) -> (ietf::FetchOk, Vec<(u64, u64, String)>) { + assert_eq!( + u64::decode(&mut buf, version).unwrap(), + ietf::FetchOk::ID, + "{version}: not a FETCH_OK" + ); + let ok = ietf::FetchOk::decode(&mut buf, version).unwrap(); + assert_eq!(u64::decode(&mut buf, version).unwrap(), FetchHeader::TYPE); + assert_eq!(FetchHeader::decode(&mut buf, version).unwrap().request_id, FETCH_ID); + + let mut objects = Vec::new(); + let mut prior: Option<(u64, u64)> = None; + while !buf.is_empty() { + let (group, object) = if version == Version::Draft14 { + let group = u64::decode(&mut buf, version).unwrap(); + assert_eq!(u64::decode(&mut buf, version).unwrap(), 0, "subgroup"); + let object = u64::decode(&mut buf, version).unwrap(); + let _priority = u8::decode(&mut buf, version).unwrap(); + let _properties = Vec::::decode(&mut buf, version).unwrap(); + (group, object) + } else { + let ietf::FetchObject::Object { group, object, .. } = + ietf::FetchObject::decode(&mut buf, version).unwrap() + else { + panic!("{version}: unexpected End of Range"); + }; + match (prior, group) { + (None, group) => (group.unwrap(), object.unwrap()), + // No Group ID: the same group, and an absent Object ID Delta is one. + (Some((group, prior)), None) => (group, prior + object.unwrap_or(1)), + (Some((prior, _)), Some(group)) => { + let group = match version { + Version::Draft15 | Version::Draft16 | Version::Draft17 => group, + _ => prior + group + 1, + }; + (group, object.unwrap()) + } + } + }; + let size = u64::decode(&mut buf, version).unwrap() as usize; + let payload = String::from_utf8(buf.split_to(size).to_vec()).unwrap(); + objects.push((group, object, payload)); + prior = Some((group, object)); + } + (ok, objects) + } + + /// The `(group, object, payload)` a range of two-frame groups holds. + fn pairs(groups: impl IntoIterator) -> Vec<(u64, u64, String)> { + groups + .into_iter() + .flat_map(|group| (0..2).map(move |object| (group, object, format!("{group}-{object}")))) + .collect() + } + + /// A whole group is answered on one fetch stream. + #[tokio::test] + async fn a_standalone_fetch_serves_one_whole_group() { + for version in FETCH_DRAFTS { + let mut h = serve(version); + publish_pairs(&mut h, 5, None); + settle().await; + + // An End Object of 0 asks for the whole End Group. + let buf = standalone_fetch( + &h, + Location { group: 2, object: 0 }, + Location { group: 2, object: 0 }, + GroupOrder::Ascending, + ) + .await; + let (ok, objects) = fetch_answer(buf, version); + assert_eq!(ok.end_location, Location { group: 3, object: 0 }, "{version}"); + assert!(!ok.end_of_track, "{version}"); + assert_eq!(objects, pairs([2]), "{version}"); + assert!(h.log.resets().is_empty(), "{version}"); + } + } + + /// The last group of a finished track ends the response one past its last object, + /// with End of Track set. + #[tokio::test] + async fn a_standalone_fetch_of_the_last_group_reports_the_end_of_track() { + for version in FETCH_DRAFTS { + let mut h = serve(version); + publish_pairs(&mut h, 5, None); + h.track.finish().unwrap(); + settle().await; + + let buf = standalone_fetch( + &h, + Location { group: 4, object: 1 }, + Location { group: 4, object: 0 }, + GroupOrder::Any, + ) + .await; + let (ok, objects) = fetch_answer(buf, version); + assert_eq!(ok.end_location, Location { group: 4, object: 2 }, "{version}"); + assert!(ok.end_of_track, "{version}"); + assert_eq!(objects, pairs([4])[1..].to_vec(), "{version}"); + } + } + + /// Draft-20 carries the range in LOCATION_FILTER, which is not read yet. + #[tokio::test] + async fn a_draft20_fetch_is_refused() { + for version in [Version::Draft20, Version::Draft21, Version::Draft22] { + let mut h = serve(version); + publish_pairs(&mut h, 3, None); + settle().await; + + let buf = standalone_fetch( + &h, + Location { group: 1, object: 0 }, + Location { group: 1, object: 0 }, + GroupOrder::Ascending, + ) + .await; + assert_eq!(fetch_refusal(buf, version), 0x3, "{version}"); + } + } + + /// A range touching several groups is refused, in either order: on a relay each + /// missing group would be its own upstream fetch. + #[tokio::test] + async fn a_standalone_fetch_of_several_groups_is_refused() { + for version in FETCH_DRAFTS { + for order in [GroupOrder::Ascending, GroupOrder::Descending] { + let mut h = serve(version); + publish_pairs(&mut h, 3, None); + settle().await; + + let buf = standalone_fetch( + &h, + Location { group: 0, object: 1 }, + Location { group: 1, object: 1 }, + order, + ) + .await; + assert_eq!(fetch_refusal(buf, version), 0x3, "{version}"); + } + } + } + + /// Read a REQUEST_ERROR (or draft-14 FETCH_ERROR) off a refused FETCH, as its code. + fn fetch_refusal(mut buf: bytes::Bytes, version: Version) -> u64 { + let id = u64::decode(&mut buf, version).unwrap(); + let code = match version { + Version::Draft14 => { + assert_eq!(id, ietf::FetchError::ID); + ietf::FetchError::decode(&mut buf, version).unwrap().error_code + } + _ => { + assert_eq!(id, ietf::RequestError::ID); + ietf::RequestError::decode(&mut buf, version).unwrap().error_code + } + }; + assert!(buf.is_empty(), "{version}: a refusal opens no fetch stream"); + code + } + + /// A group that does not exist, below the newest one or past it, is refused. + #[tokio::test] + async fn a_standalone_fetch_of_a_missing_group_is_refused() { + for version in FETCH_DRAFTS { + for missing in [2, 5] { + let mut h = serve(version); + publish_pairs(&mut h, 4, Some(2)); + settle().await; + + let buf = standalone_fetch( + &h, + Location { + group: missing, + object: 0, + }, + Location { + group: missing, + object: 0, + }, + GroupOrder::Ascending, + ) + .await; + assert_eq!( + fetch_refusal(buf, version), + does_not_exist(version), + "{version}: group {missing}" + ); + } + } + } + + /// A joining FETCH reaching back before the subscription's group is refused, like any + /// FETCH touching several groups. + #[tokio::test] + async fn a_joining_fetch_reaching_back_is_refused() { + for version in JOINING_DRAFTS { + let mut h = serve(version); + publish_pairs(&mut h, 4, None); + let mut group = h.track.create_group(group::Info { sequence: 4 }).unwrap(); + group.write_frame(timestamp(), b"4-0".as_slice()).unwrap(); + settle().await; + + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + let mut serving = std::pin::pin!( + h.publisher + .clone() + .run_subscribe_stream(stream, subscribe(Filter::NextObject, None)) + ); + registered(&h, serving.as_mut()).await; + + for fetch_type in [ + FetchType::RelativeJoining { + subscriber_request_id: RequestId(REQUEST_ID), + group_offset: 2, + }, + FetchType::AbsoluteJoining { + subscriber_request_id: RequestId(REQUEST_ID), + group_id: 1, + }, + ] { + let mark = h.log.writes.lock().unwrap().len(); + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + h.publisher + .clone() + .run_fetch_stream( + stream, + ietf::Fetch { + request_id: FETCH_ID, + subscriber_priority: 128, + group_order: GroupOrder::Ascending, + fetch_type, + }, + ) + .await + .unwrap(); + let buf = bytes::Bytes::from(h.log.writes.lock().unwrap()[mark..].to_vec()); + + assert_eq!(fetch_refusal(buf, version), 0x3, "{version}"); + } + } + } + + /// An absolute joining FETCH starting past the subscription's group names an empty range. + #[tokio::test] + async fn an_absolute_joining_fetch_past_the_subscription_is_refused() { + let version = Version::Draft16; + let mut h = serve(version); + publish_pairs(&mut h, 3, None); + settle().await; + + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + let mut serving = std::pin::pin!( + h.publisher + .clone() + .run_subscribe_stream(stream, subscribe(Filter::NextObject, None)) + ); + registered(&h, serving.as_mut()).await; + + let mark = h.log.writes.lock().unwrap().len(); + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + h.publisher + .clone() + .run_fetch_stream( + stream, + ietf::Fetch { + request_id: FETCH_ID, + subscriber_priority: 128, + group_order: GroupOrder::Ascending, + fetch_type: FetchType::AbsoluteJoining { + subscriber_request_id: RequestId(REQUEST_ID), + group_id: 5, + }, + }, + ) + .await + .unwrap(); + let buf = bytes::Bytes::from(h.log.writes.lock().unwrap()[mark..].to_vec()); + assert_eq!(fetch_refusal(buf, version), invalid_range(version)); + } + /// A fill against an empty track has an empty range: no fetch stream is owed. #[tokio::test] async fn an_empty_track_opens_no_fill_stream() { @@ -5033,8 +5580,8 @@ mod tests { (writes, h.log.resets()) } - /// Send a FETCH we don't implement, returning the same pair. - async fn fetch_unsupported(version: Version, fetch_type: FetchType<'_>) -> (Vec, Vec) { + /// Send a FETCH we refuse, returning the same pair. + async fn fetch_refused(version: Version, fetch_type: FetchType<'_>) -> (Vec, Vec) { let h = harness(version); let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); @@ -5078,8 +5625,8 @@ mod tests { /// Every FETCH we refuse goes out through its own error encoder, so it needs the same /// finish: a reset there loses the rejection the same way. #[tokio::test] - async fn unsupported_fetch_is_refused_without_resetting_the_stream() { - let unsupported = || { + async fn a_refused_fetch_does_not_reset_the_stream() { + let refused = || { [ ( "standalone", @@ -5091,25 +5638,20 @@ mod tests { }, ), ( - "relative joining with an offset", - FetchType::RelativeJoining { - subscriber_request_id: RequestId(3), - group_offset: 1, - }, - ), - ( - "absolute joining", - FetchType::AbsoluteJoining { - subscriber_request_id: RequestId(3), - group_id: 7, + "empty range", + FetchType::Standalone { + namespace: crate::Path::new("nothing/here"), + track: "video".into(), + start: Location { group: 2, object: 0 }, + end: Location { group: 1, object: 0 }, }, ), ] }; for version in [Version::Draft17, Version::Draft18, Version::Draft19, Version::Draft20] { - for (label, fetch_type) in unsupported() { - let (writes, resets) = fetch_unsupported(version, fetch_type).await; + for (label, fetch_type) in refused() { + let (writes, resets) = fetch_refused(version, fetch_type).await; assert!(!writes.is_empty(), "{version} {label}: nothing was sent"); assert_eq!( diff --git a/rs/moq-net/src/ietf/session.rs b/rs/moq-net/src/ietf/session.rs index fcb9f30007..2230852b53 100644 --- a/rs/moq-net/src/ietf/session.rs +++ b/rs/moq-net/src/ietf/session.rs @@ -204,6 +204,7 @@ where subscriber.clone(), version ))); + let mut datagrams = std::pin::pin!(err_only(run_datagrams(adapter.clone(), subscriber.clone()))); // Unsolicited PUBLISH_NAMESPACE unless the peer requires solicitation; // see `Publisher::run_publish_namespaces`. let mut pub_ns_run = std::pin::pin!(err_only(publisher.clone().run_publish_namespaces())); @@ -254,6 +255,9 @@ where if let Poll::Ready(err) = waiter.poll_future(dispatch.as_mut()) { return Poll::Ready(Err(err)); } + if let Poll::Ready(err) = waiter.poll_future(datagrams.as_mut()) { + return Poll::Ready(Err(err)); + } if task_set.poll(waiter).is_ready() { return Poll::Ready(Ok(())); } @@ -346,6 +350,7 @@ where subscriber.clone(), version ))); + let mut datagrams = std::pin::pin!(err_only(run_datagrams(session.clone(), subscriber.clone()))); let mut goaway_recv = std::pin::pin!(err_only(goaway_recv)); let mut setup = std::pin::pin!(setup); // Unsolicited PUBLISH_NAMESPACE unless the peer requires solicitation; @@ -385,6 +390,9 @@ where if let Poll::Ready(err) = waiter.poll_future(dispatch.as_mut()) { return Poll::Ready(Err(err)); } + if let Poll::Ready(err) = waiter.poll_future(datagrams.as_mut()) { + return Poll::Ready(Err(err)); + } if let Poll::Ready(err) = waiter.poll_future(goaway_recv.as_mut()) { return Poll::Ready(Err(err)); } @@ -725,6 +733,23 @@ where } } +/// Receive QUIC datagrams, each an OBJECT_DATAGRAM for one of our subscriptions. +/// +/// A transport without datagrams never delivers one, so this parks. A transport failure +/// or a malformed datagram ends the session. +async fn run_datagrams(mut session: S, subscriber: Subscriber) -> Result<(), Error> +where + S: crate::transport::poll::Boxable, +{ + if session.max_datagram_size() == 0 { + return Ok(()); + } + loop { + let payload = session.recv_datagram().await.map_err(Error::from_transport)?; + subscriber.recv_datagram(payload)?; + } +} + async fn run_uni_group( subscriber: &mut Subscriber, stream: &mut Reader, diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index 4fc5a5c68b..79f6cd0f6b 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -145,6 +145,9 @@ struct State { // Joining FETCH request ids, mapped to the SUBSCRIBE they name. fetches: HashMap, + // Group FETCH request ids, filling a cache miss. + group_fetches: HashMap>, + // Track aliases chosen by the remote publisher. aliases: TrackAliases, @@ -285,6 +288,52 @@ impl Fill { } } +/// A standalone FETCH of one group through its end, from FETCH_OK to its fetch stream. +enum GroupFetch { + /// Waiting on FETCH_OK, which the fetch stream can overtake. + Pending, + /// Accepted into the track cache, waiting for the fetch stream to write it. + Ready { + producer: group::Producer, + timescale: Option, + /// The first Object ID requested, which the producer starts at. + start: u64, + /// The exclusive last Object ID, when FETCH_OK's End Location falls inside the group. + end: Option, + }, + /// A fetch stream is writing the group. + Receiving, + /// The group is written or failed. + Done, +} + +/// A group FETCH's entry in [`State::group_fetches`], removed when its request ends. +struct GroupFetchEntry { + state: Lock, + fetch_id: RequestId, + slot: kio::Producer, +} + +impl GroupFetchEntry { + fn new(state: &Lock, fetch_id: RequestId, slot: kio::Producer) -> Self { + state.lock().group_fetches.insert(fetch_id, slot.clone()); + Self { + state: state.clone(), + fetch_id, + slot, + } + } +} + +impl Drop for GroupFetchEntry { + fn drop(&mut self) { + self.state.lock().group_fetches.remove(&self.fetch_id); + // A fetch stream that overtook a refused or unserved FETCH_OK holds its own clone + // of the slot, so only closing it wakes that stream. + let _ = self.slot.close(); + } +} + /// A pre-draft-20 joining FETCH, sent as its own request after SUBSCRIBE. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum JoiningFetch { @@ -1730,6 +1779,9 @@ where true => request.resolving_start(), false => request, }; + // Serves cache misses with a group FETCH. Registered before accepting, so a miss + // queued meanwhile waits for it rather than failing for want of a handler. + let dynamic = request.dynamic(); let mut track = request.accept(info); if live { let _ = track.start_at(largest.map(|largest| largest.group)); @@ -1776,9 +1828,12 @@ where enum End { Unused, Done(Result), + Fetch(group::Request), } let mut fetch_done = fetching.is_none(); + // Group FETCHes for cache misses, cancelled with the subscription. + let mut group_fetches = TaskSet::owned(); let cancelled = { let mut done = std::pin::pin!(Self::read_publish_done(&mut stream.reader, self.version)); loop { @@ -1789,6 +1844,11 @@ where { fetch_done = true; } + // An error is the track closing, which the arms below report. + if let Poll::Ready(Ok(request)) = dynamic.poll_requested_group(waiter) { + return Poll::Ready(End::Fetch(request)); + } + let _ = group_fetches.poll(waiter); if track.poll_unused(waiter).is_ready() { return Poll::Ready(End::Unused); } @@ -1797,6 +1857,15 @@ where .await; match end { + End::Fetch(request) => { + let fetch = self.clone().run_group_fetch( + broadcast_path.to_owned(), + track_name.clone(), + request, + timescale, + ); + group_fetches.push(fetch); + } End::Unused => match track.abort_unused(Error::Cancel) { Ok(()) => { tracing::info!(broadcast = %self.origin.absolute(&broadcast_path), track = %track_name, "subscribe cancelled"); @@ -2265,6 +2334,61 @@ where Ok(()) } + + /// Deliver one OBJECT_DATAGRAM as a datagram on its subscription's track: a + /// single-frame group at the Group ID. + /// + /// A malformed datagram is the peer breaking the protocol, so it errors. One the model + /// cannot carry is dropped like any lost datagram: an Object past ID 0 (the group would + /// need a second object), a status other than Normal, or an alias that is not bound + /// yet (the draft lets us drop rather than buffer). + pub fn recv_datagram(&self, payload: bytes::Bytes) -> Result<(), Error> { + let mut buf = payload; + let datagram = ietf::ObjectDatagram::decode(&mut buf, self.version)?; + let (alias, sequence) = (datagram.track_alias, datagram.group_id); + + if datagram.object_id.unwrap_or(0) != 0 { + tracing::debug!(alias, sequence, "dropping a datagram past object 0"); + return Ok(()); + } + let payload = match datagram.body { + ietf::DatagramBody::Payload(payload) => payload, + ietf::DatagramBody::Status(0) => bytes::Bytes::new(), + ietf::DatagramBody::Status(status) => { + tracing::debug!(alias, sequence, status, "dropping a datagram status"); + return Ok(()); + } + }; + + let mut state = self.state.lock(); + let request_id = match state.aliases.read().map.get(&alias) { + Some(Alias::Active(request_id)) => *request_id, + _ => { + tracing::debug!(alias, sequence, "dropping a datagram for an unbound alias"); + return Ok(()); + } + }; + let Some(track) = state.subscribes.get_mut(&request_id) else { + return Ok(()); + }; + + // Like a subgroup object: a track that declared no timescale is stamped on arrival. + let timestamp = match (track.timescale, &datagram.properties) { + (Some(timescale), Some(properties)) => { + ietf::decode_object_time(&mut properties.as_slice(), timescale, self.version)? + } + _ => None, + }; + let timestamp = timestamp.unwrap_or_else(|| crate::Timestamp::from(self.runtime.now())); + + let Some(producer) = track.producer.as_mut() else { + return Ok(()); + }; + if let Err(err) = producer.insert_datagram(sequence, timestamp, payload) { + tracing::debug!(%err, alias, sequence, "dropping datagram"); + } + Ok(()) + } } /// Mark where the track ends, as an END_OF_TRACK object said. @@ -2303,6 +2427,9 @@ enum Ended { Track, } +/// Object status: no object at or past this location in the group exists (draft-14, 15). +const END_OF_GROUP: u64 = 0x3; + /// Object status: no object at or past this location exists (every implemented draft). const END_OF_TRACK: u64 = 0x4; @@ -2561,6 +2688,11 @@ where let _: u64 = stream.decode().await?; let header: ietf::FetchHeader = stream.decode().await?; + let group_fetch = self.state.lock().group_fetches.get(&header.request_id).cloned(); + if let Some(slot) = group_fetch { + return self.recv_group_fetch(stream, slot).await; + } + let (subscribe_id, fill, joining, largest, _counted) = { let state = self.state.lock(); // A draft-20 fill is named by the SUBSCRIBE's request id. A pre-draft-20 joining @@ -2719,7 +2851,8 @@ where head: &mut Option<(u64, u64, crate::recv::Group)>, ) -> Result<(), Error> { let mut prior_group = None; - while let Some(object) = decode_fetch_object(stream, self.version).await? { + let mut first = true; + while let Some(object) = decode_fetch_object(stream, self.version, std::mem::take(&mut first)).await? { if !object.subgroup_ok { tracing::warn!("subgroup ID is not supported, dropping fill"); return Err(Error::Unsupported); @@ -2778,40 +2911,363 @@ where } } - // The properties carry the frame's presentation timestamp (the Timestamp Object - // Property) in the units the track declared. A track that declared none opted - // out, so its frames are stamped on arrival instead. - let timestamp = match (object.properties, timescale) { - (Some(properties), Some(timescale)) => { - let mut properties = bytes::Bytes::from(properties); - ietf::decode_object_time(&mut properties, timescale, self.version)? - } - _ => None, + let (_, next, producer) = head.as_mut().expect("the head was created above"); + if !self + .recv_fetch_payload(stream, producer, object.properties, timescale) + .await? + { + return Err(Error::Unsupported); + } + *next += 1; + } + + Ok(()) + } + + /// Read one fetch object's length and payload into `producer`, after its header. + /// + /// Returns `false` for a draft-14 or 15 end-of-group or end-of-track marker, which + /// is a status rather than a frame. + async fn recv_fetch_payload( + &self, + stream: &mut Reader, + producer: &mut group::Producer, + properties: Option>, + timescale: Option, + ) -> Result { + // The properties carry the frame's presentation timestamp (the Timestamp Object + // Property) in the units the track declared. A track that declared none opted + // out, so its frames are stamped on arrival instead. + let timestamp = match (properties, timescale) { + (Some(properties), Some(timescale)) => { + let mut properties = bytes::Bytes::from(properties); + ietf::decode_object_time(&mut properties, timescale, self.version)? + } + _ => None, + }; + let timestamp = timestamp.unwrap_or_else(|| crate::Timestamp::from(self.runtime.now())); + + // A fetch object has no status field from draft-16 on; a zero length is simply + // an empty object. Draft-14 and 15 still encode Normal (0) after a zero length. + let size: u64 = stream.decode().await?; + if size == 0 && matches!(self.version, Version::Draft14 | Version::Draft15) { + match stream.decode::().await? { + 0 => {} + END_OF_GROUP | END_OF_TRACK => return Ok(false), + _ => return Err(Error::Unsupported), + } + } + + // `create_frame_owned` is the allocation chokepoint and rejects an oversized `size` + // before allocating, so no pre-check is needed. + let mut frame = producer.create_frame_owned(frame::Info { size, timestamp })?; + if let Err(err) = std::future::poll_fn(|cx| stream.poll_read_frame(cx, &mut frame)).await { + let _ = frame.abort(err.clone()); + return Err(err); + } + frame.finish()?; + Ok(true) + } + + /// Fetch a group from the publisher to fill a cache miss, with a standalone FETCH for + /// this subscription's track. It asks from the frame the reader wants through the end + /// of the group, so a publisher that evicted the prefix can still answer. + /// + /// The group is accepted only once FETCH_OK arrives, so a refusal reaches every + /// waiting [`track::Consumer::fetch_group`] as the publisher's own error. The objects + /// arrive on a fetch stream, which [`Self::recv_fill`] routes into the group. + async fn run_group_fetch( + self, + broadcast: PathOwned, + name: String, + request: group::Request, + timescale: Option, + ) { + let sequence = request.sequence(); + let start = request.frame_start(); + if self.going_away.is_set() { + request.reject(Error::GoingAway); + return; + } + // Our FETCH still encodes the Fetch Type field that draft-20 removed. + if Filter::is_draft20(self.version) { + request.reject(Error::Unsupported); + return; + } + + let fetch_id = match self.control.next_request_id(&self.runtime).await { + Ok(id) => id, + Err(err) => return request.reject(err), + }; + + // Registered before the FETCH goes out, since its fetch stream can overtake FETCH_OK. + let slot = kio::Producer::new(GroupFetch::Pending); + let _registered = GroupFetchEntry::new(&self.state, fetch_id, slot.clone()); + + let mut stream = match Stream::open(&mut self.session.clone(), self.version).await { + Ok(stream) => stream, + Err(err) => return request.reject(err), + }; + + let res = async { + stream.writer.encode(&ietf::Fetch::ID).await?; + stream + .writer + .encode(&ietf::Fetch { + request_id: fetch_id, + subscriber_priority: super::priority::to_wire(request.priority()), + group_order: GroupOrder::Ascending, + // An End Object of 0 is the whole End Group. + fetch_type: FetchType::Standalone { + namespace: broadcast.clone(), + track: name.as_str().into(), + start: ietf::Location { + group: sequence, + object: start, + }, + end: ietf::Location { + group: sequence, + object: 0, + }, + }, + }) + .await?; + self.read_group_fetch_response(&mut stream).await + } + .await; + + let ok = match res { + Ok(ok) => ok, + Err(err) => { + tracing::debug!(%err, group = sequence, "group fetch refused"); + request.reject(err); + let _ = stream.writer.close().await; + return; + } + }; + + // The publisher knows where the track ends, which a range FETCH downstream needs. + if ok.end_of_track { + let end = ok.end_location; + let Some(final_sequence) = end.group.checked_add(u64::from(end.object > 0)) else { + request.reject(Error::ProtocolViolation); + let _ = stream.writer.close().await; + return; }; - let timestamp = timestamp.unwrap_or_else(|| crate::Timestamp::from(self.runtime.now())); - - // A fetch object has no status field from draft-16 on; a zero length is simply - // an empty object. Draft-14 and 15 still encode Normal (0) after a zero length. - let size: u64 = stream.decode().await?; - if size == 0 && matches!(self.version, Version::Draft14 | Version::Draft15) { - let status: u64 = stream.decode().await?; - if status != 0 { - return Err(Error::Unsupported); + request.finish_track_at(final_sequence); + } + + // An empty answer opens no fetch stream at all. + let end = ok.end_location; + if (end.group, end.object) <= (sequence, start) { + request.reject(Error::NotFound); + let _ = stream.writer.close().await; + return; + } + // An exclusive End Location inside the group is its Largest Object plus one, so + // the stream owes every object up to it. One past the group covers it whole. + let end = (end.group == sequence).then_some(end.object); + + // Only a track nothing has subscribed to yet takes this info, as SUBSCRIBE_OK would + // have set it. + let info = track::Info::default() + .with_timescale(Timescale::MICRO) + .with_max_age(self.origin.default_max_age()); + let mut producer = match request.accept(info) { + Ok(producer) => producer, + // Already served by a concurrent fetch, or the track closed. + Err(err) => { + tracing::debug!(%err, group = sequence, "group fetch not served"); + let _ = stream.writer.close().await; + return; + } + }; + // The objects keep the IDs they have in the group rather than restarting at 0. + if let Err(err) = producer.start_at(start) { + let _ = producer.abort(err); + let _ = stream.writer.close().await; + return; + } + if let Ok(mut state) = slot.write() { + *state = GroupFetch::Ready { + producer, + timescale, + start, + end, + }; + } + + // Hold the request open until its fetch stream is done: closing our side first is + // what a draft-14-16 adapter reads as cancelling the FETCH. A publisher that fails + // after FETCH_OK resets the request instead and owes no fetch stream, so the group + // it left waiting is aborted. A FIN is not that: the fetch stream can trail it. + let mut open = true; + let reset = kio::wait(|waiter| { + if open { + let mut cx = std::task::Context::from_waker(waiter.waker()); + match stream.reader.poll_closed(&mut cx) { + Poll::Ready(Err(err)) => return Poll::Ready(Some(err)), + Poll::Ready(Ok(())) => open = false, + Poll::Pending => {} } } + slot.poll(waiter, |state| match &**state { + GroupFetch::Done => Poll::Ready(()), + _ => Poll::Pending, + }) + .map(|_| None) + }) + .await; + if let Some(err) = reset + && let Ok(mut state) = slot.write() + && matches!(*state, GroupFetch::Ready { .. }) + && let GroupFetch::Ready { producer, .. } = std::mem::replace(&mut *state, GroupFetch::Done) + { + let _ = producer.abort(err); + } + let _ = stream.writer.close().await; + } - let (_, next, producer) = head.as_mut().expect("the head was created above"); + /// Read the answer to a group FETCH: FETCH_OK, or the publisher's refusal as an error. + async fn read_group_fetch_response(&self, stream: &mut Stream) -> Result { + let type_id: u64 = stream.reader.decode().await?; + let size: u16 = stream.reader.decode().await?; + let mut data = stream.reader.read_exact(size as usize).await?; - // `create_frame_owned` is the allocation chokepoint and rejects an oversized `size` - // before allocating, so no pre-check is needed. - let mut frame = producer.create_frame_owned(frame::Info { size, timestamp })?; - if let Err(err) = std::future::poll_fn(|cx| stream.poll_read_frame(cx, &mut frame)).await { - let _ = frame.abort(err.clone()); - return Err(err); + match type_id { + ietf::FetchOk::ID => Ok(ietf::FetchOk::decode_msg(&mut data, self.version)?), + ietf::FetchError::ID if self.version == Version::Draft14 => { + let msg = ietf::FetchError::decode_msg(&mut data, self.version)?; + Err(request::from_code(msg.error_code, request::Kind::Fetch, self.version)) } - frame.finish()?; + ietf::RequestError::ID => { + let msg = ietf::RequestError::decode_msg(&mut data, self.version)?; + Err(request::from_code(msg.error_code, request::Kind::Fetch, self.version)) + } + _ => Err(Error::UnexpectedMessage), + } + } - *next += 1; + /// Write a group FETCH's objects into the group it accepted. + async fn recv_group_fetch( + &mut self, + stream: &mut Reader, + slot: kio::Producer, + ) -> Result<(), Error> { + // FETCH_OK can trail its own fetch stream. Taking the group in the same step is + // what refuses a second stream for one request. + let taken = kio::wait(|waiter| { + match slot.poll(waiter, |state| match &**state { + GroupFetch::Pending => Poll::Pending, + _ => Poll::Ready(()), + }) { + Poll::Ready(Ok(mut state)) => { + Poll::Ready(match std::mem::replace(&mut *state, GroupFetch::Receiving) { + GroupFetch::Ready { + producer, + timescale, + start, + end, + } => Ok((producer, timescale, start, end)), + other => { + *state = other; + Err(Error::Unsupported) + } + }) + } + Poll::Ready(Err(_)) => Poll::Ready(Err(Error::Dropped)), + Poll::Pending => Poll::Pending, + } + }) + .await; + let (producer, timescale, start, end) = taken?; + + let mut producer = crate::recv::Group::new(producer); + let res = match self + .recv_group_fetch_objects(stream, &mut producer, start, end, timescale) + .await + { + Ok(()) => producer.finish(), + Err(err) => { + let _ = producer.abort(err.clone()); + Err(err) + } + }; + + if let Ok(mut state) = slot.write() { + *state = GroupFetch::Done; + } + res + } + + /// Decode one group's objects: all in the producer's group, numbered from `start` with + /// no gaps, and through `end` (exclusive) when FETCH_OK named one inside the group. + async fn recv_group_fetch_objects( + &self, + stream: &mut Reader, + producer: &mut group::Producer, + start: u64, + end: Option, + timescale: Option, + ) -> Result<(), Error> { + let sequence = producer.info().sequence; + let mut next = start; + let mut prior_group = None; + let mut ended = false; + let mut first = true; + while let Some(object) = decode_fetch_object(stream, self.version, std::mem::take(&mut first)).await? { + if ended { + tracing::warn!(sequence, "a group fetch continued past its end marker"); + return Err(Error::ProtocolViolation); + } + if !object.subgroup_ok { + tracing::warn!("subgroup ID is not supported, dropping group fetch"); + return Err(Error::Unsupported); + } + + let group = resolve_fetch_group(self.version, prior_group, object.group)?; + if let Some(group) = group { + prior_group = Some(group); + } + let id = match (object.group.is_some(), object.object) { + (true, Some(id)) => Some(id), + (false, None | Some(1)) => Some(next), + _ => None, + }; + // Another group, or an ID before the next one, is outside what was asked for. + if group.is_some_and(|group| group != sequence) || id.is_some_and(|id| id < next) { + tracing::warn!(sequence, next, group = ?group, object = ?id, "a group fetch answered outside its range"); + return Err(Error::ProtocolViolation); + } + // A skipped ID is an object that does not exist, a hole the model cannot hold. + if id != Some(next) { + tracing::warn!(sequence, next, object = ?id, "a group fetch skipped an object"); + return Err(Error::Unsupported); + } + + match self + .recv_fetch_payload(stream, producer, object.properties, timescale) + .await? + { + true => next += 1, + false => ended = true, + } + if end.is_some_and(|end| next > end) { + tracing::warn!(sequence, ?end, "a group fetch continued past FETCH_OK's End Location"); + return Err(Error::ProtocolViolation); + } + } + + // A clean FIN short of the promised end would otherwise cache a truncated group as + // a complete one. + if end.is_some_and(|end| next < end) { + tracing::warn!( + sequence, + next, + ?end, + "a group fetch ended before FETCH_OK's End Location" + ); + return Err(Error::ProtocolViolation); } Ok(()) @@ -2827,9 +3283,13 @@ struct FetchedObject { } /// Decode the next fetch object header, or `None` at stream end. +/// +/// The `first` object on a stream has no prior to inherit from, so a field it leaves to +/// the prior object is a protocol violation. async fn decode_fetch_object( stream: &mut Reader, version: Version, + first: bool, ) -> Result, Error> { if version == Version::Draft14 { let Some(group) = stream.decode_maybe::().await? else { @@ -2857,17 +3317,27 @@ async fn decode_fetch_object( subgroup, group, object, + priority, properties, - .. - }) => Some(FetchedObject { - group, - object, - subgroup_ok: matches!( - subgroup, - ietf::FetchSubgroup::Zero | ietf::FetchSubgroup::Prior | ietf::FetchSubgroup::Explicit(0) - ), - properties, - }), + }) => { + let inherits = group.is_none() + || object.is_none() + || priority.is_none() + || matches!(subgroup, ietf::FetchSubgroup::Prior | ietf::FetchSubgroup::PriorPlusOne); + if first && inherits { + tracing::warn!("the first fetch object refers to a prior object"); + return Err(Error::ProtocolViolation); + } + Some(FetchedObject { + group, + object, + subgroup_ok: matches!( + subgroup, + ietf::FetchSubgroup::Zero | ietf::FetchSubgroup::Prior | ietf::FetchSubgroup::Explicit(0) + ), + properties, + }) + } }) } @@ -3009,7 +3479,7 @@ impl GroupIngest { let frame = group.create_frame_owned(frame::Info { size: 0, timestamp })?; frame.finish()?; self.phase = IngestPhase::Delta; - } else if status == 3 && !self.has_end { + } else if status == END_OF_GROUP && !self.has_end { self.phase = IngestPhase::Finished(Ended::Group); } else if status == END_OF_TRACK { // Defined on every implemented draft, whether or not the header marks @@ -3508,6 +3978,67 @@ mod tests { ); } + /// An OBJECT_DATAGRAM at object 0 is a datagram group at its Group ID; anything the model + /// cannot carry as one is dropped, and a malformed one is the peer's violation. + #[tokio::test] + async fn an_object_datagram_is_a_datagram_group() { + use crate::coding::Encode as _; + use futures::FutureExt as _; + + let subscriber = subscriber_with_tracks(&[(RequestId(11), "cam", "audio")]); + subscriber.register_alias(RequestId(11), 7).unwrap(); + let mut consumer = { + let mut state = subscriber.state.lock(); + let track = state.subscribes.get_mut(&RequestId(11)).unwrap(); + track.timescale = Some(Timescale::default()); + track.producer.as_ref().unwrap().subscribe(None) + }; + + let timestamp = crate::Timestamp::new(96_000, Timescale::default()).unwrap(); + let datagram = |alias: u64, group_id: u64, object_id: Option, body: ietf::DatagramBody| { + let mut properties = Vec::new(); + ietf::encode_object_time(&mut properties, timestamp, Timescale::default(), Version::Draft19).unwrap(); + ietf::ObjectDatagram { + track_alias: alias, + group_id, + object_id, + publisher_priority: None, + // Only a Normal Object may carry Properties, and a status cannot end the group. + end_of_group: matches!(body, ietf::DatagramBody::Payload(_)), + properties: matches!(body, ietf::DatagramBody::Payload(_)).then_some(properties), + body, + } + .encode_bytes(Version::Draft19) + .unwrap() + }; + let payload = |bytes: &'static [u8]| ietf::DatagramBody::Payload(bytes::Bytes::from_static(bytes)); + + // Dropped: a second object in the group, an unbound alias, and a status. + subscriber + .recv_datagram(datagram(7, 4, Some(1), payload(b"no"))) + .unwrap(); + subscriber.recv_datagram(datagram(8, 4, None, payload(b"no"))).unwrap(); + subscriber + .recv_datagram(datagram(7, 4, None, ietf::DatagramBody::Status(END_OF_TRACK))) + .unwrap(); + + subscriber + .recv_datagram(datagram(7, 9, Some(0), payload(b"yes"))) + .unwrap(); + let received = consumer.recv_datagram().now_or_never().unwrap().unwrap().unwrap(); + assert_eq!(received.sequence, 9, "the Group ID is the sequence"); + assert_eq!(received.timestamp, timestamp); + assert_eq!(&received.payload[..], b"yes"); + assert!( + consumer.recv_datagram().now_or_never().is_none(), + "only one got through" + ); + + // A status datagram cannot end the group. + let malformed = bytes::Bytes::from_static(&[0x22, 0x07, 0x04, 0x00]); + assert!(is_protocol_violation(&subscriber.recv_datagram(malformed).unwrap_err())); + } + /// One alias naming two different tracks is the collision section 11.1 makes fatal. #[test] fn an_alias_reused_for_another_track_is_fatal() { @@ -6591,6 +7122,177 @@ mod stitch_tests { assert_eq!(frames[0].1, b"g7-0"); assert_eq!(frames[1].1, b"g7-1"); } + + /// A draft-14/15 end-of-group marker ends a group fetch, so an object after it is a + /// violation rather than another frame. + #[tokio::test] + async fn a_group_fetch_refuses_an_object_past_its_end_marker() { + const DRAFT: Version = Version::Draft15; + let header = |object: u64| ietf::FetchObject::Object { + subgroup: ietf::FetchSubgroup::Zero, + group: Some(SEQUENCE), + object: Some(object), + priority: Some(0), + properties: None, + }; + let mut buf = bytes::BytesMut::new(); + header(0).encode(&mut buf, DRAFT).unwrap(); + 1u64.encode(&mut buf, DRAFT).unwrap(); + buf.put_slice(b"a"); + header(1).encode(&mut buf, DRAFT).unwrap(); + 0u64.encode(&mut buf, DRAFT).unwrap(); + END_OF_GROUP.encode(&mut buf, DRAFT).unwrap(); + header(1).encode(&mut buf, DRAFT).unwrap(); + 1u64.encode(&mut buf, DRAFT).unwrap(); + buf.put_slice(b"b"); + + let mut run = GroupFetchRun::new(DRAFT, buf.to_vec()).await; + let res = run.recv_objects(0, None).await; + assert!(matches!(res, Err(Error::ProtocolViolation)), "{res:?}"); + } + + /// The first object on a fetch stream has no prior object, so leaving any field to + /// the prior one is a violation, not "the requested group, from its start". + #[tokio::test] + async fn a_group_fetch_refuses_a_first_object_that_inherits() { + use ietf::FetchSubgroup::{Prior, Zero}; + let header = |group, object, subgroup, priority| ietf::FetchObject::Object { + subgroup, + group, + object, + priority, + properties: None, + }; + let headers = [ + // Both IDs omitted: "the requested group, from its start" is exactly the bug. + header(None, None, Zero, Some(0)), + header(Some(SEQUENCE), None, Zero, Some(0)), + header(Some(SEQUENCE), Some(0), Prior, Some(0)), + header(Some(SEQUENCE), Some(0), Zero, None), + ]; + + for version in [Version::Draft15, VERSION] { + for header in &headers { + let mut buf = bytes::BytesMut::new(); + header.encode(&mut buf, version).unwrap(); + 1u64.encode(&mut buf, version).unwrap(); + buf.put_slice(b"a"); + + let mut run = GroupFetchRun::new(version, buf.to_vec()).await; + let res = run.recv_objects(0, None).await; + assert!( + matches!(res, Err(Error::ProtocolViolation)), + "{version} {header:?}: {res:?}" + ); + } + } + } + + /// FETCH_OK's End Location inside the group promises every object before it. A stream + /// that FINs short of it, or runs past it, fails the group instead of caching it. + #[tokio::test] + async fn a_group_fetch_must_reach_its_end_location() { + const END: u64 = 3; + for (count, complete) in [(2, false), (3, true), (4, false)] { + let payloads: Vec<&[u8]> = [b"a", b"b", b"c", b"d"][..count].iter().map(|p| &p[..]).collect(); + let mut run = GroupFetchRun::new(VERSION, group_fetch_objects(SEQUENCE, 0, &payloads)).await; + + let group = run.track.create_group(group::Info { sequence: SEQUENCE }).unwrap(); + let mut consumer = group.consume(); + let slot = kio::Producer::new(GroupFetch::Ready { + producer: group, + timescale: None, + start: 0, + end: Some(END), + }); + let res = run.subscriber.recv_group_fetch(&mut run.stream, slot).await; + assert_eq!(res.is_ok(), complete, "{count} objects: {res:?}"); + + let mut read = 0; + let end = loop { + match consumer.read_frame().await { + Ok(Some(_)) => read += 1, + Ok(None) => break Ok(read), + Err(err) => break Err(err), + } + }; + match complete { + true => assert_eq!(end.expect("a complete group"), END), + false => assert!(end.is_err(), "{count} objects: the group must fail, not end"), + } + } + } + + /// A group fetch's objects after the FETCH_HEADER: the first one names the group and + /// `start`, and every later one is the next object. + fn group_fetch_objects(sequence: u64, start: u64, payloads: &[&[u8]]) -> Vec { + let mut buf = bytes::BytesMut::new(); + for (index, payload) in payloads.iter().enumerate() { + let first = index == 0; + ietf::FetchObject::Object { + subgroup: ietf::FetchSubgroup::Zero, + group: first.then_some(sequence), + object: first.then_some(start), + priority: first.then_some(0), + properties: None, + } + .encode(&mut buf, VERSION) + .unwrap(); + (payload.len() as u64).encode(&mut buf, VERSION).unwrap(); + buf.put_slice(payload); + } + buf.to_vec() + } + + /// A subscriber reading one scripted group fetch stream, already past its header. + struct GroupFetchRun { + subscriber: Subscriber, + stream: Reader<::RecvStream, Version>, + track: track::Producer, + _tasks: (Tasks, TaskSet), + } + + impl GroupFetchRun { + async fn new(version: Version, objects: Vec) -> Self { + let mut session = ScriptedSession::per_stream_eof(vec![objects]); + let tasks = TaskSet::new(); + let subscriber = Subscriber::new( + crate::time::Clock::tokio(), + session.clone(), + crate::origin::Config::new(crate::Hop::new(1).unwrap()).produce(), + Control::new(None, false), + None, + peer::PeerSetup::default(), + crate::Hop::new(1).unwrap(), + None, + version, + tasks.0.clone(), + Default::default(), + ); + let (_, recv) = session.open_bi().await.unwrap(); + let track = track::Producer::new( + std::sync::Arc::new(crate::broadcast::Info::default()), + "video", + track::Info::default(), + ); + + Self { + subscriber, + stream: Reader::new(recv, version), + track, + _tasks: tasks, + } + } + + /// Decode the stream into a fresh group numbered from `start`. + async fn recv_objects(&mut self, start: u64, end: Option) -> Result<(), Error> { + let mut group = self.track.create_group(group::Info { sequence: SEQUENCE }).unwrap(); + group.start_at(start).unwrap(); + self.subscriber + .recv_group_fetch_objects(&mut self.stream, &mut group, start, end, None) + .await + } + } } /// A scripted peer answering SUBSCRIBE then FETCH, so the join is spelled on the wire. @@ -6601,6 +7303,7 @@ mod joining_fetch_tests { coding::Encode as _, lite::test_transport::ScriptedSession, model::ProduceTest, + transport::poll::Session as _, util::{TaskSet, Tasks}, }; @@ -6901,4 +7604,179 @@ mod joining_fetch_tests { let writes = log.writes.lock().unwrap().clone(); assert!(writes.is_empty(), "a refused join must not write SUBSCRIBE"); } + + /// A publisher that resets the request after FETCH_OK owes no fetch stream, so the + /// group it accepted is aborted instead of left open for every reader to wait on. + /// Without that, `run_group_fetch` never returns. + #[tokio::test(start_paused = true)] + async fn a_group_fetch_reset_after_fetch_ok_aborts_the_group() { + const VERSION: Version = Version::Draft19; + const GROUP: u64 = 4; + + let ok = message_bytes( + ietf::FetchOk::ID, + &ietf::FetchOk { + request_id: None, + group_order: GroupOrder::Ascending, + end_of_track: false, + end_location: ietf::Location { + group: GROUP + 1, + object: 0, + }, + }, + VERSION, + ); + + let session = ScriptedSession::per_stream_reset(vec![ok]); + let (tasks, _task_set) = crate::util::TaskSet::new(); + let subscriber = Subscriber::new( + crate::time::Clock::tokio(), + session.clone(), + crate::origin::Config::new(crate::Hop::new(1).unwrap()).produce(), + Control::new(None, false), + None, + peer::PeerSetup::default(), + crate::Hop::new(1).unwrap(), + None, + VERSION, + tasks, + Default::default(), + ); + + let track = track::Producer::new( + std::sync::Arc::new(crate::broadcast::Info::default()), + "video", + track::Info::default(), + ); + let dynamic = track.dynamic(); + let consumer = track.consume(); + let mut fetch = std::pin::pin!(consumer.fetch_group(GROUP, None)); + assert!(futures::poll!(fetch.as_mut()).is_pending()); + let request = dynamic.requested_group().await.expect("no group requested"); + + subscriber + .clone() + .run_group_fetch(Path::new("broadcast").to_owned(), "video".into(), request, None) + .await; + + assert!(fetch.await.is_err(), "the accepted group was aborted"); + } + + /// A cache miss for a group's tail asks upstream from the frame the reader wants and + /// numbers what arrives from there, so a publisher that evicted the prefix can answer. + #[tokio::test(start_paused = true)] + async fn a_group_fetch_asks_from_the_wanted_frame() { + const VERSION: Version = Version::Draft19; + const GROUP: u64 = 4; + const START: u64 = 2; + + let ok = message_bytes( + ietf::FetchOk::ID, + &ietf::FetchOk { + request_id: None, + group_order: GroupOrder::Ascending, + end_of_track: false, + end_location: ietf::Location { + group: GROUP + 1, + object: 0, + }, + }, + VERSION, + ); + + // The fetch stream answering our first request id, from object START on. + let mut objects = bytes::BytesMut::new(); + ietf::FetchHeader::TYPE.encode(&mut objects, VERSION).unwrap(); + ietf::FetchHeader { + request_id: RequestId(1), + } + .encode(&mut objects, VERSION) + .unwrap(); + for (index, payload) in [b"c", b"d"].iter().enumerate() { + let first = index == 0; + ietf::FetchObject::Object { + subgroup: ietf::FetchSubgroup::Zero, + group: first.then_some(GROUP), + object: first.then_some(START), + priority: first.then_some(0), + properties: None, + } + .encode(&mut objects, VERSION) + .unwrap(); + 1u64.encode(&mut objects, VERSION).unwrap(); + objects.extend_from_slice(&payload[..]); + } + + let session = ScriptedSession::per_stream_eof(vec![ok, objects.to_vec()]); + let (tasks, _task_set) = crate::util::TaskSet::new(); + let subscriber = Subscriber::new( + crate::time::Clock::tokio(), + session.clone(), + crate::origin::Config::new(crate::Hop::new(1).unwrap()).produce(), + Control::new(None, false), + None, + peer::PeerSetup::default(), + crate::Hop::new(1).unwrap(), + None, + VERSION, + tasks, + Default::default(), + ); + + let track = track::Producer::new( + std::sync::Arc::new(crate::broadcast::Info::default()), + "video", + track::Info::default(), + ); + let dynamic = track.dynamic(); + let consumer = track.consume(); + let mut fetch = std::pin::pin!(consumer.fetch_group(GROUP, group::Fetch::default().with_frame_start(START))); + assert!(futures::poll!(fetch.as_mut()).is_pending()); + let request = dynamic.requested_group().await.expect("no group requested"); + + let serving = tokio::spawn(subscriber.clone().run_group_fetch( + Path::new("broadcast").to_owned(), + "video".into(), + request, + None, + )); + settle().await; + + // The fetch stream, as the peer would open it. + let (_, recv) = session.clone().open_bi().await.unwrap(); + let mut stream = Reader::new(recv, VERSION); + subscriber.clone().recv_fill(&mut stream).await.expect("group fetch"); + serving.await.expect("run_group_fetch"); + + let mut group = fetch.await.expect("fetched"); + assert_eq!(group.index(), START, "the group starts where the reader wanted"); + let mut payloads = Vec::new(); + while let Some(frame) = group.read_frame().await.expect("complete") { + payloads.push(frame.payload.to_vec()); + } + assert_eq!(payloads, [b"c".to_vec(), b"d".to_vec()]); + + let messages = decode_messages(&session.log, VERSION); + let fetch = messages.iter().find(|(id, _)| *id == ietf::Fetch::ID).expect("FETCH"); + let mut body = fetch.1.clone(); + let msg = ietf::Fetch::decode_msg(&mut body, VERSION).unwrap(); + let FetchType::Standalone { start, end, .. } = msg.fetch_type else { + panic!("a group fetch is standalone: {:?}", msg.fetch_type); + }; + assert_eq!( + start, + ietf::Location { + group: GROUP, + object: START + } + ); + assert_eq!( + end, + ietf::Location { + group: GROUP, + object: 0 + }, + "through the end of the group" + ); + } } diff --git a/rs/moq-net/src/lite/test_transport.rs b/rs/moq-net/src/lite/test_transport.rs index ff69af76af..497cb3b326 100644 --- a/rs/moq-net/src/lite/test_transport.rs +++ b/rs/moq-net/src/lite/test_transport.rs @@ -546,6 +546,8 @@ pub struct ScriptedRecv { /// Report EOF once the script is exhausted rather than parking, so a test can drive /// a read loop all the way through its exit path. See [`ScriptedSession::eof`]. eof: bool, + /// Report a reset once the script is exhausted. See [`ScriptedSession::per_stream_reset`]. + reset: bool, log: Log, } @@ -566,6 +568,7 @@ impl poll::RecvStream for ScriptedRecv { }; match take { + 0 if self.reset => Poll::Ready(Err(SinkError)), 0 if self.eof => Poll::Ready(Ok(None)), 0 => Poll::Pending, take => Poll::Ready(Ok(Some(take))), @@ -592,6 +595,8 @@ pub struct ScriptedSession { pub log: Log, /// Whether an exhausted script reports EOF instead of parking. eof: bool, + /// Whether an exhausted script reports a reset instead of parking. + reset: bool, script: Arc>>, /// Per-stream scripts popped by `open_bi` in order; `None` shares `script` /// across every stream. @@ -611,6 +616,7 @@ impl ScriptedSession { Self { log: Log::default(), eof: false, + reset: false, script: Arc::new(Mutex::new(script)), queue: None, open_gate: None, @@ -661,6 +667,15 @@ impl ScriptedSession { } } + /// Like [`Self::per_stream`], but an exhausted script resets the stream, as a peer + /// that fails after replying would. + pub fn per_stream_reset(scripts: Vec>) -> Self { + Self { + reset: true, + ..Self::per_stream(scripts) + } + } + /// Append to the shared script: the peer sending more on a stream it already opened. /// Nothing is woken, so the test re-polls the reader itself. pub fn push(&self, bytes: &[u8]) { @@ -689,6 +704,7 @@ impl poll::Session for ScriptedSession { Poll::Ready(Ok(ScriptedRecv { script: Arc::new(Mutex::new(script)), eof: self.eof, + reset: self.reset, log: self.log.clone(), })) } @@ -720,6 +736,7 @@ impl poll::Session for ScriptedSession { ScriptedRecv { script, eof: self.eof, + reset: self.reset, log: self.log.clone(), }, ))) diff --git a/rs/moq-net/src/lite/track.rs b/rs/moq-net/src/lite/track.rs index 68f2e10d48..1d4ad41da0 100644 --- a/rs/moq-net/src/lite/track.rs +++ b/rs/moq-net/src/lite/track.rs @@ -133,7 +133,7 @@ mod test { let mut buf = Vec::new(); info.encode(&mut buf, Version::Lite05).unwrap(); - assert_eq!(buf, [0x06, 0x00, 0x00, 0x53, 0x88, 0x43, 0xe8]); + assert_eq!(buf, [0x06, 0x7f, 0x00, 0x53, 0x88, 0x43, 0xe8]); } #[test] diff --git a/rs/moq-net/src/model/bandwidth.rs b/rs/moq-net/src/model/bandwidth.rs index 735ab78339..58977552ea 100644 --- a/rs/moq-net/src/model/bandwidth.rs +++ b/rs/moq-net/src/model/bandwidth.rs @@ -232,9 +232,9 @@ impl Allocator { /// decision about what to *produce*, and there is no single subscriber /// priority to read when several are watching one track. /// - /// That last part is what carries the common case, since publishers leave - /// `priority` at its default today: one tier of audio and video still serves - /// audio's small reservation in full before video takes the remainder. + /// That last part carries a publisher that leaves `priority` at its default: + /// one tier of audio and video still serves audio's small reservation in full + /// before video takes the remainder. /// /// The reservation lasts as long as the returned [`Reservation`]: hold it for as /// long as the sender is publishing, change the ceiling with @@ -663,9 +663,8 @@ mod tests { assert_eq!(allocate(bps(1_000_000), &wants, 1), Some(bps(0))); } - /// Publishers don't set [`track::Info::priority`] today (it defaults to 0 and - /// `hang::container::track_info` leaves it there), so audio and video land in - /// one tier. That has to come out right anyway, and it does: max-min fair + /// A publisher that doesn't set [`track::Info::priority`] puts audio and video + /// in one tier. That has to come out right anyway, and it does: max-min fair /// satisfies the small claim first, so audio still gets its full reservation /// and video takes the rest. Priority only changes the answer once a tier's /// smaller claims outgrow an even split. diff --git a/rs/moq-net/src/model/datagram.rs b/rs/moq-net/src/model/datagram.rs index f305507207..a8c20b0979 100644 --- a/rs/moq-net/src/model/datagram.rs +++ b/rs/moq-net/src/model/datagram.rs @@ -8,10 +8,10 @@ //! //! Delivery is best-effort per hop: a session drops (with a debug log) any datagram whose encoded //! body exceeds the transport's datagram size, and sessions that can't carry datagrams at all -//! (IETF moq-transport, moq-lite before 05, or stream-only transports like WebSocket) never -//! deliver them. +//! (moq-lite before 05, or stream-only transports like WebSocket) never deliver them. //! -//! Wire counterpart: [`crate::lite::Datagram`]. +//! Wire counterparts: [`crate::lite::Datagram`], and on moq-transport an OBJECT_DATAGRAM at +//! object 0 whose Group ID is the sequence ([`crate::ietf::ObjectDatagram`]). use bytes::Bytes; diff --git a/rs/moq-net/src/model/resume.rs b/rs/moq-net/src/model/resume.rs index 089c361faf..5e797e0835 100644 --- a/rs/moq-net/src/model/resume.rs +++ b/rs/moq-net/src/model/resume.rs @@ -727,6 +727,13 @@ impl Consumer { self.state.read().latest() } + /// The newest segment's declared exclusive end, where fetches are routed. + pub(crate) fn final_sequence(&self) -> Option { + // Copied out: the segment's track takes its own lock. + let track = self.state.read().segments.last().map(|segment| segment.track.clone())?; + track.final_sequence() + } + /// One past the newest position across the segments: where a route taking this /// logical track over would resume. pub(crate) fn resume_position(&self) -> Option { @@ -843,11 +850,13 @@ impl kio::Pollable for Fetching { type Output = Result; fn poll(&self, waiter: &kio::Waiter) -> Poll { - if let Some(group) = (Consumer { + if let Some(mut group) = (Consumer { state: self.state.clone(), }) .cached_group(self.sequence, self.options.frame_start) { + // Sitting where the caller asked, as a fetch from the segment's own track would. + group.start_at(self.options.frame_start); return Poll::Ready(Ok(group)); } @@ -2928,6 +2937,33 @@ mod test { assert_eq!(read(&mut group), b"b4"); } + /// A cached copy is handed back sitting at the requested frame, the same as a fetch + /// the segment's own track answers. + #[tokio::test] + async fn a_cached_fetch_starts_at_the_requested_frame() { + let (track_a, consumer_a) = track_pair("a"); + let mut producer = Producer::new(); + producer.switch(&consumer_a, None).unwrap(); + + let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap(); + for payload in ["f0", "f1", "f2"] { + group + .write_frame(crate::Timestamp::ZERO, payload.as_bytes().to_vec()) + .unwrap(); + } + group.finish().unwrap(); + + let mut group = producer + .consume() + .fetch_group(0, group::Fetch::default().with_frame_start(1)) + .now_or_never() + .expect("cached fetch should resolve") + .unwrap(); + assert_eq!(group.index(), 1); + let frame = group.read_frame().now_or_never().unwrap().unwrap().unwrap(); + assert_eq!(frame.payload.as_ref(), b"f1"); + } + #[tokio::test] async fn fetch_waits_for_first_segment() { let (mut track_a, consumer_a) = track_pair("a"); diff --git a/rs/moq-net/src/model/track.rs b/rs/moq-net/src/model/track.rs index 49aa3ba2ec..c3bc6f17b0 100644 --- a/rs/moq-net/src/model/track.rs +++ b/rs/moq-net/src/model/track.rs @@ -35,6 +35,10 @@ use std::{ /// Default [`Info::max_age`] when the publisher doesn't set one. pub const DEFAULT_MAX_AGE: Duration = Duration::from_secs(5); +// The higher-first midpoint. IETF flips priority (lower first), so this goes out as 128, the +// draft's usual publisher priority, while moq-lite carries 127 as written: one urgency on both. +const DEFAULT_PRIORITY: u8 = 127; + /// Maximum number of datagrams retained in the per-track send buffer. /// /// Datagrams are a best-effort send buffer, not a replay cache (unlike groups): only the last @@ -106,6 +110,7 @@ pub struct Info { pub max_age: Duration, /// The publisher's priority for this track, used only to break ties between /// subscriptions of equal subscriber priority. Reported in TRACK_INFO (Lite05+). + /// Higher is more urgent. Defaults to 127, the midpoint. pub priority: u8, } @@ -114,7 +119,7 @@ impl Default for Info { Self { timescale: Timescale::default(), max_age: DEFAULT_MAX_AGE, - priority: 0, + priority: DEFAULT_PRIORITY, } } } @@ -1374,9 +1379,8 @@ impl Producer { /// track's groups but drawing from the same sequence namespace (so interleaving with /// [`Self::append_group`] never reuses a number). There is no group fallback: each /// session drops (with a debug log) any datagram whose encoded body exceeds the - /// transport's datagram size, and sessions that can't carry datagrams at all (IETF - /// moq-transport, moq-lite before 05, or stream-only transports like WebSocket) never - /// deliver them. Keep payloads well under the 1200-byte minimum path MTU. An origin + /// transport's datagram size, and sessions that can't carry datagrams at all (moq-lite + /// before 05, or stream-only transports like WebSocket) never deliver them. Keep payloads well under the 1200-byte minimum path MTU. An origin /// publisher uses this; a relay preserving upstream numbering uses /// [`Self::insert_datagram`]. pub fn append_datagram(&mut self, timestamp: Timestamp, payload: B) -> Result { @@ -2229,7 +2233,11 @@ impl Demand { /// The publisher's tie-break priority, as set in [`Info::priority`]. pub(crate) fn priority(&self) -> u8 { // Always Some once the track exists; a closed one reads its last value. - self.state.read().info.as_ref().map_or(0, |info| info.priority) + self.state + .read() + .info + .as_ref() + .map_or(DEFAULT_PRIORITY, |info| info.priority) } /// Whether anyone is subscribed right now, without waiting. @@ -2788,6 +2796,16 @@ impl Consumer { } } + /// The declared exclusive final sequence, or `None` while the track is open ended. + /// + /// A spliced track answers for its newest segment, which is where fetches go. + pub(crate) fn final_sequence(&self) -> Option { + match &self.inner { + ConsumerKind::Plain(state) => state.read().final_sequence, + ConsumerKind::Spliced(resume) => resume.final_sequence(), + } + } + /// The frame-precise point a replacement route should resume from: one past the /// last frame this copy produced. `None` if it produced nothing. /// @@ -2982,6 +3000,16 @@ impl group::Request { res } + /// Declare the track's exclusive final sequence, as the publisher answering this + /// fetch reported it. A no-op once the track declared one, or holds a later group. + pub(crate) fn finish_track_at(&self, final_sequence: u64) { + if let Ok(mut state) = TrackState::modify(&self.state) + && state.final_sequence.is_none() + { + let _ = state.set_final(final_sequence); + } + } + /// Reject the fetch, resolving every joined [`Consumer::fetch_group`] with `err`. pub fn reject(mut self, err: Error) { self.done = true; diff --git a/rs/moq-net/tests/datagram.rs b/rs/moq-net/tests/datagram.rs index 4cc78b1297..d974cc580e 100644 --- a/rs/moq-net/tests/datagram.rs +++ b/rs/moq-net/tests/datagram.rs @@ -1,4 +1,4 @@ -//! MoQ Lite datagram delivery over the in-memory mock transport. +//! Datagram delivery over the in-memory mock transport, on moq-lite and moq-transport. //! //! Covers the whole receive path end to end: publisher encoding, the transport's //! datagram channel, the subscriber's receive loop and `route_datagram`, and @@ -9,7 +9,6 @@ mod support; use std::time::Duration; -use futures::FutureExt as _; use moq_net::{Hop, Timestamp, Version}; use support::harness::{MockConnectOptions, MockPair, connect_mock}; @@ -110,53 +109,74 @@ async fn datagrams_reach_the_subscriber_in_order() { .expect("timed out"); } -/// MoQ Transport has no datagram mapping: groups still flow, inserted datagrams do not. +/// MoQ Transport carries a datagram as an OBJECT_DATAGRAM at object 0, so it arrives as a +/// datagram with its sequence, alongside the groups on streams. #[tokio::test] -async fn ietf_does_not_deliver_datagrams() { - tokio::time::timeout(TEST_TIMEOUT, async { - let publisher = produce_origin(1); - let consumer_origin = produce_origin(2); +async fn ietf_delivers_datagrams() { + for version in [ + "moq-transport-14", + "moq-transport-16", + "moq-transport-17", + "moq-transport-20", + ] { + tokio::time::timeout(TEST_TIMEOUT, ietf_delivers_datagrams_on(version)) + .await + .unwrap_or_else(|_| panic!("{version}: timed out")); + } +} - let broadcast = publisher.create_broadcast("bench").unwrap(); - let mut producer = broadcast.create_track("datagrams", None).unwrap(); - broadcast.announce(Default::default()).unwrap(); +async fn ietf_delivers_datagrams_on(version: &str) { + let publisher = produce_origin(1); + let consumer_origin = produce_origin(2); - let mut options = MockConnectOptions::new("moq-transport-19".parse::().unwrap()); - options.server_publish = Some(publisher.consume()); - options.client_subscribe = Some(consumer_origin.clone()); - let _pair = connect_mock(options).await; + let broadcast = publisher.create_broadcast("bench").unwrap(); + let mut producer = broadcast.create_track("datagrams", None).unwrap(); + broadcast.announce(Default::default()).unwrap(); - let consumer = consumer_origin.consume(); - consumer.routed("bench").await.unwrap(); - let remote = consumer.request_broadcast("bench").await.unwrap(); - let mut subscriber = remote.track("datagrams").unwrap().subscribe(None).await.unwrap(); + let mut options = MockConnectOptions::new(version.parse::().unwrap()); + options.server_publish = Some(publisher.consume()); + options.client_subscribe = Some(consumer_origin.clone()); + let _pair = connect_mock(options).await; - producer - .write_frame(Timestamp::from_millis(0).unwrap(), &b"before"[..]) - .unwrap(); - let before = subscriber.recv_group().await.unwrap().unwrap(); - assert_eq!(before.sequence, 0); + let consumer = consumer_origin.consume(); + consumer.routed("bench").await.unwrap(); + let remote = consumer.request_broadcast("bench").await.unwrap(); + let mut subscriber = remote.track("datagrams").unwrap().subscribe(None).await.unwrap(); - producer - .insert_datagram( - 0, - Timestamp::from_millis(7).unwrap(), - bytes::Bytes::from_static(PAYLOAD), - ) - .unwrap(); - producer - .write_frame(Timestamp::from_millis(1).unwrap(), &b"after"[..]) - .unwrap(); - let after = subscriber.recv_group().await.unwrap().unwrap(); - assert_eq!(after.sequence, 1); + // A group first, so the subscription's alias is bound before any datagram lands. + producer + .write_frame(Timestamp::from_millis(0).unwrap(), &b"before"[..]) + .unwrap(); + let mut before = subscriber.recv_group().await.unwrap().unwrap(); + assert_eq!(before.sequence, 0, "{version}"); + // Drafts without a track timescale stamp objects on arrival, datagrams included. + let stamped = before.next_frame().await.unwrap().unwrap().timestamp.value() == 0; - assert!( - subscriber.recv_datagram().now_or_never().is_none(), - "MoQ Transport must not map datagrams" + producer + .insert_datagram( + 5, + Timestamp::from_millis(7).unwrap(), + bytes::Bytes::from_static(PAYLOAD), + ) + .unwrap(); + let datagram = subscriber.recv_datagram().await.unwrap().unwrap(); + assert_eq!(datagram.sequence, 5, "{version}: the relay must not renumber"); + assert_eq!(&datagram.payload[..], PAYLOAD, "{version}"); + if stamped { + let expected = Timestamp::from_millis(7).unwrap(); + assert_eq!( + datagram.timestamp.convert(expected.scale()).unwrap(), + expected, + "{version}" ); - }) - .await - .expect("timed out"); + } + + // The group sequence continues past the datagram. + producer + .write_frame(Timestamp::from_millis(8).unwrap(), &b"after"[..]) + .unwrap(); + let after = subscriber.recv_group().await.unwrap().unwrap(); + assert_eq!(after.sequence, 6, "{version}"); } /// Explicit insert keeps the origin sequence on the lite wire, including a gap. diff --git a/rs/moq-tokio/tests/broadcast.rs b/rs/moq-tokio/tests/broadcast.rs index 23a16c1fde..2ab11f6844 100644 --- a/rs/moq-tokio/tests/broadcast.rs +++ b/rs/moq-tokio/tests/broadcast.rs @@ -376,6 +376,143 @@ async fn broadcast_moq_lite_05_fetch_webtransport() { lite05_fetch_roundtrip("https").await; } +/// A cache miss over moq-transport is a standalone FETCH of the one group, served from +/// the publisher's cache, and a group the publisher lacks comes back as the publisher's +/// own refusal. Without `served`, the version has no FETCH and the miss is refused. +async fn transport_fetch_roundtrip(version: &str, served: bool) { + let pub_origin = moq_tokio::origin::spawn(); + let broadcast = pub_origin.create_broadcast("test").expect("failed to create broadcast"); + broadcast + .announce(Default::default()) + .expect("failed to announce broadcast"); + let track = broadcast.create_track("video", None).expect("failed to create track"); + for sequence in 0..3u64 { + let mut group = track.append_group().expect("failed to append group"); + for frame in 0..2 { + group + .write_frame( + moq_net::Timestamp::ZERO, + bytes::Bytes::from(format!("{sequence}-{frame}")), + ) + .expect("failed to write frame"); + } + group.finish().expect("failed to finish group"); + } + + let mut server_config = moq_tokio::listen::Config::default(); + server_config.bind = Some("[::]:0".parse().unwrap()); + server_config.tls.generate = vec!["localhost".into()]; + server_config.version = vec![version.parse().unwrap()]; + let server = server_config.init(Default::default()).expect("failed to init server"); + let mut server = server.listen().await.expect("failed to listen"); + let addr = server.local_addr().expect("failed to get local addr"); + + let sub_origin = moq_tokio::origin::spawn(); + let sub_consumer = sub_origin.consume(); + let mut announcements = sub_consumer.announced(); + + let mut client_config = moq_tokio::connect::Config::default(); + client_config.tls.insecure = Some(true); + client_config.version = vec![version.parse().unwrap()]; + let client = client_config.init(Default::default()).expect("failed to init client"); + let url: url::Url = format!("moqt://localhost:{}", addr.port()).parse().unwrap(); + + let server_handle = tokio::spawn(async move { + let request = server.accept().await.expect("no incoming connection"); + let session = request.with_publisher(&pub_origin).ok().await?; + let _broadcast = broadcast; + let _track = track; + let _ = session.closed().await; + Ok::<_, anyhow::Error>(()) + }); + + let client = client.with_subscriber(sub_origin); + let (_client, connection) = tokio::time::timeout(TIMEOUT, connect_once(client, url)) + .await + .expect("client connect timed out") + .expect("client connect failed"); + + tokio::time::timeout(TIMEOUT, announcements.next()) + .await + .expect("announce timed out") + .expect("origin closed"); + let bc = tokio::time::timeout(TIMEOUT, sub_consumer.request_broadcast("test")) + .await + .expect("request timed out") + .expect("announced broadcast resolves"); + + let fetched = tokio::time::timeout(TIMEOUT, async { bc.track("video").unwrap().fetch_group(1, None).await }) + .await + .expect("fetch timed out"); + if !served { + assert!( + matches!(fetched, Err(moq_net::Error::Unsupported)), + "expected an unsupported FETCH, got {:?}", + fetched.err() + ); + drop(connection); + server_handle + .await + .expect("server task panicked") + .expect("server task failed"); + return; + } + let mut group = fetched.expect("fetch failed"); + assert_eq!(group.sequence, 1); + for frame in 0..2 { + let payload = tokio::time::timeout(TIMEOUT, group.read_frame()) + .await + .expect("read timed out") + .expect("read failed") + .expect("group ended early"); + assert_eq!(payload.payload, bytes::Bytes::from(format!("1-{frame}"))); + } + let end = tokio::time::timeout(TIMEOUT, group.read_frame()) + .await + .expect("read timed out") + .expect("read failed"); + assert!(end.is_none(), "the group ends after its frames"); + + let refused = tokio::time::timeout(TIMEOUT, async { bc.track("video").unwrap().fetch_group(7, None).await }) + .await + .expect("fetch timed out"); + assert!( + matches!(refused, Err(moq_net::Error::NotFound)), + "expected the publisher's refusal, got {:?}", + refused.err() + ); + + drop(connection); + server_handle + .await + .expect("server task panicked") + .expect("server task failed"); +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_14() { + transport_fetch_roundtrip("moq-transport-14", true).await; +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_16() { + transport_fetch_roundtrip("moq-transport-16", true).await; +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_18() { + transport_fetch_roundtrip("moq-transport-18", true).await; +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_20() { + transport_fetch_roundtrip("moq-transport-20", false).await; +} + /// A fetch must be served while a live subscription is active on the same track. /// The relay subscribes starting at the latest group, so an older group isn't /// cached and the fetch has to issue a wire FETCH concurrently with the diff --git a/swift/README.md b/swift/README.md index 613831a581..3317296c37 100644 --- a/swift/README.md +++ b/swift/README.md @@ -72,7 +72,7 @@ Incoming server `Request` values expose the query-free `path` before acceptance. Raw tracks also expose best-effort datagrams: `TrackProducer.appendDatagram(_:timestampUs:)` returns the assigned sequence number, `TrackConsumer.recvDatagram()` receives one datagram, and `TrackConsumer.datagrams` streams them in arrival order. Payloads are capped at 1200 bytes. -Datagrams require a datagram-capable transport and lite-05 or newer moq-lite; IETF moq-transport, +Datagrams require a datagram-capable transport and lite-05 or newer moq-lite, or moq-transport; pre-lite-05, WebSocket, and TCP paths do not deliver them, and there is no stream fallback. JSON tracks carry your own `Codable` types with the framing handled for you. You opt into one of two