From 1cbc868c437f43967cf4b851178158256aca9c26 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Fri, 25 Sep 2026 18:24:19 -0700 Subject: [PATCH 1/3] chore: claim JS fetch answer quest Co-Authored-By: GPT-6 From 5cfebab0c9ec9a728682f8f2a15898a754a1387e Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Fri, 25 Sep 2026 18:30:29 -0700 Subject: [PATCH 2/3] fix(net): await the publisher answer for lite fetches Share fetch acceptance across coalesced callers and report local group misses with NotFound. Co-Authored-By: GPT-6 --- doc/lib/js/net.md | 2 +- js/net/src/broadcast.ts | 5 +- js/net/src/integration.test.ts | 90 ++++++++++++++++++++++++++++++++++ js/net/src/lite/subscriber.ts | 45 ++++++++++------- quest/m1/README.md | 1 - quest/m1/js-fetch-answer.md | 32 ------------ 6 files changed, 120 insertions(+), 55 deletions(-) delete mode 100644 quest/m1/js-fetch-answer.md diff --git a/doc/lib/js/net.md b/doc/lib/js/net.md index f086c11cd9..61e64b17b1 100644 --- a/doc/lib/js/net.md +++ b/doc/lib/js/net.md @@ -53,7 +53,7 @@ for (;;) { - **Discovery** by any pattern scope (`origin.announced(scope)`, such as `room/*/chat`; default everything). Each event's `prefix` is the covered prefix relative to the origin, `captures` reports what the scope's wildcards matched when the prefix pins them, and `kind` says whether it was announced, updated, or retracted. The consumer is an async iterable. `origin.broadcasts(scope)` is a live `Getter>` of the same covered prefixes for UIs that need the current set. A borrowed `Connection.origin` also exposes `dynamic(prefix, route)` for serving paths on demand. - **Subscriptions** carry a priority, a `Time.Milli` max age, and optional `groups` bounds. Groups arrive out of order and are read frame by frame, with `Error.TooFarBehind` when a reader asks for a frame the group never held and `Error.GroupTooLarge` when a write exceeds the cache budget and aborts the group. - **Track ends**: `close()` ends a track at its live edge, while `finishAt(n)` declares the exclusive end ahead of it and still accepts the groups below. A subscriber reads the end with `final()` or awaits `finished()`. A remote track ends only once every group below its end has arrived or was dropped; one reset before its header arrived is skipped after the subscription's max age on moq-lite (one second without one), or after one second on IETF. -- **Datagrams** on moq-lite 05+ and fetch-by-sequence for history. +- **Datagrams** on moq-lite 05+ and fetch-by-sequence for history. `track.fetchGroup(sequence)` on moq-lite resolves when the publisher sends the first response byte or finishes an empty group. A missing group rejects the fetch with `Error.StreamCode.NotFound`, including every concurrent caller sharing that fetch. - **Errors** live under one namespace: a stream reset throws `Error.Stream` with a `StreamCode`, while a session close gives `Error.Session` with a `SessionCode`. The registries are disjoint, so the same number means different things in each, and 64+ is yours. Named conditions such as `Error.TooFarBehind`, `Error.FrameTooLarge`, and `Error.GroupTooLarge` subclass `Error.Stream`, so one `code` check handles a condition raised here or reported by the peer. IETF streams use their own mapping: cancellation sends CANCELLED, other local failures send INTERNAL\_ERROR, and received codes remain opaque. - **Paths** with `Path.relative` for the cross-broadcast catalog references hang uses. Path patterns (`Path.Pattern`, `Path.Patterns`) are re-exported from [`@moq/pattern`](https://www.npmjs.com/package/@moq/pattern). Literal `Path` stays a coordinate. diff --git a/js/net/src/broadcast.ts b/js/net/src/broadcast.ts index 41cbeca2ef..32bc42ad8e 100644 --- a/js/net/src/broadcast.ts +++ b/js/net/src/broadcast.ts @@ -4,6 +4,7 @@ * @module */ import { type GetPromise, Once, Signal } from "@moq/signals"; +import { NotFound } from "./error.ts"; import type { Consumer as GroupConsumer } from "./group.ts"; import { Route } from "./hop.ts"; import { hooks, type TrackSequence } from "./internal.ts"; @@ -131,7 +132,7 @@ async function fetchGroup( try { for (;;) { const group = await subscriber.recvGroup(); - if (!group) throw new Error(`group not found: ${sequence}`); + if (!group) throw new NotFound(`group ${sequence}`); if (group.sequence === sequence) { // Close the subscription when the returned group finishes, not now: an // in-progress group must keep receiving frames for its lifetime (mirrors @@ -141,7 +142,7 @@ async function fetchGroup( } group.close(); - if (group.sequence > sequence) throw new Error(`group not found: ${sequence}`); + if (group.sequence > sequence) throw new NotFound(`group ${sequence}`); } } catch (err) { subscriber.close(); diff --git a/js/net/src/integration.test.ts b/js/net/src/integration.test.ts index aaed493344..249f4b1bb9 100644 --- a/js/net/src/integration.test.ts +++ b/js/net/src/integration.test.ts @@ -10,6 +10,7 @@ import { type Established, } from "./connection/index.ts"; import { SessionCode, SessionError, StreamCode, StreamError, TooFarBehind } from "./error.ts"; +import { Producer as GroupProducer } from "./group.ts"; import * as Ietf from "./ietf/index.ts"; import * as Lite from "./lite/index.ts"; import { createMockTransportPair } from "./mock.ts"; @@ -868,6 +869,95 @@ test("integration: lite draft-05 fetches a cached group", async () => { server.close(); }); +test.each(["gap", "end"])("integration: lite fetch rejects coalesced misses at %s with NotFound", async (missing) => { + const pair = createMockTransportPair(Lite.ALPN_05); + const origin = new OriginProducer(); + const [client, server] = await Promise.all([ + connect(url, { transport: pair.client }), + accept(pair.server, url, { publish: origin.consume() }), + ]); + const broadcast = publish(origin, Path.from("test")); + let serving: Promise | undefined; + if (missing === "gap") { + const producer = broadcast.createTrack("video"); + const group = new GroupProducer(1); + producer.writeGroup(group); + group.close(); + } else { + serving = (async () => { + for (;;) { + const request = await wireOf(broadcast).requested(); + if (!request) return; + request.accept().close(); + } + })(); + } + const remote = wireOf(client).consume(Path.from("test")); + try { + const track = remote.track("video"); + const results = await Promise.allSettled([track.fetchGroup(0), track.fetchGroup(0)]); + for (const result of results) { + expect(result.status).toBe("rejected"); + if (result.status !== "rejected") throw new Error("missing group was accepted"); + expect(result.reason).toBeInstanceOf(StreamError); + expect(result.reason.code).toBe(StreamCode.NotFound); + } + if (results[0].status === "rejected" && results[1].status === "rejected") { + expect(results[0].reason).toBe(results[1].reason); + } + } finally { + broadcast.close(); + await serving; + remote.close(); + client.close(); + server.close(); + origin.close(); + } +}); + +test.each(["frame", "FIN"])("integration: lite fetch waits for the publisher's first %s", async (answer) => { + const pair = createMockTransportPair(Lite.ALPN_05); + const origin = new OriginProducer(); + const [client, server] = await Promise.all([ + connect(url, { transport: pair.client }), + accept(pair.server, url, { publish: origin.consume() }), + ]); + const broadcast = publish(origin, Path.from("test")); + const producer = broadcast.createTrack("video"); + const group = producer.appendGroup(); + const remote = wireOf(client).consume(Path.from("test")); + try { + let settled = 0; + const fetch = () => + remote + .track("video") + .fetchGroup(0) + .then((consumer) => { + settled++; + return consumer; + }); + const a = fetch(); + const b = fetch(); + // Let the in-memory peer process the request with no response available yet. + await sleep(0); + expect(settled).toBe(0); + if (answer === "frame") group.writeString("accepted"); + else group.close(); + const consumers = await Promise.all([a, b]); + for (const consumer of consumers) { + expect(await consumer.readString()).toBe(answer === "frame" ? "accepted" : undefined); + consumer.close(); + } + } finally { + group.close(); + broadcast.close(); + remote.close(); + client.close(); + server.close(); + origin.close(); + } +}); + test("integration: lite draft-05 coalesces concurrent fetches of one group", async () => { const enc = new TextEncoder(); const dec = new TextDecoder(); diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index eb5aff5569..913abf2998 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -132,7 +132,7 @@ export class Subscriber { // Dedup in-flight one-shot fetches, keyed by [broadcast, track, sequence]. Concurrent (or // repeat, while still open) fetchGroup() calls for the same group share one FETCH stream and // each get an independent mirror; the entry is evicted once the group closes. - #fetches = new Map(); + #fetches = new Map }>(); // The peer's PROBE estimates, written as they arrive (Lite03+ only). #probe?: Signal; @@ -704,7 +704,7 @@ export class Subscriber { // Open a FETCH stream for one group and stream its bare frames into a group, for the // ConsumeBroadcast backing track.Consumer.fetchGroup() (lite-05+). - fetchGroup( + async fetchGroup( broadcast: Path.Valid, track: string, sequence: number, @@ -713,29 +713,37 @@ export class Subscriber { // Coalesce onto a still-open fetch of the same group so we don't open a second FETCH // stream (and re-download it); each caller reads an independent mirror. const key = JSON.stringify([broadcast, track, sequence]); - const existing = this.#fetches.get(key); - if (existing && !existing.isClosed) return Promise.resolve(existing.mirror()); - - // Create and cache the group synchronously (before any await) so a concurrent fetch for - // the same group finds it and coalesces rather than racing to open its own stream. - const group = new netGroup.Producer(sequence); - this.#fetches.set(key, group); - void group.closed.then(() => { - if (this.#fetches.get(key) === group) this.#fetches.delete(key); - }); + let entry = this.#fetches.get(key); + if (!entry || entry.group.isClosed) { + const group = new netGroup.Producer(sequence); + entry = { group, accepted: this.#runFetch(broadcast, track, sequence, options, group) }; + this.#fetches.set(key, entry); + void group.closed.then(() => { + if (this.#fetches.get(key)?.group === group) this.#fetches.delete(key); + }); + } - return this.#runFetch(broadcast, track, sequence, options, group); + // Reserve each caller's mirror before awaiting acceptance so the pump sees demand, + // and a fast FIN cannot discard frames before these callers receive their handles. + const consumer = entry.group.mirror(); + try { + await entry.accepted; + return consumer; + } catch (err) { + consumer.close(); + throw err; + } } // Open the FETCH stream and pump the response into the shared group. Setup errors close the - // group (so coalesced mirrors observe them and the entry evicts) and reject this caller. + // group, evict the entry, and reject every caller waiting for acceptance. async #runFetch( broadcast: Path.Valid, track: string, sequence: number, options: track.FetchGroupOptions, group: netGroup.Producer, - ): Promise { + ): Promise { try { if (!supportsTrackStream(this.version)) { throw new Error("fetch group requires moq-lite-05 or newer"); @@ -751,16 +759,15 @@ export class Subscriber { stream.writer, this.version, ); + // A byte or an empty-group FIN accepts the fetch; a reset rejects it. + // done() buffers that byte so the response pump can decode it normally. + await stream.reader.done(); } catch (err: unknown) { stream.abort(error(err)); throw err; } - // Mint this caller's reader before starting the pump, so the group has demand when the - // pump begins watching it (an abandoned fetch cancels once every reader has left). - const consumer = group.mirror(); void this.#runFetchResponse(stream, group, Time.Timescale(info.timescale)); - return consumer; } catch (err: unknown) { group.close(error(err)); throw err; diff --git a/quest/m1/README.md b/quest/m1/README.md index 2077fb1231..5aef7064c5 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -17,7 +17,6 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. ## Quests -- [JS fetch answer](/quest/m1/js-fetch-answer.md) - js/net's lite fetch settles on the publisher's answer, and a JS publisher's miss resets with NotFound - [libmoq hidden opt-in](/quest/m1/libmoq-hidden.md) - `moq_origin_announced` takes a `hidden` flag so C callers can list `.`-named broadcasts - [lite-07 count settle](/quest/m1/lite-count-settle.md) - moq-lite-07 subscribers stop waiting for a subscription's tail once SUBSCRIBE_END's stream count is reached - [Dropped sources](/quest/m1/dropped-sources.md) - consumers see the producer's real error on every end path, never `Dropped` diff --git a/quest/m1/js-fetch-answer.md b/quest/m1/js-fetch-answer.md deleted file mode 100644 index 0e9bc28b3b..0000000000 --- a/quest/m1/js-fetch-answer.md +++ /dev/null @@ -1,32 +0,0 @@ -# [S] JS lite fetch waits for the publisher's answer - -## Goal - -`js/net`'s lite `fetchGroup` resolves only once the publisher has answered: -the first response byte, or a FIN for an empty group. A missing group rejects -the fetch itself, and every coalesced caller sees the same rejection, instead -of receiving a group whose first `readFrame()` fails. A JS publisher that -cannot serve a group resets the stream with `NotFound`, not a generic error, -so a Rust or JS subscriber can tell a miss from a failure. - -## Plan - -- Rust already behaves this way since #4164, which waits in the lite - subscriber before accepting and rejects on reset. Mirror it: resolve after - the first response byte or an empty-group FIN, and reject on reset. Today - the fetch path returns its mirror before the response arrives; the stream - reader can already block until data or FIN and throw on reset. -- The publisher side throws a plain error for a local miss, which reaches the - wire as a generic reset code. Give it the `NotFound` code the Rust side uses. -- The IETF JS path refuses `fetchGroup` outright and is out of scope. -- Tests in the lite integration suite: a missing group rejects the fetch, a - coalesced second caller rejects too, an existing group is unchanged, and a - JS publisher's miss reaches a subscriber as `NotFound`. - -Public API: none; a behavior change in when `fetchGroup` settles. Wire: no -format change; a miss resets with the existing `NotFound` code instead of a -generic one. - -## Related - -- [#4164](https://github.com/moq-dev/moq/pull/4164) - the same fix in Rust From f392919cc7cd96d48c761637831d773c731e28c1 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Fri, 25 Sep 2026 20:32:02 -0700 Subject: [PATCH 3/3] docs(net): name the fetch miss code StreamCode.NotFound Error.StreamCode is not exported. The code lives at the package root. Co-authored-by: Grok 4.7 --- doc/lib/js/net.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/doc/lib/js/net.md b/doc/lib/js/net.md index 61e64b17b1..833961ec74 100644 --- a/doc/lib/js/net.md +++ b/doc/lib/js/net.md @@ -53,7 +53,7 @@ for (;;) { - **Discovery** by any pattern scope (`origin.announced(scope)`, such as `room/*/chat`; default everything). Each event's `prefix` is the covered prefix relative to the origin, `captures` reports what the scope's wildcards matched when the prefix pins them, and `kind` says whether it was announced, updated, or retracted. The consumer is an async iterable. `origin.broadcasts(scope)` is a live `Getter>` of the same covered prefixes for UIs that need the current set. A borrowed `Connection.origin` also exposes `dynamic(prefix, route)` for serving paths on demand. - **Subscriptions** carry a priority, a `Time.Milli` max age, and optional `groups` bounds. Groups arrive out of order and are read frame by frame, with `Error.TooFarBehind` when a reader asks for a frame the group never held and `Error.GroupTooLarge` when a write exceeds the cache budget and aborts the group. - **Track ends**: `close()` ends a track at its live edge, while `finishAt(n)` declares the exclusive end ahead of it and still accepts the groups below. A subscriber reads the end with `final()` or awaits `finished()`. A remote track ends only once every group below its end has arrived or was dropped; one reset before its header arrived is skipped after the subscription's max age on moq-lite (one second without one), or after one second on IETF. -- **Datagrams** on moq-lite 05+ and fetch-by-sequence for history. `track.fetchGroup(sequence)` on moq-lite resolves when the publisher sends the first response byte or finishes an empty group. A missing group rejects the fetch with `Error.StreamCode.NotFound`, including every concurrent caller sharing that fetch. +- **Datagrams** on moq-lite 05+ and fetch-by-sequence for history. `track.fetchGroup(sequence)` on moq-lite resolves when the publisher sends the first response byte or finishes an empty group. A missing group rejects the fetch with `StreamCode.NotFound`, including every concurrent caller sharing that fetch. - **Errors** live under one namespace: a stream reset throws `Error.Stream` with a `StreamCode`, while a session close gives `Error.Session` with a `SessionCode`. The registries are disjoint, so the same number means different things in each, and 64+ is yours. Named conditions such as `Error.TooFarBehind`, `Error.FrameTooLarge`, and `Error.GroupTooLarge` subclass `Error.Stream`, so one `code` check handles a condition raised here or reported by the peer. IETF streams use their own mapping: cancellation sends CANCELLED, other local failures send INTERNAL\_ERROR, and received codes remain opaque. - **Paths** with `Path.relative` for the cross-broadcast catalog references hang uses. Path patterns (`Path.Pattern`, `Path.Patterns`) are re-exported from [`@moq/pattern`](https://www.npmjs.com/package/@moq/pattern). Literal `Path` stays a coordinate.