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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 10 additions & 4 deletions drafts/draft-lcurley-moq-lite.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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.
Expand Down
6 changes: 2 additions & 4 deletions js/net/src/integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down Expand Up @@ -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
Expand Down
49 changes: 42 additions & 7 deletions js/net/src/lite/publisher.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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.
Expand Down Expand Up @@ -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();

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

Expand Down
18 changes: 11 additions & 7 deletions js/net/src/lite/publisher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
37 changes: 28 additions & 9 deletions js/net/src/lite/subscribe.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand All @@ -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 () => {
Expand Down
9 changes: 4 additions & 5 deletions js/net/src/lite/subscribe.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

/**
Expand Down
10 changes: 8 additions & 2 deletions js/net/src/lite/subscriber.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ import {
hasDatagrams,
hasProbeRtt,
hasStreamCount,
resolvesStart,
updateSupported,
Version,
waitsForSubscriberFin,
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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
);
Expand Down
9 changes: 9 additions & 0 deletions js/net/src/track.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) });
Expand Down
1 change: 0 additions & 1 deletion quest/m1/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
67 changes: 0 additions & 67 deletions quest/m1/lite-late-lower-group.md

This file was deleted.

5 changes: 3 additions & 2 deletions quest/m1/lite-live.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading
Loading