diff --git a/js/net/src/lite/subscriber.test.ts b/js/net/src/lite/subscriber.test.ts index a26242cf66..f9dec8b901 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 } from "../error.ts"; +import { error, 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"; @@ -9,6 +9,7 @@ import * as Time from "../time.ts"; import { type AnnounceBroadcast, AnnounceInit, AnnounceOk, encodeAnnounceBroadcast } from "./announce.ts"; import { Probe } from "./probe.ts"; import { Subscriber } from "./subscriber.ts"; +import { TrackInfo } from "./track.ts"; import { Version } from "./version.ts"; test("closing the subscriber suppresses probe stream warnings", async () => { @@ -691,3 +692,114 @@ test("a draft-05 duplicate start is still a restart", async () => { announced.close(); subscriber.close(); }); + +interface FakeStream { + inbound: ReadableStreamDefaultController; + // Resolves once the subscriber waits on a read the test has not answered. + reading: Promise; + aborted: Promise; + // Hands the stream to the subscriber, for an open the session was told to park. + release: () => void; +} + +// A session whose streams the test answers by hand and that never fails them on its own, so +// only Subscriber.close() can end a wait. Opens numbered in `park` wait for `release()`. +function fakeSession(park: number[] = []) { + const streams: FakeStream[] = []; + const quic = { + createBidirectionalStream: () => { + let inbound!: ReadableStreamDefaultController; + let onRead!: () => void; + let onAbort!: (reason: unknown) => void; + let release!: () => void; + const reading = new Promise((resolve) => (onRead = resolve)); + const aborted = new Promise((resolve) => (onAbort = resolve)); + // No high water mark, so pull() means the subscriber is blocked on a read. + const readable = new ReadableStream( + { + start: (controller) => { + inbound = controller; + }, + pull: () => onRead(), + }, + { highWaterMark: 0 }, + ); + const writable = new WritableStream({ abort: (reason) => void onAbort(reason) }); + const opened = new Promise((resolve) => (release = () => resolve({ readable, writable }))); + if (!park.includes(streams.length)) release(); + streams.push({ inbound, reading, aborted, release }); + return opened; + }, + } as unknown as WebTransport; + return { quic, streams }; +} + +async function answerTrackInfo(stream: FakeStream): Promise { + const chunks: Uint8Array[] = []; + const writer = new Writer( + new WritableStream({ write: (chunk) => void chunks.push(new Uint8Array(chunk)) }), + ); + await new TrackInfo({}).encode(writer, Version.DRAFT_05); + for (const chunk of chunks) stream.inbound.enqueue(chunk); + stream.inbound.close(); +} + +function expectCut(err: unknown, cause: Error | undefined) { + if (cause) { + expect(err).toBe(cause); + } else { + expect(err).toBeInstanceOf(StreamError); + expect((err as StreamError).code).toBe(StreamCode.SessionClosed); + } +} + +// Lite has no FETCH_OK, so a publisher that never answers holds each setup stage until the +// subscriber closes. The stream that stage opened is reset, even one opening after the close. +test.each([ + ["the TRACK_INFO", "track", undefined], + ["the FETCH", "fetch", undefined], + ["the FETCH, on a session error", "fetch", new Error("session died")], + ["a stream slot for the FETCH", "open", undefined], +] as const)("closing the subscriber rejects a fetch waiting on %s", async (_, stage, cause) => { + const { quic, streams } = fakeSession(stage === "open" ? [1] : []); + const subscriber = new Subscriber(quic, Version.DRAFT_05, HopSchema.parse(1n)); + + let settled = false; + const fetch = subscriber.fetchGroup(Path.from("room"), "video", 0).then( + () => { + settled = true; + return undefined; + }, + (err: unknown) => { + settled = true; + return err; + }, + ); + + await drainUntil(() => streams.length === 1); + if (stage === "track") { + await streams[0].reading; + } else { + await answerTrackInfo(streams[0]); + await drainUntil(() => streams.length === 2); + if (stage === "fetch") await streams[1].reading; + } + expect(settled).toBe(false); + + subscriber.close(cause); + expectCut(await fetch, cause); + + const stuck = streams[stage === "track" ? 0 : 1]; + stuck.release(); + await stuck.aborted; +}); + +test("a fetch started after the subscriber closes rejects without opening a stream", async () => { + const { quic, streams } = fakeSession(); + const subscriber = new Subscriber(quic, Version.DRAFT_05, HopSchema.parse(1n)); + subscriber.close(); + + const err = await subscriber.fetchGroup(Path.from("room"), "video", 0).catch((err: unknown) => err); + expectCut(err, undefined); + expect(streams.length).toBe(0); +}); diff --git a/js/net/src/lite/subscriber.ts b/js/net/src/lite/subscriber.ts index eb55837234..61c3e7aaf2 100644 --- a/js/net/src/lite/subscriber.ts +++ b/js/net/src/lite/subscriber.ts @@ -8,7 +8,7 @@ import * as netGroup from "../group.ts"; import { Cost, type Hop, MAX_HOPS, type Route, routesEqual, UNKNOWN_HOP } from "../hop.ts"; import { groupBounds, hiddenBelow, scopeCaptures, scopeHead, scopeOverlaps } from "../internal.ts"; import * as Path from "../path.ts"; -import { type Reader, Stream } from "../stream.ts"; +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"; @@ -672,17 +672,33 @@ export class Subscriber { // Opens a TRACK stream, reads the single TRACK_INFO, and FINs. Lite-05+ only. async #trackInfo(broadcast: Path.Valid, track: string): Promise { - const stream = await Stream.open(this.#quic); - try { + return this.#exchange(undefined, async (stream) => { await stream.writer.u53(StreamId.Track); await new TrackMessage(broadcast, track).encode(stream.writer, this.version); const info = await TrackInfo.decode(stream.reader, this.version); // The publisher FINs after TRACK_INFO; FIN our side too. stream.close(); return info; + }); + } + + // Opens a stream and runs a request/response exchange on it, resetting the stream if `run` + // fails. Subscriber.close() also resets it while `run` is pending, so a peer that never + // answers cannot hold it open, and a stream that opens after the close is reset at once. + async #exchange(options: OpenOptions | undefined, run: (stream: Stream) => Promise): Promise { + const closed = this.#closed.signal; + closed.throwIfAborted(); + const stream = await Stream.open(this.#quic, options); + const abort = () => stream.abort(error(closed.reason)); + closed.addEventListener("abort", abort); + try { + closed.throwIfAborted(); + return await run(stream); } catch (err) { stream.abort(error(err)); throw err; + } finally { + closed.removeEventListener("abort", abort); } } @@ -754,31 +770,48 @@ export class Subscriber { throw new Error("fetch group requires moq-lite-05 or newer"); } - const info = await this.#trackInfo(broadcast, track); - const priority = options.priority ?? 0; - const stream = await Stream.open(this.#quic, { sendOrder: sendOrder({ priority }) }); - + // 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); + let accepted: { stream: Stream; info: TrackInfo }; try { - 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(); + accepted = await untilClosed(group, setup); } catch (err: unknown) { - stream.abort(error(err)); + // A setup that finishes just after the close hands back a stream nobody will read. + void setup.then( + ({ stream }) => stream.abort(error(err)), + () => void 0, + ); throw err; } - void this.#runFetchResponse(stream, group, Time.Timescale(info.timescale)); + void this.#runFetchResponse(accepted.stream, group, Time.Timescale(accepted.info.timescale)); } catch (err: unknown) { group.close(error(err)); throw err; } } + // Resolve the track's timescale, then open the FETCH stream and wait for it to be accepted. + async #fetchSetup( + broadcast: Path.Valid, + track: string, + sequence: number, + options: track.FetchGroupOptions, + ): Promise<{ stream: Stream; info: TrackInfo }> { + const info = await this.#trackInfo(broadcast, track); + const priority = options.priority ?? 0; + 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(); + return { stream, info }; + }); + } + // Read the FETCH response (bare zigzag-delta-timestamped frames) into the group, then // FIN. A stream-level failure aborts the group so its reader observes the gap. async #runFetchResponse(stream: Stream, group: netGroup.Producer, timescale: Time.Timescale): Promise { @@ -1135,16 +1168,33 @@ export class Subscriber { * session died, since those tracks were cut off rather than ended. */ close(err?: Error) { - this.#closed.abort(); + // A fetch or setup exchange cut off by the session is incomplete even on a deliberate + // close, so it always ends with an error. + const cut = err ?? new StreamError(StreamCode.SessionClosed, { message: "session closed" }); + this.#closed.abort(cut); for (const { track } of this.#subscribes.values()) { track.close(err); } this.#subscribes.clear(); + + // This also releases callers still awaiting acceptance. + for (const { group } of this.#fetches.values()) { + group.close(cut); + } } } +// Settles with `step`, or rejects with the group's error once it closes first. A publisher +// may never answer a FETCH, so Subscriber.close() closing the group is what releases it. +async function untilClosed(group: netGroup.Producer, step: Promise): Promise { + const value = await race([step, group.closed]); + const closed = group.closed.peek(); + if (closed !== undefined) throw closed ?? new Error("fetch closed before it was accepted"); + return value as T; +} + /** * 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