diff --git a/drafts/draft-lcurley-moq-lite.md b/drafts/draft-lcurley-moq-lite.md index 05148b5838..dc963c4150 100644 --- a/drafts/draft-lcurley-moq-lite.md +++ b/drafts/draft-lcurley-moq-lite.md @@ -1108,7 +1108,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. @@ -1217,9 +1217,13 @@ 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 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. +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)). The subscriber derives the start frame from `Group` and its own request: @@ -1419,6 +1423,8 @@ The `Message Length` describes the payload size on the wire. ## moq-lite-07 +- 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. - The subscriber FINs its Subscribe Stream after settling its tail; graceful session close waits for that FIN or reset. diff --git a/js/net/src/integration.test.ts b/js/net/src/integration.test.ts index be94ee8f8e..a8f0a7987a 100644 --- a/js/net/src/integration.test.ts +++ b/js/net/src/integration.test.ts @@ -171,8 +171,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, maxDelay: Milli(250), @@ -235,8 +234,7 @@ test("integration: lite carries a fractional maxDelay as a whole millisecond", a 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({ maxDelay: 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 8245444087..634e7ca98d 100644 --- a/js/net/src/lite/publisher.test.ts +++ b/js/net/src/lite/publisher.test.ts @@ -693,9 +693,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 { @@ -717,6 +717,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. @@ -758,8 +791,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(); @@ -810,8 +844,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 e087b6a18f..e408b8918c 100644 --- a/js/net/src/lite/publisher.ts +++ b/js/net/src/lite/publisher.ts @@ -971,13 +971,17 @@ 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. + if (bounds.startGroup === undefined && !resolvesStart(this.version)) { + hooks.replaceGroups(track, { + start: { included: group.sequence }, + end: bounds.endGroup === undefined ? undefined : { 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 ea847db059..7310913f27 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 ec6b4b13d8..76ee866698 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 afbab2cca2..cb6b9c603a 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -49,6 +49,7 @@ import { hasDatagrams, hasProbeRtt, hasStreamCount, + resolvesStart, updateSupported, Version, waitsForSubscriberFin, @@ -1046,8 +1047,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"); @@ -1166,9 +1167,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. 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) && + 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 faeec7e24e..54f9501a3e 100644 --- a/js/net/src/track.test.ts +++ b/js/net/src/track.test.ts @@ -395,6 +395,15 @@ test("multiple subscriber options aggregate like Rust", async () => { expect(await none).toBeUndefined(); }); +// 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: 4 } } }); + producer.subscribe({}); + expect(producer.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({ maxDelay: Milli(10_000) }); diff --git a/quest/m1/README.md b/quest/m1/README.md index e06cb20693..2b39d0a3c5 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -57,7 +57,6 @@ ladders. - [Subscribe ranges](/quest/m1/subscribe-ranges/README.md) - a lite-07 SUBSCRIBE asks for past and live ranges in either order and replaces FETCH; relays fill misses by range, including over moq-transport - [Cross-relay bursts re-run](/quest/m1/cross-relay-bursts.md) - condition: the #4349 reporter re-runs their A/B/C comparison against current cdn.moq.pro - [Untimed lite-07](/quest/m1/lite-untimed.md) - lite-07 carries an untimed track in both languages; lite-05/06 write send time -- [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 - [A watch and publish release ships assets()](/quest/m1/assets-release.md) - the release that lets the sites host the worklets - [Dogfood hosted worklets](/quest/m1/dogfood-assets.md) - the moq.pro dashboard hosts the worklets and calls `assets()` after the release - [More tests under load](/quest/m1/test-flakes-2/README.md) - the second round of load-only failures, one quest per flake, fixed at the cause diff --git a/quest/m1/lite-late-lower-group.md b/quest/m1/lite-late-lower-group.md deleted file mode 100644 index a646bad88c..0000000000 --- a/quest/m1/lite-late-lower-group.md +++ /dev/null @@ -1,67 +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_delay` 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_delay`, matching - moq-transport, draft-ietf-moq-transport ("subscriptions without a filter pass - all Objects"), and moq-mux's floor handling (#3258). -- A live-only subscription (no floor, per - [lite-07 Live flag](/quest/m1/lite-live.md)) 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` docs in moq-net (which say `None` equals a floor of - group 0) and the matching js/net docs. -- `raise_start_to` applies only to a live-only subscription - (`floor.is_none()`). With an explicit floor, merged with `live` or not, - `SUBSCRIBE_START` still reports the first served group, but nothing below - it is suppressed; each subscriber's cursor filters locally. -- Aggregation follows [lite-07 Live flag](/quest/m1/lite-live.md): letting - an explicit floor survive mixing with no floor starved a floorless subscriber - until the floor's group existed (#5000's review), so the floor and `Live` - become separate fields and merge as min and OR. -- `start_floor_suppresses_late_lower_arrivals` stays as the live-only case: it - subscribes with no floor, 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 live-only subscriber already receiving group 1, then one with a floor -of group 0, and a fresh group 0 reaches only the second. - -## Required - -- [lite-07 Live flag](/quest/m1/lite-live.md) - floors and `Live` are separate fields, so merging them never starves a subscriber - -## 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 a live-only join can be named instead of vanishing 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/lite/publisher.rs b/rs/moq-net/src/lite/publisher.rs index a7679fa926..d65ab3c297 100644 --- a/rs/moq-net/src/lite/publisher.rs +++ b/rs/moq-net/src/lite/publisher.rs @@ -1835,10 +1835,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; @@ -1873,6 +1872,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_delay(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"); @@ -2930,12 +2962,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 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 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. + 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 2dc96c4047..41708df312 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/lite/subscriber.rs b/rs/moq-net/src/lite/subscriber.rs index 3e3e1008af..747cf8385b 100644 --- a/rs/moq-net/src/lite/subscriber.rs +++ b/rs/moq-net/src/lite/subscriber.rs @@ -4650,8 +4650,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 3f88398321..7577281041 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 @@ -1932,8 +1932,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 diff --git a/rs/moq-net/tests/history_groups.rs b/rs/moq-net/tests/history_groups.rs index 0cf508dd2c..8268385290 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..395eabb0c7 --- /dev/null +++ b/rs/moq-net/tests/late_lower_group.rs @@ -0,0 +1,196 @@ +//! A subscriber with an explicit floor receives a group created below the first group +//! that was served, while it is still inside `max_delay`. +//! +//! 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", None)) + .await + .expect("resolve timeout") + .expect("broadcast resolves"); + + let subscription = Subscription::default().with_max_delay(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_delay(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:#?}"); +} 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(()) }