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
2 changes: 1 addition & 1 deletion doc/lib/js/net.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<ReadonlyMap<Path.Valid, Route>>` 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.

Expand Down
16 changes: 16 additions & 0 deletions js/net/src/broadcast.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
});
4 changes: 3 additions & 1 deletion js/net/src/broadcast.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 */
Expand Down Expand Up @@ -129,11 +130,12 @@ async function fetchGroup(
sequence: number,
options: track.FetchGroupOptions = {},
): Promise<GroupConsumer> {
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
Expand Down
84 changes: 83 additions & 1 deletion js/net/src/lite/subscriber.test.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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();
});
50 changes: 36 additions & 14 deletions js/net/src/lite/subscriber.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -731,24 +732,32 @@ export class Subscriber {
sequence: number,
options: track.FetchGroupOptions = {},
): Promise<netGroup.Consumer> {
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();
Expand All @@ -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<void> {
try {
Expand All @@ -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(
Expand All @@ -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 };
});
}
Expand Down Expand Up @@ -1195,6 +1206,17 @@ async function untilClosed<T>(group: netGroup.Producer, step: Promise<T>): 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<T>(group: netGroup.Producer, step: Promise<T>): Promise<T> {
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
Expand Down
8 changes: 8 additions & 0 deletions js/net/src/track.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

/**
Expand Down
18 changes: 18 additions & 0 deletions js/net/src/util/abort.ts
Original file line number Diff line number Diff line change
@@ -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<T>(promise: Promise<T>, signal?: AbortSignal): Promise<T> {
if (!signal) return promise;
signal.throwIfAborted();

const { promise: aborted, reject } = Promise.withResolvers<never>();
const onAbort = () => reject(signal.reason);
signal.addEventListener("abort", onAbort, { once: true });
try {
return await race([promise, aborted]);
} finally {
signal.removeEventListener("abort", onAbort);
}
}
1 change: 0 additions & 1 deletion quest/m1/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
25 changes: 0 additions & 25 deletions quest/m1/js-fetch-cancel.md

This file was deleted.

Loading