diff --git a/doc/lib/js/net.md b/doc/lib/js/net.md index b0d4fda217..9a7b4dd01e 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 `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. `fetchGroup(sequence, { signal })` abandons a pending fetch with the signal's reason; the shared stream is cancelled only once every caller has left. Once resolved, close the group instead. - **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. A `Broadcast.Consumer` names its broadcast by `path`: the path it was requested at relative to the origin handle's root, or empty for a standalone broadcast. Those references resolve against it. 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.test.ts b/js/net/src/broadcast.test.ts index f1640b4af3..9c77d51b06 100644 --- a/js/net/src/broadcast.test.ts +++ b/js/net/src/broadcast.test.ts @@ -487,3 +487,19 @@ test("a fetch waits for a group still to come", async () => { broadcast.close(); }); + +test("aborting a fetch rejects with the signal's reason", async () => { + const broadcast = new BroadcastProducer(); + broadcast.createTrack("video"); + + const early = new Error("early"); + await expect(wireOf(broadcast).fetchGroup("video", 0, { signal: AbortSignal.abort(early) })).rejects.toBe(early); + + const controller = new AbortController(); + const pending = wireOf(broadcast).fetchGroup("video", 0, { signal: controller.signal }); + const late = new Error("late"); + controller.abort(late); + await expect(pending).rejects.toBe(late); + + broadcast.close(); +}); diff --git a/js/net/src/broadcast.ts b/js/net/src/broadcast.ts index cf1a099510..fa94ed42dd 100644 --- a/js/net/src/broadcast.ts +++ b/js/net/src/broadcast.ts @@ -10,6 +10,7 @@ import { Route } from "./hop.ts"; import { hooks, type TrackSequence } from "./internal.ts"; import * as Path from "./path.ts"; import * as track from "./track.ts"; +import { untilAborted } from "./util/abort.ts"; import { registerWire, trackOf, type Broadcast as Wire } from "./wire.ts"; /** The origin callback a created broadcast uses to advertise its exact path. @internal */ @@ -129,11 +130,12 @@ async function fetchGroup( sequence: number, options: track.FetchGroupOptions = {}, ): Promise { + options.signal?.throwIfAborted(); const subscriber = subscribe(state, name, { priority: options.priority }); hooks.exemptFetch(subscriber); try { for (;;) { - const group = await subscriber.recvGroup(); + const group = await untilAborted(subscriber.recvGroup(), options.signal); if (!group) throw new NotFound(`group ${sequence}`); if (group.sequence === sequence) { // Close the subscription when the returned group finishes, not now: an diff --git a/js/net/src/lite/subscriber.test.ts b/js/net/src/lite/subscriber.test.ts index f9dec8b901..48a931839e 100644 --- a/js/net/src/lite/subscriber.test.ts +++ b/js/net/src/lite/subscriber.test.ts @@ -1,7 +1,7 @@ import { expect, spyOn, test } from "bun:test"; import { Signal } from "@moq/signals"; import type { Probe as ProbeStats } from "../connection/stats.ts"; -import { error, reason, StreamCode, StreamError } from "../error.ts"; +import { error, fromTransport, reason, StreamCode, StreamError } from "../error.ts"; import { HopSchema, isAnonymous, MAX_HOPS, Route, UNKNOWN_HOP } from "../hop.ts"; import * as Path from "../path.ts"; import { Writer } from "../stream.ts"; @@ -803,3 +803,85 @@ test("a fetch started after the subscriber closes rejects without opening a stre expectCut(err, undefined); expect(streams.length).toBe(0); }); + +test("an already-aborted fetch rejects without opening a stream", async () => { + const { quic, streams } = fakeSession(); + const subscriber = new Subscriber(quic, Version.DRAFT_05, HopSchema.parse(1n)); + const cause = new Error("gone"); + + const err = await subscriber + .fetchGroup(Path.from("room"), "video", 0, { signal: AbortSignal.abort(cause) }) + .catch((err: unknown) => err); + expect(err).toBe(cause); + expect(streams.length).toBe(0); +}); + +// Coalesced fetches share one FETCH stream. An abort releases only that caller's share; the +// stream is cancelled once the last sharer leaves, before its FETCH is sent if it can be. +test("one of two fetch sharers aborting leaves the other's fetch", async () => { + const { quic, streams } = fakeSession(); + const subscriber = new Subscriber(quic, Version.DRAFT_05, HopSchema.parse(1n)); + + const controller = new AbortController(); + const a = subscriber.fetchGroup(Path.from("room"), "video", 0, { signal: controller.signal }); + const b = subscriber.fetchGroup(Path.from("room"), "video", 0); + + await drainUntil(() => streams.length === 1); + await answerTrackInfo(streams[0]); + await drainUntil(() => streams.length === 2); + await streams[1].reading; + + const cause = new Error("gone"); + controller.abort(cause); + expect(await a.catch((err: unknown) => err)).toBe(cause); + + let aborted = false; + void streams[1].aborted.then(() => { + aborted = true; + }); + // An empty-group FIN accepts the fetch. + streams[1].inbound.close(); + const group = await b; + expect(await group.readFrame()).toBeUndefined(); + expect(aborted).toBe(false); + + subscriber.close(); +}); + +test.each([ + ["the TRACK_INFO", "track"], + ["the FETCH", "fetch"], +] as const)("the last fetch sharer aborting during %s cancels it", async (_, stage) => { + const { quic, streams } = fakeSession(); + const subscriber = new Subscriber(quic, Version.DRAFT_05, HopSchema.parse(1n)); + + const first = new AbortController(); + const second = new AbortController(); + const a = subscriber.fetchGroup(Path.from("room"), "video", 0, { signal: first.signal }); + const b = subscriber.fetchGroup(Path.from("room"), "video", 0, { signal: second.signal }); + + await drainUntil(() => streams.length === 1); + await streams[0].reading; + if (stage === "fetch") { + await answerTrackInfo(streams[0]); + await drainUntil(() => streams.length === 2); + await streams[1].reading; + } + + first.abort(new Error("first")); + second.abort(new Error("second")); + expect(((await a.catch((err: unknown) => err)) as Error).message).toBe("first"); + expect(((await b.catch((err: unknown) => err)) as Error).message).toBe("second"); + + if (stage === "track") { + // The TRACK_INFO still completes, but no FETCH is sent for the abandoned group. + await answerTrackInfo(streams[0]); + for (let i = 0; i < 100; i++) await Promise.resolve(); + expect(streams.length).toBe(1); + } else { + const err = fromTransport(await streams[1].aborted) as StreamError; + expect(err.code).toBe(StreamCode.Cancel); + } + + subscriber.close(); +}); diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index 61c3e7aaf2..b7d9a8b7b0 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -12,6 +12,7 @@ import { type OpenOptions, type Reader, Stream } from "../stream.ts"; import { TAIL_GRACE_MS, Tail } from "../tail.ts"; import * as Time from "../time.ts"; import type * as track from "../track.ts"; +import { untilAborted } from "../util/abort.ts"; import { TimeoutError, withTimeout } from "../util/timeout.ts"; import { overrideBroadcastWire, wireOf } from "../wire.ts"; import { @@ -731,24 +732,32 @@ export class Subscriber { sequence: number, options: track.FetchGroupOptions = {}, ): Promise { + options.signal?.throwIfAborted(); + // 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. + // + // Reserve each caller's mirror before the fetch starts or is awaited: the fetch watches + // demand from the start, and a fast FIN cannot discard frames before these callers + // receive their handles. An abort closes only this caller's mirror, so the stream is + // cancelled once the last one leaves. const key = JSON.stringify([broadcast, track, sequence]); let entry = this.#fetches.get(key); - if (!entry || entry.group.isClosed) { + let consumer: netGroup.Consumer; + if (entry && !entry.group.isClosed) { + consumer = entry.group.mirror(); + } else { const group = new netGroup.Producer(sequence); - entry = { group, accepted: this.#runFetch(broadcast, track, sequence, options, group) }; + consumer = group.mirror(); + entry = { group, accepted: this.#runFetch(broadcast, track, sequence, options.priority ?? 0, group) }; this.#fetches.set(key, entry); void group.closed.then(() => { if (this.#fetches.get(key)?.group === group) this.#fetches.delete(key); }); } - // 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; + await untilAborted(entry.accepted, options.signal); return consumer; } catch (err) { consumer.close(); @@ -757,12 +766,13 @@ export class Subscriber { } // Open the FETCH stream and pump the response into the shared group. Setup errors close the - // group, evict the entry, and reject every caller waiting for acceptance. + // group, evict the entry, and reject every caller waiting for acceptance. A setup every caller + // has abandoned is cancelled the same way. async #runFetch( broadcast: Path.Valid, track: string, sequence: number, - options: track.FetchGroupOptions, + priority: number, group: netGroup.Producer, ): Promise { try { @@ -773,10 +783,10 @@ export class Subscriber { // Lite has no FETCH_OK, so a publisher that never answers would hold the setup forever. // Subscriber.close() closing the group releases every caller at any stage, and resets // the streams the setup opened. - const setup = this.#fetchSetup(broadcast, track, sequence, options); + const setup = this.#fetchSetup(broadcast, track, sequence, priority, group); let accepted: { stream: Stream; info: TrackInfo }; try { - accepted = await untilClosed(group, setup); + accepted = await untilAbandoned(group, setup); } catch (err: unknown) { // A setup that finishes just after the close hands back a stream nobody will read. void setup.then( @@ -794,20 +804,21 @@ export class Subscriber { } // Resolve the track's timescale, then open the FETCH stream and wait for it to be accepted. + // Closing the group during that wait resets the stream. async #fetchSetup( broadcast: Path.Valid, track: string, sequence: number, - options: track.FetchGroupOptions, + priority: number, + group: netGroup.Producer, ): Promise<{ stream: Stream; info: TrackInfo }> { - const info = await this.#trackInfo(broadcast, track); - const priority = options.priority ?? 0; + const info = await untilClosed(group, this.#trackInfo(broadcast, track)); return this.#exchange({ sendOrder: sendOrder({ priority }) }, async (stream) => { await stream.writer.u53(StreamId.Fetch); await new FetchMessage({ broadcast, track, priority, group: sequence }).encode(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(); + await untilClosed(group, stream.reader.done()); return { stream, info }; }); } @@ -1195,6 +1206,17 @@ async function untilClosed(group: netGroup.Producer, step: Promise): Promi return value as T; } +// Like untilClosed, but also cancels once every reader has left. Demand is level-triggered, so a +// caller that coalesces onto the group before the check re-arms it. +async function untilAbandoned(group: netGroup.Producer, step: Promise): Promise { + const idle: unique symbol = Symbol("idle"); + for (;;) { + const value = await untilClosed(group, race([step, group.unused().then((): typeof idle => idle)])); + if (value !== idle) return value as T; + if (!group.used.peek()) throw new StreamError(StreamCode.Cancel, { message: "cancel" }); + } +} + /** * A broadcast consumed from a lite session. It resolves `track.Consumer.query()` and * `.fetchGroup()` over the wire (lite-05+ TRACK / FETCH streams) by reaching into the diff --git a/js/net/src/track.ts b/js/net/src/track.ts index 09af757c9e..52dc4034e0 100644 --- a/js/net/src/track.ts +++ b/js/net/src/track.ts @@ -243,6 +243,14 @@ export class Request { export interface FetchGroupOptions { /** Delivery priority for the fetch stream. Defaults to `0`. */ priority?: number; + + /** + * Abandons this fetch, rejecting with the signal's reason. Concurrent fetches of the same + * group share one stream, cancelled only once every caller has left. An already-aborted + * signal rejects before anything is sent, and aborting after the group resolves has no + * effect; close the group instead. + */ + signal?: AbortSignal; } /** diff --git a/js/net/src/util/abort.ts b/js/net/src/util/abort.ts new file mode 100644 index 0000000000..ed1f32ba8b --- /dev/null +++ b/js/net/src/util/abort.ts @@ -0,0 +1,18 @@ +import { race } from "@moq/signals"; + +// Settle with `promise`, or reject with the signal's reason once it aborts. An already-aborted +// signal rejects at once. There's no cancellation of `promise` itself; the caller releases +// whatever it holds when this rejects. +export async function untilAborted(promise: Promise, signal?: AbortSignal): Promise { + if (!signal) return promise; + signal.throwIfAborted(); + + const { promise: aborted, reject } = Promise.withResolvers(); + const onAbort = () => reject(signal.reason); + signal.addEventListener("abort", onAbort, { once: true }); + try { + return await race([promise, aborted]); + } finally { + signal.removeEventListener("abort", onAbort); + } +} diff --git a/quest/m1/README.md b/quest/m1/README.md index 03a1dec8be..3c57823667 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -67,7 +67,6 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. - [JS IETF datagrams](/quest/m1/js-ietf-datagram.md) - `@moq/net` sends and receives datagram groups over moq-transport, like Rust - [#2991](/quest/m1/2991-net-coalesce-dynamic-tracks-and-preserve-sequences-across.md) - one dynamic producer per track name in both languages, with the sequence namespace surviving a replacement - [JavaScript FETCH](/quest/m1/js-fetch.md) - generic on-demand group serving and IETF FETCH for browser publishers -- [JS fetch cancel](/quest/m1/js-fetch-cancel.md) - `FetchGroupOptions.signal` abandons one pending fetch without closing the track - [Archive](/quest/m1/archive/README.md) - record selected tracks to any object_store and replay them over FETCH or derived HLS; the catalog entry and format may break in place, since no archives exist - [Tooling](/quest/m1/tooling/README.md) - justfiles become a one-line menu over `sh/`, one impact map scopes CI, and every workflow step runs a recipe - [Path patterns](/quest/m1/path-patterns.md) - one matcher for every predicate over broadcast paths: tokens, origins, interest diff --git a/quest/m1/js-fetch-cancel.md b/quest/m1/js-fetch-cancel.md deleted file mode 100644 index e7371fb4fd..0000000000 --- a/quest/m1/js-fetch-cancel.md +++ /dev/null @@ -1,25 +0,0 @@ -# [S] A JS fetch can be cancelled - -## Goal - -A `@moq/net` caller can abandon one pending group fetch without closing the -track or session, the follow-up #4357 left. - -## Plan - -Decided 2026-09-28: `FetchGroupOptions` (`js/net/src/track.ts`) gains -`signal?: AbortSignal`, the standard JS idiom, and additive. - -- Thread it through `js/net/src/broadcast.ts` and `lite/subscriber.ts`. - `ietf/subscriber.ts` rejects `fetchGroup` today (moq-transport has no - one-shot group fetch), so it stays unchanged. -- Fetches for the same group share one stream (`lite/subscriber.ts`), so an - abort releases this caller's share and rejects its promise with the - signal's reason. The stream is cancelled only when the last sharer leaves. -- An already-aborted signal rejects before anything is sent. -- Tests: one of two sharers aborts and the other still receives the group; - the last sharer aborting cancels the stream. -- Document the option where fetch is documented under `doc/`. - -Public API: additive `FetchGroupOptions.signal`. Wire: none new; the -existing cancel path.