From cde18e2e91855d6aa6d1c70b8f3d2b80d2eda0e2 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Wed, 7 Oct 2026 02:02:05 -0700 Subject: [PATCH 1/4] fix(net): deliver a fresh group below the first served group when the floor allows it A pre-06 subscribe that named no group still pins its floor to the first group served. An explicit floor, including the group 0 that lite-06 encodes as 0, keeps that floor, so a later group inside the max age is delivered. An omitted floor no longer clears an explicit one when subscriptions are aggregated. Co-authored-by: Grok 4.7 --- drafts/draft-lcurley-moq-lite.md | 13 +- js/net/src/integration.test.ts | 6 +- js/net/src/lite/publisher.test.ts | 49 ++++- js/net/src/lite/publisher.ts | 24 ++- js/net/src/lite/subscribe.test.ts | 37 +++- js/net/src/lite/subscribe.ts | 9 +- js/net/src/lite/subscriber.ts | 3 + js/net/src/track.test.ts | 12 ++ js/net/src/track.ts | 28 ++- quest/m1/README.md | 1 - quest/m1/lite-late-lower-group.md | 63 ------ rs/moq-net/src/ietf/subscriber.rs | 35 ++- rs/moq-net/src/lite/publisher.rs | 77 +++++-- rs/moq-net/src/lite/subscribe.rs | 34 ++- rs/moq-net/src/model/subscription.rs | 57 ++--- rs/moq-net/tests/history_groups.rs | 5 +- rs/moq-net/tests/late_lower_group.rs | 307 +++++++++++++++++++++++++++ 17 files changed, 597 insertions(+), 163 deletions(-) delete mode 100644 quest/m1/lite-late-lower-group.md create mode 100644 rs/moq-net/tests/late_lower_group.rs diff --git a/drafts/draft-lcurley-moq-lite.md b/drafts/draft-lcurley-moq-lite.md index e9bccd9014..8f61a5dc9f 100644 --- a/drafts/draft-lcurley-moq-lite.md +++ b/drafts/draft-lcurley-moq-lite.md @@ -1069,7 +1069,7 @@ A value of 0 means the whole group (default). A non-zero value is the absolute frame index + 1, matching `Group End`. MUST be 0 when `Group End` is 0, since an unbounded subscription has no end group to qualify. -`Group Start` and `Group End` are offset by 1 only so 0 can mean "absent"; every other group field in this document is a plain absolute sequence. +`Group End` and `Frame End` are offset by 1 so 0 can mean "absent". `Group Start` is an absolute sequence, and 0 is group 0. ## SUBSCRIBE_UPDATE A subscriber can modify a subscription with a SUBSCRIBE_UPDATE message. @@ -1178,9 +1178,12 @@ SUBSCRIBE_OK Message { Set to 0x0 to indicate a SUBSCRIBE_OK message. **Group**: -The absolute sequence number of the first group that will be delivered. -It MUST be greater than or equal to the requested start group; any groups in between are unavailable. -A subscriber that requested the latest group learns the resolved sequence here. +The absolute sequence number of the first group served when the subscription resolves. +It MUST be greater than or equal to the requested `Group Start`. +Groups between the requested floor and this group are unavailable. +A group at or above the requested floor that is still within `Subscriber Max Age`, created or becoming available later, is delivered. +This group is not itself a new floor. +A subscriber whose `Subscriber Max Age` resolves the start to the latest group learns that sequence here. There is no matching frame field, because the start frame is never in doubt: a partial group is only delivered when it was asked for, so the subscription starts either exactly where it asked or at the beginning of a later group (see [Positions](#positions)). The subscriber derives the start frame from `Group` and its own request: @@ -1380,6 +1383,8 @@ The `Message Length` describes the payload size on the wire. ## moq-lite-07 +- SUBSCRIBE_OK `Group` names the first group served when the subscription resolves. Groups between the requested floor and that group are unavailable. A later group at or above the floor and within Subscriber Max Age is still delivered. The announced group is not a new floor. +- Corrected the SUBSCRIBE note that offset `Group Start` by 1. Only `Group End` and `Frame End` are offset so 0 can mean absent. `Group Start` is an absolute sequence, and 0 is group 0. - The subscriber FINs its Subscribe Stream after settling its tail; graceful session close waits for that FIN or reset. - A refusal is not retried at another route of the same prefix either. - Made TRACK_INFO Publisher Max Age optional, encoded as milliseconds plus one with zero meaning no limit. diff --git a/js/net/src/integration.test.ts b/js/net/src/integration.test.ts index 3811b540b8..a411da6fde 100644 --- a/js/net/src/integration.test.ts +++ b/js/net/src/integration.test.ts @@ -166,8 +166,7 @@ test("integration: lite subscription options and updates reach the publisher", a })(); const remote = wireOf(client).consume(Path.from("test")); - // A floor of 1, not 0: a pre-06 wire folds a vacuous floor of 0 back to absent, since - // its encoding of group 0 means "replay from the beginning" instead. + // A floor of group 1, a concrete group on every draft. const subscriber = remote.track("video").subscribe({ priority: 3, maxAge: Milli(250), @@ -225,8 +224,7 @@ test("integration: lite carries a fractional maxAge as a whole millisecond", asy const remote = wireOf(client).consume(Path.from("test")); // A varint cannot encode 38.75, so an unrounded value fails the SUBSCRIBE outright and - // nothing resubscribes. The publisher must see the budget rounded up instead. A floor - // of 1, not 0: a pre-06 wire folds a vacuous floor of 0 back to absent. + // nothing resubscribes. The publisher must see the budget rounded up instead. const subscriber = remote.track("video").subscribe({ maxAge: Milli(38.75), groups: { start: { included: 1 } } }); // A failed subscribe never reaches the publisher, so race its closure to report the diff --git a/js/net/src/lite/publisher.test.ts b/js/net/src/lite/publisher.test.ts index 81492ed6a9..db34a12f1f 100644 --- a/js/net/src/lite/publisher.test.ts +++ b/js/net/src/lite/publisher.test.ts @@ -643,9 +643,9 @@ test("lite draft-05: a late-arriving older group is still served", async () => { } }); -// SUBSCRIBE_START promises nothing below the announced sequence will be delivered, so -// the floor is pinned there: a straggler below the first served group is dropped even -// though arrival-order serving would otherwise surface it. +// An unfloored pre-06 subscribe joins at the publisher's start. SUBSCRIBE_START makes +// that group the floor: a straggler below it is dropped even though arrival-order +// serving would otherwise surface it. test("lite draft-05: a straggler below the announced start group is not served", async () => { const sub = await servedSubscription(); try { @@ -667,6 +667,39 @@ test("lite draft-05: a straggler below the announced start group is not served", } }); +// An explicit floor stays where it was named. SUBSCRIBE_START still reports the first +// served group, and a later group at or above the floor is delivered. +test("lite draft-05: an explicit floor still serves a group below the first one", async () => { + const sub = await servedSubscription({ startGroup: 0 }); + try { + sub.serve(2); + expect(await sub.servedSequence()).toBe(2); + + const resp = await decodeSubscribeResponse(sub.client.reader, Version.DRAFT_05); + if (!("start" in resp)) throw new Error("expected SUBSCRIBE_START"); + expect(resp.start.group).toBe(2); + + sub.serve(0); + expect(await sub.servedSequence()).toBe(0); + } finally { + await sub.close(); + } +}); + +// Lite-06 encodes a floor of group 0 as 0, including a subscribe that omitted Group Start. +// That is a floor, so a later group 0 is still served. +test("lite draft-06: an omitted group start still serves a group below the first one", async () => { + const sub = await servedSubscription({ version: Version.DRAFT_06 }); + try { + sub.serve(2); + expect(await sub.servedSequence()).toBe(2); + sub.serve(0); + expect(await sub.servedSequence()).toBe(0); + } finally { + await sub.close(); + } +}); + // A group pop is the linearization point. Once the publisher reaches the SUBSCRIBE_START write, // the group and its range have already been decided, so an update queued during the write applies // to the next pop rather than reaching backward into this one. @@ -708,8 +741,9 @@ test("lite draft-06: a queued floor update applies before the next group pop", a await replayUpdate({ priority: 0, startGroup: 2 }).encode(sub.client.writer, Version.DRAFT_06); sub.serve(2); await flush(); - expect(ranges).toHaveBeenCalledTimes(1); - expect(lastGroups(ranges)).toEqual({ start: { included: 0 }, end: undefined }); + // Draft-06 group 0 is a floor, so the first group does not raise the cursor. + // The queued update is still waiting on the parked write. + expect(ranges).not.toHaveBeenCalled(); sub.release(); @@ -760,8 +794,9 @@ test("lite draft-06: a queued frame floor applies before the next group pop", as await replayUpdate({ priority: 0, startGroup: 1, startFrame: 2 }).encode(sub.client.writer, Version.DRAFT_06); await flush(); - expect(ranges).toHaveBeenCalledTimes(1); - expect(lastGroups(ranges)).toEqual({ start: { included: 0 }, end: undefined }); + // Draft-06 group 0 is a floor, so serving the first group does not move the cursor. + // The queued update is still waiting on the parked write. + expect(ranges).not.toHaveBeenCalled(); sub.release(); diff --git a/js/net/src/lite/publisher.ts b/js/net/src/lite/publisher.ts index 09a67a90a4..f1c7c8c999 100644 --- a/js/net/src/lite/publisher.ts +++ b/js/net/src/lite/publisher.ts @@ -591,11 +591,14 @@ export class Publisher { } const endGroup = exclusiveGroupEnd(msg.endGroup); + // Lite-06 encodes an omitted floor as 0, and that 0 is group 0. Store it as an + // explicit floor so it widens a higher resume. A pre-06 absent start stays absent. + const startGroup = msg.startGroup === undefined && resolvesStart(this.version) ? 0 : msg.startGroup; const track = wireOf(front).subscribe(msg.track, { priority: msg.priority, maxAge: Milli(servingMaxAge(this.version, msg.maxAge)), groups: { - start: msg.startGroup === undefined ? undefined : { included: msg.startGroup }, + start: startGroup === undefined ? undefined : { included: startGroup }, end: endGroup === undefined ? undefined : { excluded: endGroup }, }, }); @@ -917,13 +920,18 @@ export class Publisher { if (emitRange && !startSent) { startSent = true; - // SUBSCRIBE_START promises nothing below this sequence will be delivered. - // Arrival-order serving could later surface a straggler below the first - // group, so pin the floor to what was announced. - hooks.replaceGroups(track, { - start: { included: group.sequence }, - end: bounds.endGroup === undefined ? undefined : { included: bounds.endGroup }, - }); + // SUBSCRIBE_START names the first group served now. A later group at or + // above an explicit floor is still delivered, so the cursor stays put. + // A pre-06 subscribe that named no group is the exception: an absent + // Group Start is the latest group, and this sequence becomes the floor. + // Lite-06 encodes a floor of group 0 as 0, so an omitted start is not pinned. + const pin = bounds.startGroup === undefined && !resolvesStart(this.version); + if (pin || bounds.endGroup !== undefined) { + hooks.replaceGroups(track, { + ...(pin ? { start: { included: group.sequence } } : {}), + ...(bounds.endGroup === undefined ? {} : { end: { included: bounds.endGroup } }), + }); + } if ( !(await controls.response( encodeSubscribeResponse( diff --git a/js/net/src/lite/subscribe.test.ts b/js/net/src/lite/subscribe.test.ts index 79c6737f7a..8bcf18947f 100644 --- a/js/net/src/lite/subscribe.test.ts +++ b/js/net/src/lite/subscribe.test.ts @@ -111,14 +111,23 @@ test("Subscribe round-trips every option including startGroup 0", async () => { expect(resumed.startFrame).toBe(4); message.startFrame = 0; - // A pre-06 wire folds the vacuous floor back to absent: an explicit group 0 there - // would mean "replay from the beginning", which is not what a floor of 0 asks for. - const folded = await Subscribe.decode( + // A pre-06 wire encodes an explicit group 0 as sequence + 1, so it round-trips. + // Omitting the floor is still the latest group, not group 0. + const explicit = await Subscribe.decode( new Reader(undefined, await encodeMessage(Version.DRAFT_05, message), Version.DRAFT_05), Version.DRAFT_05, ); - expect(folded.startGroup).toBeUndefined(); - expect(folded.endGroup).toBe(9); + expect(explicit.startGroup).toBe(0); + expect(explicit.endGroup).toBe(9); + const absent = await Subscribe.decode( + new Reader( + undefined, + await encodeMessage(Version.DRAFT_05, new Subscribe({ ...message, startGroup: undefined })), + Version.DRAFT_05, + ), + Version.DRAFT_05, + ); + expect(absent.startGroup).toBeUndefined(); }); test("SubscribeUpdate round-trips every option including startGroup 0", async () => { @@ -137,13 +146,23 @@ test("SubscribeUpdate round-trips every option including startGroup 0", async () expect(got.startGroup).toBeUndefined(); expect(got.endGroup).toBe(12); - // The same fold as SUBSCRIBE on a pre-06 wire. - const folded = await SubscribeUpdate.decode( + // The same pre-06 encoding as SUBSCRIBE: explicit group 0 round-trips, and an + // omitted floor stays omitted. + const explicit = await SubscribeUpdate.decode( new Reader(undefined, await encodeMessage(Version.DRAFT_05, message), Version.DRAFT_05), Version.DRAFT_05, ); - expect(folded.startGroup).toBeUndefined(); - expect(folded.endGroup).toBe(12); + expect(explicit.startGroup).toBe(0); + expect(explicit.endGroup).toBe(12); + const absent = await SubscribeUpdate.decode( + new Reader( + undefined, + await encodeMessage(Version.DRAFT_05, new SubscribeUpdate({ ...message, startGroup: undefined })), + Version.DRAFT_05, + ), + Version.DRAFT_05, + ); + expect(absent.startGroup).toBeUndefined(); }); test("SubscribeStart round-trips on draft-05", async () => { diff --git a/js/net/src/lite/subscribe.ts b/js/net/src/lite/subscribe.ts index 084f68f40e..fa3b4aac2b 100644 --- a/js/net/src/lite/subscribe.ts +++ b/js/net/src/lite/subscribe.ts @@ -9,17 +9,16 @@ import { hasFrameBounds, hasGroupOrder, hasLargest, hasStreamCount, resolvesStar /** * Encode the `Group Start` field shared by SUBSCRIBE and SUBSCRIBE_UPDATE. * - * Lite-06 writes the raw floor (`undefined` and 0 are the same absence of a constraint), - * while a pre-06 wire encodes the sequence + 1 and gets a vacuous floor folded back to - * absent: an explicit group 0 there means "replay from the beginning", which is not what - * a floor of 0 asks for. + * Lite-06 writes the raw floor (`undefined` and 0 are both group 0). A pre-06 wire + * encodes the sequence + 1: omitting the floor is 0 (the latest group, where the + * publisher starts) and an explicit 0 is 1, replay from the beginning. */ async function encodeStartGroup(w: Writer, version: Version, startGroup?: number) { if (resolvesStart(version)) { await w.u53(startGroup ?? 0); return; } - await w.u53(startGroup !== undefined && startGroup > 0 ? startGroup + 1 : 0); + await w.u53(startGroup === undefined ? 0 : startGroup + 1); } /** diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index eb2bae2d21..309aeac2b1 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -1042,9 +1042,12 @@ export class Subscriber { #sameSubscription(a: track.Subscription, b: track.Subscription): boolean { const ag = groupBounds(a.groups); const bg = groupBounds(b.groups); + // `groupBounds` reads an omitted start as 0. On a pre-06 wire those are different + // subscriptions: omitted joins at the publisher's start, and 0 is group 0. return ( (a.priority ?? 0) === (b.priority ?? 0) && (a.maxAge ?? 0) === (b.maxAge ?? 0) && + (a.groups?.start === undefined) === (b.groups?.start === undefined) && ag.start === bg.start && ag.end === bg.end ); diff --git a/js/net/src/track.test.ts b/js/net/src/track.test.ts index 5900603037..6f2a4ad8b1 100644 --- a/js/net/src/track.test.ts +++ b/js/net/src/track.test.ts @@ -359,6 +359,18 @@ test("multiple subscriber options aggregate like Rust", async () => { expect(await none).toBeUndefined(); }); +test("an omitted floor leaves another subscriber's floor in place", () => { + const producer = new TrackProducer("test"); + producer.subscribe({ groups: { start: { included: 10 } } }); + producer.subscribe({}); + expect(producer.subscription.peek()?.groups?.start).toEqual({ included: 10 }); + + const unfloored = new TrackProducer("test"); + unfloored.subscribe({}); + unfloored.subscribe({}); + expect(unfloored.subscription.peek()?.groups?.start).toBeUndefined(); +}); + test("the producer aggregate is clamped without changing subscriber options", async () => { const producer = new TrackProducer("test"); const track = producer.subscribe({ maxAge: Milli(10_000) }); diff --git a/js/net/src/track.ts b/js/net/src/track.ts index 1ad168ad06..6045680a55 100644 --- a/js/net/src/track.ts +++ b/js/net/src/track.ts @@ -139,9 +139,14 @@ export interface Subscription { * The lowest group the publisher may deliver (a floor), or omit for none. * * A floor, not a request: only {@link maxAge} asks for data, and the floor bounds how - * far back it may reach. Omitting it and a floor of 0 mean the same thing, and a floor - * above the live edge simply waits there (a resumed subscription naming where it left - * off). + * far back it may reach. Omitting it joins where the publisher starts: the first group + * served becomes the floor, so a group created below it is not delivered. A floor of 0 + * is explicit, and a later group at or above it is still delivered while {@link maxAge} + * considers it fresh. A floor above the live edge simply waits there (a resumed + * subscription naming where it left off). + * + * On lite-06 and later the wire's 0 is group 0, so an omitted floor and a floor of 0 + * are one subscription once they are encoded. */ groups?: Groups; } @@ -174,15 +179,22 @@ function combineSubscriptions(states: Iterable): Subscription | unde combined.priority = Math.max(combined.priority ?? 0, subscription.priority ?? 0); combined.maxAge = Milli(Math.max(combined.maxAge ?? Milli.zero, subscription.maxAge ?? Milli.zero)); - // A floor only restricts, so a subscriber without one clears the aggregate: - // its budget may reach below any floor the others set. + // An omitted floor joins where the publisher starts. It does not clear an explicit + // floor: the loosest explicit one wins, and each subscriber's cursor still filters + // what it reads. const a = groupBounds(combined.groups ?? {}); const b = groupBounds(subscription.groups ?? {}); + const aStart = combined.groups?.start; + const bStart = subscription.groups?.start; combined.groups = { start: - combined.groups?.start === undefined || subscription.groups?.start === undefined - ? undefined - : { included: Math.min(a.start, b.start) }, + aStart === undefined + ? bStart === undefined + ? undefined + : { included: b.start } + : bStart === undefined + ? { included: a.start } + : { included: Math.min(a.start, b.start) }, end: a.end === undefined || b.end === undefined ? undefined : { excluded: Math.max(a.end, b.end) }, }; } diff --git a/quest/m1/README.md b/quest/m1/README.md index 5b3caebb26..f7e3211a58 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -58,7 +58,6 @@ blocks. The quests that gated m0 lines moved under them. - [Flapping peer drill](/quest/m1/cross-relay-flap-drill.md) - the cluster burst drill survives a flapping peer link, the regression for 0.15.6 route-flap drops - [Untimed lite-07](/quest/m1/lite-untimed.md) - lite-07 carries an untimed track in both languages; lite-05/06 write send time - [Finalize moq-lite-07](/quest/m1/lite07-finalize.md) - when the maintainer cuts it, lite-07 negotiates as `moq-lite-07` and the next release ships it -- [Late lower groups](/quest/m1/lite-late-lower-group.md) - a moq-lite subscriber with a floor receives a group created below the first served one, as moq-transport does - [JSON stream budget](/quest/m1/json-stream-budget.md) - an oversized JSON stream record is refused without ending the log, and a JS subscribe to a gone track answers NotFound - [Flate stream budget](/quest/m1/flate-stream-budget.md) - a flate stream refuses an oversized append without ending, sharing one DEFLATE bound with json - [JS track takeover](/quest/m1/js-track-takeover.md) - JS `createTrack` answers a queued request and continues its sequences, as Rust does diff --git a/quest/m1/lite-late-lower-group.md b/quest/m1/lite-late-lower-group.md deleted file mode 100644 index 451ed6ec01..0000000000 --- a/quest/m1/lite-late-lower-group.md +++ /dev/null @@ -1,63 +0,0 @@ -# [S] moq-lite delivers a late group above the subscriber's floor - -## Goal - -A moq-lite subscriber with an explicit floor receives every group at or above -that floor that its `max_age` still considers fresh, whatever order the -publisher created them in, as moq-transport already does. Today the lite -publisher silently drops a group created below the first group it served: the -publisher's `create_group` and `write_frame` succeed, and the subscriber sees a -clean end without it. - -## Plan - -`TrackRun::start` (`rs/moq-net/src/lite/publisher.rs`) sends `SUBSCRIBE_START` -for the first served group and then calls `raise_start_to(start)`, which -suppresses every lower group regardless of the subscription's floor and -budget. #4387 added it to resolve a relayed subscription's start from its -source; keep that fix (see the note in -[Track tail interop](/quest/m1/track-tail-interop.md)) while honoring an -explicit floor. - -Decided: - -- An explicit floor delivers late lower groups within `max_age`, matching - moq-transport, draft-ietf-moq-transport ("subscriptions without a filter pass - all Objects"), and moq-mux's floor handling (#3258). -- `start: None` joins where the publisher starts: the first served group - becomes the floor, so a group created below it is dropped by design. This is - what moq-mux's `container::Consumer` and the IETF publisher already assume. - Fix the `Subscription::start` docs in moq-net (which say `None` equals a - floor of group 0) and the matching js/net docs. -- `raise_start_to` applies only when the subscription has no explicit floor - (`start.is_none()`). With an explicit floor, `SUBSCRIBE_START` still reports - the first served group, but nothing below it is suppressed. -- Aggregation follows: `Subscription::poll_combined` - (`rs/moq-net/src/model/subscription.rs`) lets any `None` clear another - subscriber's explicit floor, so a relay forwards "join at the start" upstream - and the explicit-floor subscriber still loses group 0. An explicit floor - survives mixing with `None`; the loosest explicit floor wins. Each - subscriber's own cursor still filters what it sees. -- `start_floor_suppresses_late_lower_arrivals` stays as the `None` case: it - subscribes with `None`, receives group 7, raises the start to 7 as - `SUBSCRIBE_START` does, and group 5 is still suppressed. Add its explicit-floor - twin, which delivers group 5. - -If `drafts/draft-lcurley-moq-lite.md` describes `SUBSCRIBE_START` as an -implicit drop below the start regardless of the floor, update it in the same -PR. - -Test: port the reporter's deterministic repro -(`rs/moq-net/tests/late_lower_group.rs` on kidq330's fork, paused time and the -mock transport) for lite-05, lite-06, and lite-07-wip, direct and through a -relay, with the in-process and moq-transport controls. Add a mixed case through -a relay: a `None` subscriber already receiving group 1, then one with a floor -of group 0, and a fresh group 0 reaches only the second. - -## Closes - -- [#4595](https://github.com/moq-dev/moq/issues/4595) - moq-lite: a group created below the first served group is silently dropped - -## Related - -- [SUBSCRIBE_DROP](/quest/m1/subscribe-drop.md) - once publishers name every undelivered group, a group dropped under the None floor can be named instead of vanishing diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index 13e3a6da23..05312c0e73 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -1766,7 +1766,8 @@ where let subscription = target.subscription(); let start = subscription.as_ref().and_then(|s| s.start); // A live join delivers nothing below the group SUBSCRIBE_OK names as Largest. - let live = start.is_none(); + // A floor at group 0 is that same join on pre-draft-20; see [`subscribe_join`]. + let live = joins_live(start, self.version); let join = match subscribe_join( subscription.as_ref().and_then(|s| s.start), subscription.as_ref().and_then(|s| s.end), @@ -6664,6 +6665,19 @@ struct Join { fetch: Option, } +/// Whether `start` joins at the publisher's live edge. +/// +/// An absent floor does. So does a floor at the first frame of group 0 on a pre-draft-20 +/// session: that draft's absolute fetch names only group 0, and the current group's head +/// would otherwise arrive on neither stream. Draft-20 spells group 0 as an unfiltered +/// subscription, which keeps its own start. +fn joins_live(start: Option, version: Version) -> bool { + if Filter::is_draft20(version) { + return start.is_none(); + } + matches!(start, None | Some(track::Position { group: 0, frame: 0 })) +} + /// What a moq-lite subscription's group range asks for on the wire. /// /// moq-lite joins a track at the *start* of the current group, which is a decodable point. @@ -6674,7 +6688,9 @@ struct Join { /// /// Earlier drafts have no fill parameter. Every subscription is Largest Object, followed /// by a joining FETCH: relative at offset 0 for a live join, absolute at the requested -/// group for an explicit group-aligned and unbounded start. A frame-level start or a +/// group for an explicit group-aligned and unbounded start. A floor at the first frame +/// of group 0 excludes nothing, so it uses the live join: an absolute fetch of group 0 +/// would leave the current group's head on neither stream. A frame-level start or a /// bounded end has no joining-FETCH spelling, so those shapes are refused rather than /// rounded down or left open. /// @@ -6695,7 +6711,7 @@ fn subscribe_join( filter: Filter::NextObject, fill: None, fetch: Some(match start { - None => JoiningFetch::Relative { group_offset: 0 }, + None | Some(track::Position { group: 0, frame: 0 }) => JoiningFetch::Relative { group_offset: 0 }, Some(start) => JoiningFetch::Absolute { group_id: start.group }, }), }); @@ -6907,6 +6923,19 @@ mod filter_tests { } } + /// A floor at group 0 excludes no group, so the join is the live one. Fetching only + /// group 0 would drop the head of whichever group is current. + #[test] + fn older_drafts_join_live_at_group_zero() { + for version in JOINING_DRAFTS { + assert_eq!( + subscribe_join(Some(track::Position::group(0)), None, version).unwrap(), + subscribe_join(None, None, version).unwrap(), + "{version}" + ); + } + } + /// An explicit group-aligned unbounded start is an absolute joining FETCH at that group. #[test] fn older_drafts_absolute_join_at_the_start_group() { diff --git a/rs/moq-net/src/lite/publisher.rs b/rs/moq-net/src/lite/publisher.rs index fcf39467f4..50da93f0e2 100644 --- a/rs/moq-net/src/lite/publisher.rs +++ b/rs/moq-net/src/lite/publisher.rs @@ -1251,7 +1251,7 @@ impl Request for SubscribeServe { let subscription = crate::track::Subscription { priority: msg.priority, max_age: serving_max_age(shared.version, msg.max_age), - ..Bounds::from(&msg).positions() + ..Bounds::from(&msg).subscription(shared.version) }; // One subscriber for the whole subscription: the run loop polls its groups and its @@ -1596,10 +1596,9 @@ mod test { assert_eq!(drain(&mut declared), vec![0, 1, 2]); } - /// The run_track contract behind SUBSCRIBE_OK's implicit drop: once the first - /// served group resolves the start, the cursor floor rises to it, so a lower - /// group arriving late (unordered delivery) is never served after the range - /// was declared dropped. + /// An unfloored join adopts the first served group: the cursor rises to it, so a + /// group created below it is not served. An explicit floor does not; see + /// [`explicit_floor_delivers_a_late_lower_group`]. #[moq_net_sim::test] async fn start_floor_suppresses_late_lower_arrivals() { use futures::FutureExt; @@ -1634,6 +1633,39 @@ mod test { ); } + /// The explicit-floor twin of [`start_floor_suppresses_late_lower_arrivals`]. + /// The subscription named group 0, so group 5 is still inside the floor when + /// it is created after group 7. + #[moq_net_sim::test] + async fn explicit_floor_delivers_a_late_lower_group() { + let mut producer = track_producer("test"); + let mut subscriber = producer.subscribe( + track::Subscription::default() + .with_start(track::Position::group(0)) + .with_max_age(std::time::Duration::from_secs(60)), + ); + + let write = |producer: &mut track::Producer, sequence: u64| { + let mut group = producer.create_group(crate::group::Info { sequence }).unwrap(); + group + .write_frame(Timestamp::from_millis(1).unwrap(), b"x".to_vec()) + .unwrap(); + group.finish().unwrap(); + }; + + write(&mut producer, 7); + match recv_next(&mut subscriber, false, false).await.unwrap() { + Recv::Group(group) => assert_eq!(group.sequence, 7), + _ => panic!("expected the first group"), + } + + write(&mut producer, 5); + match recv_next(&mut subscriber, false, false).await.unwrap() { + Recv::Group(group) => assert_eq!(group.sequence, 5), + _ => panic!("expected the late group above the floor"), + } + } + #[moq_net_sim::test] async fn recv_next_drains_datagram_before_finished() { let mut producer = track_producer("test"); @@ -2195,6 +2227,20 @@ impl Bounds { sub } + /// The range to store on the serving track. + /// + /// Lite-06 encodes an omitted floor as 0, and that 0 is group 0. Store it as an + /// explicit floor so it widens a higher resume instead of leaving that resume in + /// place: a fresh reader still receives a finished group the resume skipped. A + /// pre-06 absent start stays absent, because those drafts join at the latest group. + fn subscription(&self, version: Version) -> crate::track::Subscription { + let mut sub = self.positions(); + if version.resolves_start() && sub.start.is_none() { + sub.start = Some(track::Position::group(0)); + } + sub + } + /// The group to apply a frame offset to and the offset itself, or `None` when /// delivery starts on a group boundary. fn start_frame(&self) -> Option<(u64, u64)> { @@ -2397,7 +2443,7 @@ impl TrackRun { let _ = self.track.update(crate::track::Subscription { priority: upd.priority, max_age: serving_max_age(self.ctx.version, upd.max_age), - ..bounds.positions() + ..bounds.subscription(self.ctx.version) }); // An explicit start moves the read cursor. Lite-06+ encodes an absent // start as no floor, so a subscription that had one above group 0 has to @@ -2568,12 +2614,19 @@ impl TrackRun { group: start, largest, }))?; - // SUBSCRIBE_START is an implicit drop of everything below the resolved start (the - // subscriber records it as a permanent miss), so a lower group arriving late must - // not be served after all. A widening SUBSCRIBE_UPDATE re-lowers the floor, - // renegotiating the resolved start along with the demand. Raised, not assigned: an - // update that landed while the group was held may already have raised it past. - self.track.raise_start_to(start); + // SUBSCRIBE_START names the first group served now. Groups already skipped between + // the requested floor and this group stay unavailable. A later group at or above an + // explicit floor, still inside the subscriber's max age, is delivered, so the cursor + // stays at that floor. + // + // A pre-06 subscription that named no group is the exception: those drafts define an + // absent Group Start as the latest group, and the first group served becomes the + // floor. Lite-06 encodes a floor of group 0 as 0, which decodes as no named floor; + // that 0 is group 0, so it is not pinned. Raised, not assigned: an update that + // landed while the group was held may already have raised it past. + if self.track.subscription().start.is_none() && !self.ctx.version.resolves_start() { + self.track.raise_start_to(start); + } self.serve(group); Ok(()) } diff --git a/rs/moq-net/src/lite/subscribe.rs b/rs/moq-net/src/lite/subscribe.rs index b81cc2c4ea..9bed7c5193 100644 --- a/rs/moq-net/src/lite/subscribe.rs +++ b/rs/moq-net/src/lite/subscribe.rs @@ -157,15 +157,15 @@ fn canonical_start_group(version: Version, start_group: Option, start_frame /// Encode the `Group Start` field shared by SUBSCRIBE and SUBSCRIBE_UPDATE. /// -/// The inverse of [`decode_start_group`]: lite-06 writes the raw floor (`None` and -/// `Some(0)` are the same absence of a constraint), while a pre-06 wire gets `Some(0)` -/// folded back to absent. On those wires an explicit group 0 means "replay from the -/// beginning", which is not what a vacuous floor asks for. +/// The inverse of [`decode_start_group`]. Lite-06 writes the raw floor, so `None` and +/// `Some(0)` are the same group 0. A pre-06 wire encodes the sequence + 1: `None` is 0 +/// (the latest group, where the publisher starts) and `Some(0)` is 1, replay from the +/// beginning. fn encode_start_group(w: &mut Encoder<'_>, version: Version, start_group: Option) -> Result<(), EncodeError> { if version.resolves_start() { return w.varint(start_group.unwrap_or(0)); } - w.varint_opt(start_group.filter(|&group| group > 0)) + w.varint_opt(start_group) } /// Decode the trailing `Frame Start` / `Frame End` pair shared by SUBSCRIBE, @@ -773,8 +773,8 @@ mod test { } /// Lite06 carries the raw floor; pre-06 wires encode the sequence + 1 with 0 meaning - /// the latest group. A vacuous floor folds to absent on the old wire, where an - /// explicit group 0 would mean "replay from the beginning" instead. + /// the latest group. An explicit group 0 is that sequence, so it round-trips as group 0 + /// rather than as the latest group. #[test] fn group_start_is_absolute_on_lite06() { let mut msg = subscribe_sample(); @@ -809,11 +809,23 @@ mod test { let got = crate::coding::decode_buf(&mut zero.as_slice(), Version::Lite06, Subscribe::decode_msg).unwrap(); assert_eq!(got.start_group, None); - // On the pre-06 wire the vacuous floor folds to absent (the latest group). - let mut folded = Vec::new(); - msg.encode_msg(&mut Encoder::new(&mut folded, Version::Lite05.into()), Version::Lite05) + // On the pre-06 wire an explicit group 0 is sequence + 1, distinct from an + // absent start (the latest group). + let mut explicit = Vec::new(); + msg.encode_msg( + &mut Encoder::new(&mut explicit, Version::Lite05.into()), + Version::Lite05, + ) + .unwrap(); + let got = crate::coding::decode_buf(&mut explicit.as_slice(), Version::Lite05, Subscribe::decode_msg).unwrap(); + assert_eq!(got.start_group, Some(0)); + + msg.start_group = None; + let mut latest = Vec::new(); + msg.encode_msg(&mut Encoder::new(&mut latest, Version::Lite05.into()), Version::Lite05) .unwrap(); - let got = crate::coding::decode_buf(&mut folded.as_slice(), Version::Lite05, Subscribe::decode_msg).unwrap(); + assert_ne!(explicit, latest); + let got = crate::coding::decode_buf(&mut latest.as_slice(), Version::Lite05, Subscribe::decode_msg).unwrap(); assert_eq!(got.start_group, None); } diff --git a/rs/moq-net/src/model/subscription.rs b/rs/moq-net/src/model/subscription.rs index cec094b20e..e3f9f0fc49 100644 --- a/rs/moq-net/src/model/subscription.rs +++ b/rs/moq-net/src/model/subscription.rs @@ -56,18 +56,22 @@ pub struct Subscription { /// in three reads as three. The publisher's copy is stamped as it produces, so the /// gate there still holds; it is just the coarser of the two. pub max_age: Duration, - /// The lowest [`Position`] the publisher may deliver, or `None` for no floor. + /// The lowest [`Position`] the publisher may deliver, or `None` to join where the + /// publisher starts. /// /// A floor, not a request: only [`Self::max_age`] asks for data, and the floor bounds - /// how far back it may reach. `None` and a floor of group 0 mean the same thing, since - /// nothing sits below group 0. Delivery starts at the oldest group at or above the - /// floor that the budget still considers fresh, so a floor above the live edge simply - /// waits there (a resumed subscription naming where it left off). + /// how far back it may reach. `None` is not a floor of group 0. The first group served + /// to an unfloored subscriber becomes its floor, so a group created below that group + /// is not delivered. An explicit floor, including group 0, still delivers a later + /// group at or above it while the budget considers it fresh. A floor above the live + /// edge simply waits there (a resumed subscription naming where it left off). /// - /// Aggregated across every live subscriber (the loosest floor wins, and any subscriber - /// without one clears it), so it says what the publisher sends, not what any one - /// subscriber sees. [`crate::track::Subscriber::set_groups`] is the local read cursor; - /// setting one does not imply the other. See [Local cursor vs wire + /// Aggregated across every live subscriber: the loosest explicit floor wins, and a + /// subscriber without one leaves that floor in place. All of them unfloored stays + /// unfloored. It says what the publisher sends, not what any one subscriber sees. + /// Each subscriber's read cursor still filters its own view. + /// [`crate::track::Subscriber::set_groups`] is that cursor; setting one does not imply + /// the other. See [Local cursor vs wire /// preference](crate::track::Subscriber#local-cursor-vs-wire-preference). pub start: Option, /// First [`Position`] the publisher should *not* deliver, or `None` for no end. @@ -161,8 +165,9 @@ impl Subscription { max_age: self.max_age.max(combined.max_age), // Bounds fold as whole positions. Two subscribers starting in the same group // are separated only by their frame, so folding group and frame independently - // would invent a bound neither asked for. - start: min_floored(self.start, combined.start), + // would invent a bound neither asked for. An omitted floor does not clear an + // explicit one: the loosest explicit floor wins. + start: min_some(self.start, combined.start), end: max_unbounded(self.end, combined.end), }; @@ -299,12 +304,12 @@ pub(super) fn before_end(sequence: u64, end: Option) -> bool { } // Combining two optional bounds comes in two families, and they disagree on what `None` -// means. `_some` treats it as the neutral element (the other side wins), for intersecting -// two ranges that each restrict independently. `_floored` / `_unbounded` treat it as -// absorbing (the result is `None` too), for aggregating across subscribers, where one -// subscriber asking for everything makes the aggregate everything. Picking the wrong -// family silently narrows or widens what the publisher sends, so the suffix, not the -// `min`/`max`, is the part to read. +// means. `_some` treats it as the neutral element (the other side wins). `_unbounded` +// treats it as absorbing (the result is `None` too). A start uses the neutral family: +// omitting one joins where the publisher starts, and the loosest explicit floor still +// wins. An end uses the absorbing family: one subscriber with no end keeps the aggregate +// open. Picking the wrong family silently narrows or widens what the publisher sends, so +// the suffix, not the `min`/`max`, is the part to read. /// The higher of two optional bounds, `None` neutral. pub(super) fn max_some(a: Option, b: Option) -> Option { @@ -315,13 +320,12 @@ pub(super) fn max_some(a: Option, b: Option) -> Option { } } -/// The lower of two optional floors, `None` absorbing (no floor). The mirror of -/// [`max_unbounded`]: both bounds only ever *restrict*, so a subscriber without one keeps -/// the aggregate unrestricted. -pub(super) fn min_floored(a: Option, b: Option) -> Option { +/// The lower of two optional floors, `None` neutral. The mirror of [`max_some`]. +pub(super) fn min_some(a: Option, b: Option) -> Option { match (a, b) { (Some(a), Some(b)) => Some(a.min(b)), - (None, _) | (_, None) => None, + (Some(a), None) | (None, Some(a)) => Some(a), + (None, None) => None, } } @@ -413,10 +417,13 @@ mod tests { let combined = combine(&[catchup.clone(), older_catchup]).unwrap(); assert_eq!(combined.start, Some(Position::group(5))); - // A subscriber with no floor at all clears the aggregate: its budget may reach - // below any floor the others set. + // An omitted floor joins where the publisher starts. It leaves an explicit floor + // in place; each subscriber's own cursor still filters what it reads. let unfloored = Subscription::default(); - let combined = combine(&[catchup, unfloored]).unwrap(); + let combined = combine(&[catchup, unfloored.clone()]).unwrap(); + assert_eq!(combined.start, Some(Position::group(10))); + + let combined = combine(&[unfloored.clone(), unfloored]).unwrap(); assert_eq!(combined.start, None); } diff --git a/rs/moq-net/tests/history_groups.rs b/rs/moq-net/tests/history_groups.rs index 62db30709f..adaa03161a 100644 --- a/rs/moq-net/tests/history_groups.rs +++ b/rs/moq-net/tests/history_groups.rs @@ -113,9 +113,8 @@ async fn round(version: &str, relay: bool, newest_first: bool) -> Vec<(u64, usiz .await .expect("no subscriber appeared") .unwrap(); - // Pre-06 wires can't name group 0 (it reads as the latest group), so the publisher - // must place its cursor before the groups exist. Simulated time advances only once - // every task is idle, so this settles the subscription first. + // Simulated time advances only once every task is idle, so this settles the + // subscription before the groups exist. moq_net_sim::sleep(Duration::from_millis(10)).await; // The hop the publisher's group streams cross first: into the relay, when there is one. diff --git a/rs/moq-net/tests/late_lower_group.rs b/rs/moq-net/tests/late_lower_group.rs new file mode 100644 index 0000000000..53a54234e8 --- /dev/null +++ b/rs/moq-net/tests/late_lower_group.rs @@ -0,0 +1,307 @@ +//! A subscriber with an explicit floor receives a group created below the first group +//! that was served, while it is still inside `max_age`. +//! +//! An unfloored pre-06 subscribe joins where the publisher starts, so that same group +//! is dropped. Lite-06 and later encode a floor of group 0 as 0, which is the only +//! spelling of both an omitted floor and group 0; the draft says that 0 is group 0. + +mod support; + +use std::time::Duration; + +use moq_net::track::{Info, Position, Subscription}; +use moq_net::{Timestamp, Version}; +use support::harness::{MockConnectOptions, connect_mock}; + +const TIMEOUT: Duration = Duration::from_secs(10); + +/// Long enough that two groups written one after the other are both fresh. +const BUDGET: Duration = Duration::from_secs(60); + +/// A budget no group outlives, on the publisher. +const FOREVER: Duration = Duration::from_millis((1 << 53) - 1); + +const LITE: &[&str] = &["moq-lite-05", "moq-lite-06", "moq-lite-07-wip"]; + +const TRANSPORT: &[&str] = &["moq-transport-14", "moq-transport-17", "moq-transport-22"]; + +fn produce_origin(hop: u64) -> moq_net::origin::Producer { + let (producer, driver) = + moq_net::origin::Producer::new(moq_net::origin::Config::new(moq_net::Hop::new(hop).unwrap())); + support::harness::spawn(driver); + producer +} + +fn write_group(track: &moq_net::track::Producer, sequence: u64) { + let mut group = track.create_group(moq_net::group::Info { sequence }).unwrap(); + for _ in 0..4 { + group.write_frame(Timestamp::now(), &b"x"[..]).unwrap(); + } + group.finish().unwrap(); +} + +async fn read_all(sub: &mut moq_net::track::Subscriber) -> Vec<(u64, usize)> { + let mut got = Vec::new(); + loop { + let next = moq_net_sim::timeout(TIMEOUT, sub.recv_group()) + .await + .expect("recv_group timed out"); + let Some(mut group) = next.expect("recv_group") else { + break; + }; + let mut frames = 0; + loop { + let frame = moq_net_sim::timeout(TIMEOUT, group.read_frame()) + .await + .expect("read_frame timed out") + .expect("read_frame"); + match frame { + Some(_) => frames += 1, + None => break, + } + } + got.push((group.sequence, frames)); + } + got +} + +/// Group 1, then a fresh group 0, over `version`. `relay` inserts one hop. +async fn session(version: &str, relay: bool, start: Option) -> Vec<(u64, usize)> { + let version: Version = version.parse().unwrap(); + let publisher = produce_origin(1); + let broadcast = publisher.create_broadcast("bcast").unwrap(); + let track = broadcast + .create_track("late", Info::default().with_max_age(FOREVER)) + .unwrap(); + broadcast.announce(Default::default()).unwrap(); + + let relay_origin = produce_origin(2); + let upstream = match relay { + true => { + let mut options = MockConnectOptions::new(version); + options.server_publish = Some(publisher.consume()); + options.client_subscribe = Some(relay_origin.clone()); + Some(connect_mock(options).await) + } + false => None, + }; + + let subscriber = produce_origin(3); + let mut options = MockConnectOptions::new(version); + options.server_publish = Some(match relay { + true => relay_origin.consume(), + false => publisher.consume(), + }); + options.client_subscribe = Some(subscriber.clone()); + let pair = connect_mock(options).await; + + let consumer = subscriber.consume(); + moq_net_sim::timeout(TIMEOUT, consumer.routed("bcast")) + .await + .expect("announce timeout") + .expect("routed"); + let remote = moq_net_sim::timeout(TIMEOUT, consumer.request_broadcast("bcast")) + .await + .expect("resolve timeout") + .expect("broadcast resolves"); + + let subscription = Subscription::default().with_max_age(BUDGET).with_start(start); + let reader = moq_net_sim::spawn(async move { + let mut sub = remote + .track("late") + .unwrap() + .subscribe(subscription) + .await + .expect("subscribe"); + read_all(&mut sub).await + }); + + moq_net_sim::timeout(TIMEOUT, track.demand().used()) + .await + .expect("no subscriber appeared") + .unwrap(); + // Idle time lets the subscription settle, and again lets SUBSCRIBE_START land, + // before the lower group is created. + moq_net_sim::sleep(Duration::from_millis(10)).await; + write_group(&track, 1); + moq_net_sim::sleep(Duration::from_millis(10)).await; + write_group(&track, 0); + track.finish().unwrap(); + + let got = moq_net_sim::timeout(TIMEOUT * 3, reader) + .await + .expect("reader hung") + .expect("reader panicked"); + drop((track, pair, upstream, relay_origin, broadcast, publisher, subscriber)); + got +} + +async fn in_process(start: Option) -> Vec<(u64, usize)> { + let publisher = produce_origin(1); + let broadcast = publisher.create_broadcast("bcast").unwrap(); + let track = broadcast + .create_track("late", Info::default().with_max_age(FOREVER)) + .unwrap(); + let subscription = Subscription::default().with_max_age(BUDGET).with_start(start); + let mut sub = track.subscribe(subscription); + let reader = moq_net_sim::spawn(async move { read_all(&mut sub).await }); + moq_net_sim::sleep(Duration::from_millis(10)).await; + write_group(&track, 1); + moq_net_sim::sleep(Duration::from_millis(10)).await; + write_group(&track, 0); + track.finish().unwrap(); + let got = moq_net_sim::timeout(TIMEOUT * 3, reader) + .await + .expect("reader hung") + .expect("reader panicked"); + drop((track, broadcast, publisher)); + got +} + +#[moq_net_sim::test] +async fn an_explicit_floor_receives_a_group_created_below_the_first_served_group() { + let want = vec![(1, 4), (0, 4)]; + let mut failures = Vec::new(); + + let got = in_process(Some(Position::group(0))).await; + if got != want { + failures.push(format!("in-process: {got:?}")); + } + + for version in LITE.iter().chain(TRANSPORT) { + for relay in [false, true] { + let got = session(version, relay, Some(Position::group(0))).await; + if got != want { + failures.push(format!("{version} relay={relay}: {got:?}")); + } + } + } + assert!(failures.is_empty(), "{failures:#?}"); +} + +/// Pre-06, an omitted floor joins at the first served group. Group 0, created after +/// group 1 was served, is below that join. Lite-06 cannot spell this separately from +/// a floor of group 0; that case is the explicit-floor test. +#[moq_net_sim::test] +async fn an_unfloored_lite05_join_drops_a_group_created_below_the_first_served_group() { + let want = vec![(1, 4)]; + let mut failures = Vec::new(); + for relay in [false, true] { + let got = session("moq-lite-05", relay, None).await; + if got != want { + failures.push(format!("relay={relay}: {got:?}")); + } + } + assert!(failures.is_empty(), "{failures:#?}"); +} + +/// A `None` subscriber is already reading group 1. A second subscriber names group 0. +/// The fresh group 0 reaches only the second. Lite-05 is the draft that can tell the +/// two subscriptions apart on the wire. +#[moq_net_sim::test] +async fn a_relayed_explicit_floor_survives_an_unfloored_subscriber() { + let version: Version = "moq-lite-05".parse().unwrap(); + let publisher = produce_origin(1); + let broadcast = publisher.create_broadcast("bcast").unwrap(); + let track = broadcast + .create_track("late", Info::default().with_max_age(FOREVER)) + .unwrap(); + broadcast.announce(Default::default()).unwrap(); + + let relay = produce_origin(2); + let mut options = MockConnectOptions::new(version); + options.server_publish = Some(publisher.consume()); + options.client_subscribe = Some(relay.clone()); + let upstream = connect_mock(options).await; + + let unfloored_origin = produce_origin(3); + let mut options = MockConnectOptions::new(version); + options.server_publish = Some(relay.consume()); + options.client_subscribe = Some(unfloored_origin.clone()); + let unfloored_pair = connect_mock(options).await; + let consumer = unfloored_origin.consume(); + moq_net_sim::timeout(TIMEOUT, consumer.routed("bcast")) + .await + .expect("announce timeout") + .unwrap(); + let remote = consumer.request_broadcast("bcast").await.unwrap(); + let unfloored_sub = Subscription::default().with_max_age(BUDGET); + let reader = moq_net_sim::spawn(async move { + let mut sub = remote.track("late").unwrap().subscribe(unfloored_sub).await.unwrap(); + let mut group = moq_net_sim::timeout(TIMEOUT, sub.recv_group()) + .await + .expect("group 1") + .unwrap() + .expect("group 1 ended"); + assert_eq!(group.sequence, 1); + for _ in 0..4 { + assert!(group.read_frame().await.unwrap().is_some()); + } + assert!(group.read_frame().await.unwrap().is_none()); + sub + }); + + track.demand().used().await.unwrap(); + moq_net_sim::sleep(Duration::from_millis(10)).await; + write_group(&track, 1); + let mut unfloored = moq_net_sim::timeout(TIMEOUT, reader) + .await + .expect("unfloored reader") + .unwrap(); + + let floored_origin = produce_origin(4); + let mut options = MockConnectOptions::new(version); + options.server_publish = Some(relay.consume()); + options.client_subscribe = Some(floored_origin.clone()); + let floored_pair = connect_mock(options).await; + let consumer = floored_origin.consume(); + consumer.routed("bcast").await.unwrap(); + let remote = consumer.request_broadcast("bcast").await.unwrap(); + let floored_sub = Subscription::default() + .with_max_age(BUDGET) + .with_start(Position::group(0)); + let floored_reader = moq_net_sim::spawn(async move { + let mut sub = remote.track("late").unwrap().subscribe(floored_sub).await.unwrap(); + let mut saw_zero = false; + for _ in 0..2 { + let Some(mut group) = moq_net_sim::timeout(TIMEOUT, sub.recv_group()) + .await + .expect("floored recv") + .unwrap() + else { + break; + }; + let sequence = group.sequence; + while group.read_frame().await.unwrap().is_some() {} + if sequence == 0 { + saw_zero = true; + break; + } + } + assert!(saw_zero, "explicit floor missed group 0"); + }); + + // The spawned subscribe runs to its first wait, which forwards the new floor. + moq_net_sim::sleep(Duration::from_millis(10)).await; + write_group(&track, 0); + moq_net_sim::timeout(TIMEOUT, floored_reader) + .await + .expect("floored reader") + .unwrap(); + + let late = moq_net_sim::timeout(Duration::from_millis(50), unfloored.recv_group()).await; + assert!(late.is_err(), "unfloored subscriber received a group below its join"); + + drop(( + unfloored, + unfloored_pair, + floored_pair, + upstream, + track, + broadcast, + publisher, + relay, + unfloored_origin, + floored_origin, + )); +} From 5a1daa0a8ffdfd59b073b7e219b3decfb8b9a8ad Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Fri, 9 Oct 2026 10:06:56 -0700 Subject: [PATCH 2/4] fix(net): keep an unfloored subscriber absorbing the aggregate floor A resume floor above the live edge must not hide the latest group from a reader that names no floor. Restore the absorbing floor fold in Rust and JS, which also makes storing lite-06's 0 as an explicit floor and the moq-transport group-0 live-join special case unnecessary. The mixed relay case that needs both an explicit floor and a live join moves to the lite-07 Live flag quest. Add the reverse-order regression: a peer resumes past a quiet catalog's only group, then an unfloored reader at the relay still receives it. Co-Authored-By: Claude Opus 5.5 --- js/net/src/lite/publisher.ts | 5 +- js/net/src/lite/subscriber.ts | 9 ++- js/net/src/track.test.ts | 13 ++- quest/m1/lite-live.md | 5 +- rs/moq-net/src/ietf/subscriber.rs | 35 +-------- rs/moq-net/src/lite/publisher.rs | 18 +---- rs/moq-net/tests/late_lower_group.rs | 113 +-------------------------- rs/moq-net/tests/quiet_catalog.rs | 86 +++++++++++--------- 8 files changed, 68 insertions(+), 216 deletions(-) diff --git a/js/net/src/lite/publisher.ts b/js/net/src/lite/publisher.ts index 5f9318a091..f08cacf6e5 100644 --- a/js/net/src/lite/publisher.ts +++ b/js/net/src/lite/publisher.ts @@ -619,14 +619,11 @@ export class Publisher { } const endGroup = exclusiveGroupEnd(msg.endGroup); - // Lite-06 encodes an omitted floor as 0, and that 0 is group 0. Store it as an - // explicit floor so it widens a higher resume. A pre-06 absent start stays absent. - const startGroup = msg.startGroup === undefined && resolvesStart(this.version) ? 0 : msg.startGroup; const track = wireOf(front).subscribe(msg.track, { priority: msg.priority, maxDelay: Milli(servingMaxDelay(this.version, msg.maxDelay)), groups: { - start: startGroup === undefined ? undefined : { included: startGroup }, + start: msg.startGroup === undefined ? undefined : { included: msg.startGroup }, end: endGroup === undefined ? undefined : { excluded: endGroup }, }, }); diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index 12685c8b91..571c8666c9 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -49,6 +49,7 @@ import { hasDatagrams, hasProbeRtt, hasStreamCount, + resolvesStart, restartSupported, Version, waitsForSubscriberFin, @@ -1127,12 +1128,14 @@ export class Subscriber { #sameSubscription(a: track.Subscription, b: track.Subscription): boolean { const ag = groupBounds(a.groups); const bg = groupBounds(b.groups); - // `groupBounds` reads an omitted start as 0. On a pre-06 wire those are different - // subscriptions: omitted joins at the publisher's start, and 0 is group 0. + // `groupBounds` reads an omitted start as 0. A pre-06 wire tells them apart: omitted + // joins at the publisher's start, and 0 is group 0. Lite-06 encodes both as 0. + const sameStart = + resolvesStart(this.version) || (a.groups?.start === undefined) === (b.groups?.start === undefined); return ( (a.priority ?? 0) === (b.priority ?? 0) && (a.maxDelay ?? 0) === (b.maxDelay ?? 0) && - (a.groups?.start === undefined) === (b.groups?.start === undefined) && + sameStart && ag.start === bg.start && ag.end === bg.end ); diff --git a/js/net/src/track.test.ts b/js/net/src/track.test.ts index 946b54526a..54f9501a3e 100644 --- a/js/net/src/track.test.ts +++ b/js/net/src/track.test.ts @@ -395,16 +395,13 @@ test("multiple subscriber options aggregate like Rust", async () => { expect(await none).toBeUndefined(); }); -test("an omitted floor leaves another subscriber's floor in place", () => { +// A resume floor above the live edge must not hide the latest group from a reader that +// joins after it: on a quiet track that group may be the only one for a long time. +test("an unfloored subscriber clears a resume floor above the live edge", () => { const producer = new TrackProducer("test"); - producer.subscribe({ groups: { start: { included: 10 } } }); + producer.subscribe({ groups: { start: { included: 4 } } }); producer.subscribe({}); - expect(producer.subscription.peek()?.groups?.start).toEqual({ included: 10 }); - - const unfloored = new TrackProducer("test"); - unfloored.subscribe({}); - unfloored.subscribe({}); - expect(unfloored.subscription.peek()?.groups?.start).toBeUndefined(); + expect(producer.subscription.peek()?.groups?.start).toBeUndefined(); }); test("the producer aggregate is clamped without changing subscriber options", async () => { diff --git a/quest/m1/lite-live.md b/quest/m1/lite-live.md index 710fc3e773..deadbc9ebe 100644 --- a/quest/m1/lite-live.md +++ b/quest/m1/lite-live.md @@ -83,10 +83,11 @@ merged `{floor: 2, live}` through a lite-03 to 05 upstream with groups 2-4 retained that delivers all three, a buffered `live` merged with a floor-5 subscriber over a moq-transport upstream with fresh groups 2-4 that still delivers 2-4, an untimed `live` merged with a floor, and the codec mapping -of each older wire. +of each older wire. Also the mixed case #5000 left out: through a lite-05 +relay, a `live` subscriber already reading group 1, then one with a floor of +group 0, and a fresh group 0 reaches only the second. ## Related - [One max_age meaning](/quest/m1/cache-max-age.md) - decides where an untimed `Live` starts -- [Late lower group](/quest/m1/lite-late-lower-group.md) - requires this; #5000 rebases onto it and drops its own floor-combining change - [SUBSCRIBE_DROP](/quest/m1/subscribe-drop.md) - names undelivered groups once publishers report them diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index 481fce112c..d57fc40419 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -1844,8 +1844,7 @@ where let subscription = target.subscription(); let start = subscription.as_ref().and_then(|s| s.start); // A live join delivers nothing below the group SUBSCRIBE_OK names as Largest. - // A floor at group 0 is that same join on pre-draft-20; see [`subscribe_join`]. - let live = joins_live(start, self.version); + let live = start.is_none(); let join = match subscribe_join( subscription.as_ref().and_then(|s| s.start), subscription.as_ref().and_then(|s| s.end), @@ -7218,19 +7217,6 @@ struct Join { fetch: Option, } -/// Whether `start` joins at the publisher's live edge. -/// -/// An absent floor does. So does a floor at the first frame of group 0 on a pre-draft-20 -/// session: that draft's absolute fetch names only group 0, and the current group's head -/// would otherwise arrive on neither stream. Draft-20 spells group 0 as an unfiltered -/// subscription, which keeps its own start. -fn joins_live(start: Option, version: Version) -> bool { - if Filter::is_draft20(version) { - return start.is_none(); - } - matches!(start, None | Some(track::Position { group: 0, frame: 0 })) -} - /// What a moq-lite subscription's group range asks for on the wire. /// /// moq-lite joins a track at the *start* of the current group, which is a decodable point. @@ -7241,9 +7227,7 @@ fn joins_live(start: Option, version: Version) -> bool { /// /// Earlier drafts have no fill parameter. Every subscription is Largest Object, followed /// by a joining FETCH: relative at offset 0 for a live join, absolute at the requested -/// group for an explicit group-aligned and unbounded start. A floor at the first frame -/// of group 0 excludes nothing, so it uses the live join: an absolute fetch of group 0 -/// would leave the current group's head on neither stream. A frame-level start or a +/// group for an explicit group-aligned and unbounded start. A frame-level start or a /// bounded end has no joining-FETCH spelling, so those shapes are refused rather than /// rounded down or left open. /// @@ -7264,7 +7248,7 @@ fn subscribe_join( filter: Filter::NextObject, fill: None, fetch: Some(match start { - None | Some(track::Position { group: 0, frame: 0 }) => JoiningFetch::Relative { group_offset: 0 }, + None => JoiningFetch::Relative { group_offset: 0 }, Some(start) => JoiningFetch::Absolute { group_id: start.group }, }), }); @@ -7476,19 +7460,6 @@ mod filter_tests { } } - /// A floor at group 0 excludes no group, so the join is the live one. Fetching only - /// group 0 would drop the head of whichever group is current. - #[test] - fn older_drafts_join_live_at_group_zero() { - for version in JOINING_DRAFTS { - assert_eq!( - subscribe_join(Some(track::Position::group(0)), None, version).unwrap(), - subscribe_join(None, None, version).unwrap(), - "{version}" - ); - } - } - /// An explicit group-aligned unbounded start is an absolute joining FETCH at that group. #[test] fn older_drafts_absolute_join_at_the_start_group() { diff --git a/rs/moq-net/src/lite/publisher.rs b/rs/moq-net/src/lite/publisher.rs index 62fd317451..92cfcd71e8 100644 --- a/rs/moq-net/src/lite/publisher.rs +++ b/rs/moq-net/src/lite/publisher.rs @@ -1432,7 +1432,7 @@ impl Request for SubscribeServe { let subscription = crate::track::Subscription { priority: msg.priority, max_delay: serving_max_delay(shared.version, msg.max_delay), - ..Bounds::from(&msg).subscription(shared.version) + ..Bounds::from(&msg).positions() }; // One subscriber for the whole subscription: the run loop polls its groups and its @@ -2439,20 +2439,6 @@ impl Bounds { sub } - /// The range to store on the serving track. - /// - /// Lite-06 encodes an omitted floor as 0, and that 0 is group 0. Store it as an - /// explicit floor so it widens a higher resume instead of leaving that resume in - /// place: a fresh reader still receives a finished group the resume skipped. A - /// pre-06 absent start stays absent, because those drafts join at the latest group. - fn subscription(&self, version: Version) -> crate::track::Subscription { - let mut sub = self.positions(); - if version.resolves_start() && sub.start.is_none() { - sub.start = Some(track::Position::group(0)); - } - sub - } - /// The group to apply a frame offset to and the offset itself, or `None` when /// delivery starts on a group boundary. fn start_frame(&self) -> Option<(u64, u64)> { @@ -2681,7 +2667,7 @@ impl TrackRun { let _ = self.track.update(crate::track::Subscription { priority: upd.priority, max_delay: serving_max_delay(self.ctx.version, upd.max_delay), - ..bounds.subscription(self.ctx.version) + ..bounds.positions() }); // An explicit start moves the read cursor. Lite-06+ encodes an absent // start as no floor, so a subscription that had one above group 0 has to diff --git a/rs/moq-net/tests/late_lower_group.rs b/rs/moq-net/tests/late_lower_group.rs index 7b84dd1d55..395eabb0c7 100644 --- a/rs/moq-net/tests/late_lower_group.rs +++ b/rs/moq-net/tests/late_lower_group.rs @@ -100,7 +100,7 @@ async fn session(version: &str, relay: bool, start: Option) -> Vec<(u6 .await .expect("announce timeout") .expect("routed"); - let remote = moq_net_sim::timeout(TIMEOUT, consumer.request_broadcast("bcast")) + let remote = moq_net_sim::timeout(TIMEOUT, consumer.request_broadcast("bcast", None)) .await .expect("resolve timeout") .expect("broadcast resolves"); @@ -194,114 +194,3 @@ async fn an_unfloored_lite05_join_drops_a_group_created_below_the_first_served_g } assert!(failures.is_empty(), "{failures:#?}"); } - -/// A `None` subscriber is already reading group 1. A second subscriber names group 0. -/// The fresh group 0 reaches only the second. Lite-05 is the draft that can tell the -/// two subscriptions apart on the wire. -#[moq_net_sim::test] -async fn a_relayed_explicit_floor_survives_an_unfloored_subscriber() { - let version: Version = "moq-lite-05".parse().unwrap(); - let publisher = produce_origin(1); - let broadcast = publisher.create_broadcast("bcast").unwrap(); - let track = broadcast - .create_track("late", Info::default().with_max_age(FOREVER)) - .unwrap(); - broadcast.announce(Default::default()).unwrap(); - - let relay = produce_origin(2); - let mut options = MockConnectOptions::new(version); - options.server_publish = Some(publisher.consume()); - options.client_subscribe = Some(relay.clone()); - let upstream = connect_mock(options).await; - - let unfloored_origin = produce_origin(3); - let mut options = MockConnectOptions::new(version); - options.server_publish = Some(relay.consume()); - options.client_subscribe = Some(unfloored_origin.clone()); - let unfloored_pair = connect_mock(options).await; - let consumer = unfloored_origin.consume(); - moq_net_sim::timeout(TIMEOUT, consumer.routed("bcast")) - .await - .expect("announce timeout") - .unwrap(); - let remote = consumer.request_broadcast("bcast").await.unwrap(); - let unfloored_sub = Subscription::default().with_max_delay(BUDGET); - let reader = moq_net_sim::spawn(async move { - let mut sub = remote.track("late").unwrap().subscribe(unfloored_sub).await.unwrap(); - let mut group = moq_net_sim::timeout(TIMEOUT, sub.recv_group()) - .await - .expect("group 1") - .unwrap() - .expect("group 1 ended"); - assert_eq!(group.sequence, 1); - for _ in 0..4 { - assert!(group.read_frame().await.unwrap().is_some()); - } - assert!(group.read_frame().await.unwrap().is_none()); - sub - }); - - track.demand().used().await.unwrap(); - moq_net_sim::sleep(Duration::from_millis(10)).await; - write_group(&track, 1); - let mut unfloored = moq_net_sim::timeout(TIMEOUT, reader) - .await - .expect("unfloored reader") - .unwrap(); - - let floored_origin = produce_origin(4); - let mut options = MockConnectOptions::new(version); - options.server_publish = Some(relay.consume()); - options.client_subscribe = Some(floored_origin.clone()); - let floored_pair = connect_mock(options).await; - let consumer = floored_origin.consume(); - consumer.routed("bcast").await.unwrap(); - let remote = consumer.request_broadcast("bcast").await.unwrap(); - let floored_sub = Subscription::default() - .with_max_delay(BUDGET) - .with_start(Position::group(0)); - let floored_reader = moq_net_sim::spawn(async move { - let mut sub = remote.track("late").unwrap().subscribe(floored_sub).await.unwrap(); - let mut saw_zero = false; - for _ in 0..2 { - let Some(mut group) = moq_net_sim::timeout(TIMEOUT, sub.recv_group()) - .await - .expect("floored recv") - .unwrap() - else { - break; - }; - let sequence = group.sequence; - while group.read_frame().await.unwrap().is_some() {} - if sequence == 0 { - saw_zero = true; - break; - } - } - assert!(saw_zero, "explicit floor missed group 0"); - }); - - // The spawned subscribe runs to its first wait, which forwards the new floor. - moq_net_sim::sleep(Duration::from_millis(10)).await; - write_group(&track, 0); - moq_net_sim::timeout(TIMEOUT, floored_reader) - .await - .expect("floored reader") - .unwrap(); - - let late = moq_net_sim::timeout(Duration::from_millis(50), unfloored.recv_group()).await; - assert!(late.is_err(), "unfloored subscriber received a group below its join"); - - drop(( - unfloored, - unfloored_pair, - floored_pair, - upstream, - track, - broadcast, - publisher, - relay, - unfloored_origin, - floored_origin, - )); -} diff --git a/rs/moq-net/tests/quiet_catalog.rs b/rs/moq-net/tests/quiet_catalog.rs index 17eb3d8a79..94b8c31246 100644 --- a/rs/moq-net/tests/quiet_catalog.rs +++ b/rs/moq-net/tests/quiet_catalog.rs @@ -139,46 +139,54 @@ async fn a_quiet_catalog_reaches_a_late_reader_after_its_route_dies() { /// group. A fresh reader must still receive that group. #[moq_net_sim::test] async fn a_quiet_catalog_reaches_a_fresh_reader_when_a_peer_resumes_past_it() { - each_version(LITE, |version| async move { - let publisher = produce_origin(1); - let relay = produce_origin(2); - let resuming = produce_origin(3); - let fresh = produce_origin(4); - let (_broadcast, mut track) = publish(&publisher); - - let _upstream = link(version, &publisher, &relay).await; - let _resume_link = link(version, &relay, &resuming).await; - let _fresh_link = link(version, &relay, &fresh).await; - - let resume_remote = request(&resuming).await; - // Held until the scenario returns. Dropping it would cancel the upstream - // subscription and let the fresh reader open a new one at the live edge. - let _resume = moq_net_sim::spawn(async move { - let _sub = resume_remote - .track("catalog.json") - .unwrap() - .subscribe(track::Subscription::default().with_start(track::Position::group(1))) - .await - .expect("resume subscribe"); - moq_net_sim::sleep(Duration::from_secs(60)).await; - }); - - // Wait for the resume to reach the publisher, so the fresh reader widens it - // rather than opening its own floorless subscription. - let resumed = Some(track::Position::group(1)); - moq_net_sim::timeout(TIMEOUT, async { - while track.subscription().map(|sub| sub.start) != Some(resumed) { - track.subscription_changed().await.unwrap(); - } - }) - .await - .map_err(|_| "the resume never reached the publisher")?; + each_version(LITE, |version| resume_then_fresh(version, false)).await; +} - let remote = request(&fresh).await; - read_snapshot(&remote) +/// The same, with the fresh reader in-process at the relay. It names no floor at all, +/// so the aggregate must not keep the resume's floor above the only group. +#[moq_net_sim::test] +async fn a_quiet_catalog_reaches_an_unfloored_relay_reader_when_a_peer_resumes_past_it() { + each_version(LITE, |version| resume_then_fresh(version, true)).await; +} + +async fn resume_then_fresh(version: Version, in_process: bool) -> Result<(), String> { + let publisher = produce_origin(1); + let relay = produce_origin(2); + let resuming = produce_origin(3); + let fresh = produce_origin(4); + let (_broadcast, mut track) = publish(&publisher); + + let _upstream = link(version, &publisher, &relay).await; + let _resume_link = link(version, &relay, &resuming).await; + let _fresh_link = link(version, &relay, &fresh).await; + + let resume_remote = request(&resuming).await; + // Held until the scenario returns. Dropping it would cancel the upstream + // subscription and let the fresh reader open a new one at the live edge. + let _resume = moq_net_sim::spawn(async move { + let _sub = resume_remote + .track("catalog.json") + .unwrap() + .subscribe(track::Subscription::default().with_start(track::Position::group(1))) .await - .map_err(|err| format!("fresh reader: {err}"))?; - Ok(()) + .expect("resume subscribe"); + moq_net_sim::sleep(Duration::from_secs(60)).await; + }); + + // Wait for the resume to reach the publisher, so the fresh reader widens it + // rather than opening its own floorless subscription. + let resumed = Some(track::Position::group(1)); + moq_net_sim::timeout(TIMEOUT, async { + while track.subscription().map(|sub| sub.start) != Some(resumed) { + track.subscription_changed().await.unwrap(); + } }) - .await; + .await + .map_err(|_| "the resume never reached the publisher")?; + + let remote = request(if in_process { &relay } else { &fresh }).await; + read_snapshot(&remote) + .await + .map_err(|err| format!("fresh reader: {err}"))?; + Ok(()) } From 8bc3e7f049145c7341551b9e8e26067affd345e5 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Fri, 9 Oct 2026 11:55:57 -0700 Subject: [PATCH 3/4] docs(net): SUBSCRIBE_OK Group is not waited below, not unavailable The publisher now delivers a group below SUBSCRIBE_OK's Group when the floor allows it, so the draft and subscriber comments that called those groups unavailable contradicted the behavior. Also collapse the JS publisher's pin to a single branch. Co-Authored-By: Claude Opus 5.5 --- drafts/draft-lcurley-moq-lite.md | 8 ++++---- js/net/src/lite/publisher.ts | 7 +++---- js/net/src/lite/subscriber.ts | 4 ++-- rs/moq-net/src/lite/publisher.rs | 2 +- rs/moq-net/src/lite/subscriber.rs | 5 +++-- rs/moq-net/src/model/track.rs | 8 ++++---- 6 files changed, 17 insertions(+), 17 deletions(-) diff --git a/drafts/draft-lcurley-moq-lite.md b/drafts/draft-lcurley-moq-lite.md index ed9043519a..172b305496 100644 --- a/drafts/draft-lcurley-moq-lite.md +++ b/drafts/draft-lcurley-moq-lite.md @@ -1188,9 +1188,9 @@ Set to 0x0 to indicate a SUBSCRIBE_OK message. **Group**: The absolute sequence number of the first group served when the subscription resolves. It MUST be greater than or equal to the requested `Group Start`. -Groups between the requested floor and this group are unavailable. -A group at or above the requested floor that is still within `Subscriber Max Age`, created or becoming available later, is delivered. -This group is not itself a new floor. +This group is not a new floor. +A subscriber SHOULD NOT wait for a group between the requested floor and this group. +The publisher still delivers such a group if it is created or becomes available later, while it is within `Subscriber Max Age`. A subscriber whose `Subscriber Max Age` resolves the start to the latest group learns that sequence here. There is no matching frame field, because the start frame is never in doubt: a partial group is only delivered when it was asked for, so the subscription starts either exactly where it asked or at the beginning of a later group (see [Positions](#positions)). @@ -1391,7 +1391,7 @@ The `Message Length` describes the payload size on the wire. ## moq-lite-07 -- SUBSCRIBE_OK `Group` names the first group served when the subscription resolves. Groups between the requested floor and that group are unavailable. A later group at or above the floor and within Subscriber Max Age is still delivered. The announced group is not a new floor. +- SUBSCRIBE_OK `Group` names the first group served when the subscription resolves. It is not a new floor. A subscriber does not wait for a group between the requested floor and it, but a publisher still delivers one that arrives later within Subscriber Max Age. - Corrected the SUBSCRIBE note that offset `Group Start` by 1. Only `Group End` and `Frame End` are offset so 0 can mean absent. `Group Start` is an absolute sequence, and 0 is group 0. - A subscription's range bounds datagrams like groups, and FETCH never returns a datagram. - Assigned 0x3A NOT_FETCHABLE in the stream error table: a FETCH for a group delivered only as a datagram. diff --git a/js/net/src/lite/publisher.ts b/js/net/src/lite/publisher.ts index f08cacf6e5..3f6e711a95 100644 --- a/js/net/src/lite/publisher.ts +++ b/js/net/src/lite/publisher.ts @@ -960,11 +960,10 @@ export class Publisher { // A pre-06 subscribe that named no group is the exception: an absent // Group Start is the latest group, and this sequence becomes the floor. // Lite-06 encodes a floor of group 0 as 0, so an omitted start is not pinned. - const pin = bounds.startGroup === undefined && !resolvesStart(this.version); - if (pin || bounds.endGroup !== undefined) { + if (bounds.startGroup === undefined && !resolvesStart(this.version)) { hooks.replaceGroups(track, { - ...(pin ? { start: { included: group.sequence } } : {}), - ...(bounds.endGroup === undefined ? {} : { end: { included: bounds.endGroup } }), + start: { included: group.sequence }, + end: bounds.endGroup === undefined ? undefined : { included: bounds.endGroup }, }); } if ( diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index 571c8666c9..96c306c63a 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -1008,8 +1008,8 @@ export class Subscriber { if ("start" in resp) { entry.start = resp.start.group; - // The groups the SUBSCRIBE asked for below it are unavailable, whatever the - // demand asks later. + // The groups the SUBSCRIBE asked for below it are not waited for, whatever the + // demand asks later. One that still arrives is delivered. if (entry.requested !== undefined) entry.tail.account(entry.requested, entry.start); } else if ("end" in resp) { if (entry.end !== undefined) throw new ProtocolViolation("duplicate SUBSCRIBE_END"); diff --git a/rs/moq-net/src/lite/publisher.rs b/rs/moq-net/src/lite/publisher.rs index 92cfcd71e8..f7407b8390 100644 --- a/rs/moq-net/src/lite/publisher.rs +++ b/rs/moq-net/src/lite/publisher.rs @@ -2845,7 +2845,7 @@ impl TrackRun { largest, }))?; // SUBSCRIBE_START names the first group served now. Groups already skipped between - // the requested floor and this group stay unavailable. A later group at or above an + // the requested floor and this group are not served. A later group at or above an // explicit floor, still inside the subscriber's max age, is delivered, so the cursor // stays at that floor. // diff --git a/rs/moq-net/src/lite/subscriber.rs b/rs/moq-net/src/lite/subscriber.rs index f7c8ed0470..beffd67947 100644 --- a/rs/moq-net/src/lite/subscriber.rs +++ b/rs/moq-net/src/lite/subscriber.rs @@ -4577,8 +4577,9 @@ impl ServeLoop { let _ = self.serving.start_at(start.group); } active.served = Some(start.group); - // The groups the SUBSCRIBE asked for below it are - // unavailable, whatever the demand asks later. + // The groups the SUBSCRIBE asked for below it are not + // waited for, whatever the demand asks later. One that + // still arrives is delivered. if let Some(requested) = active.requested && let Ok(mut tail) = active.tail.write() { diff --git a/rs/moq-net/src/model/track.rs b/rs/moq-net/src/model/track.rs index c6c0716d64..dfa1833771 100644 --- a/rs/moq-net/src/model/track.rs +++ b/rs/moq-net/src/model/track.rs @@ -290,8 +290,8 @@ pub(crate) struct TrackState { closed: bool, // The first sequence the live feed serves, once the publisher declared one - // (the wire's SUBSCRIBE_START). Lower groups never arrive on their own; a - // fetch can still create them. + // (the wire's SUBSCRIBE_START). Lower groups are not waited for, though one + // may still arrive late or be fetched. start_sequence: Option, // Whether `start_sequence` is only the floor a subscription asked for, still @@ -1875,8 +1875,8 @@ impl Producer { /// Declare the first group the live feed serves (the wire's SUBSCRIBE_START, /// or the start the subscription itself requested): groups below `sequence` - /// will never arrive on their own, so a reader waiting for one fails over - /// instead of stalling. A fetch can still retrieve them. + /// are not waited for, so a reader waiting for one fails over instead of + /// stalling. One may still arrive late, and a fetch can still retrieve them. /// /// Scoped to the current subscription's demand, so a later declaration /// replaces this one in either direction: a re-subscription may start From d79f41991942ecd01a57ae85ffe6fc1ecef3ba38 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Fri, 9 Oct 2026 16:26:02 -0700 Subject: [PATCH 4/4] docs(net): SUBSCRIBE_OK Group is where delivery starts, not the first group sent Co-Authored-By: Claude Opus 5.5 --- drafts/draft-lcurley-moq-lite.md | 5 +++-- rs/moq-net/src/lite/publisher.rs | 4 ++-- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/drafts/draft-lcurley-moq-lite.md b/drafts/draft-lcurley-moq-lite.md index 2391a1a745..dc963c4150 100644 --- a/drafts/draft-lcurley-moq-lite.md +++ b/drafts/draft-lcurley-moq-lite.md @@ -1217,7 +1217,8 @@ SUBSCRIBE_OK Message { Set to 0x0 to indicate a SUBSCRIBE_OK message. **Group**: -The absolute sequence number of the first group served when the subscription resolves. +The absolute sequence number where delivery starts when the subscription resolves. +Group streams can arrive out of order, so it need not be the first group sent. It MUST be greater than or equal to the requested `Group Start`. This group is not a new floor. A subscriber SHOULD NOT wait for a group between the requested floor and this group. @@ -1422,7 +1423,7 @@ The `Message Length` describes the payload size on the wire. ## moq-lite-07 -- SUBSCRIBE_OK `Group` names the first group served when the subscription resolves. It is not a new floor. A subscriber does not wait for a group between the requested floor and it, but a publisher still delivers one that arrives later within Subscriber Max Age. +- SUBSCRIBE_OK `Group` names where delivery starts when the subscription resolves, which need not be the first group sent. It is not a new floor. A subscriber does not wait for a group between the requested floor and it, but a publisher still delivers one that arrives later within Subscriber Max Age. - Corrected the SUBSCRIBE note that offset `Group Start` by 1. Only `Group End` and `Frame End` are offset so 0 can mean absent. `Group Start` is an absolute sequence, and 0 is group 0. - A subscription's range bounds datagrams like groups, and FETCH never returns a datagram. - Assigned 0x3A NOT_FETCHABLE in the stream error table: a FETCH for a group delivered only as a datagram. diff --git a/rs/moq-net/src/lite/publisher.rs b/rs/moq-net/src/lite/publisher.rs index b625adf85e..fe3f9f3051 100644 --- a/rs/moq-net/src/lite/publisher.rs +++ b/rs/moq-net/src/lite/publisher.rs @@ -2958,13 +2958,13 @@ impl TrackRun { group: start, largest, }))?; - // SUBSCRIBE_START names the first group served now. Groups already skipped between + // SUBSCRIBE_START names where delivery starts. Groups already skipped between // the requested floor and this group are not served. A later group at or above an // explicit floor, still inside the subscriber's max age, is delivered, so the cursor // stays at that floor. // // A pre-06 subscription that named no group is the exception: those drafts define an - // absent Group Start as the latest group, and the first group served becomes the + // absent Group Start as the latest group, and the resolved start becomes the // floor. Lite-06 encodes a floor of group 0 as 0, which decodes as no named floor; // that 0 is group 0, so it is not pinned. Raised, not assigned: an update that // landed while the group was held may already have raised it past.