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
97 changes: 92 additions & 5 deletions js/net/src/lite/subscriber.test.ts
Original file line number Diff line number Diff line change
@@ -1,15 +1,16 @@
import { expect, spyOn, test } from "bun:test";
import { expect, jest, spyOn, test } from "bun:test";
import { Signal } from "@moq/signals";
import type { Probe as ProbeStats } from "../connection/stats.ts";
import * as Epoch from "../epoch.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";
import { type Reader, Writer } from "../stream.ts";
import * as Time from "../time.ts";
import { type AnnounceBroadcast, AnnounceInit, AnnounceOk, encodeAnnounceBroadcast } from "./announce.ts";
import { Group } from "./group.ts";
import { Probe } from "./probe.ts";
import { Subscriber } from "./subscriber.ts";
import { SUBSCRIBE_SETUP_TIMEOUT_MS, Subscriber } from "./subscriber.ts";
import { TrackInfo } from "./track.ts";
import { Version } from "./version.ts";

Expand Down Expand Up @@ -785,6 +786,8 @@ interface FakeStream {
// Resolves once the subscriber waits on a read the test has not answered.
reading: Promise<void>;
aborted: Promise<unknown>;
// Every chunk the subscriber wrote.
written: Uint8Array[];
// Hands the stream to the subscriber, for an open the session was told to park.
release: () => void;
}
Expand All @@ -811,10 +814,14 @@ function fakeSession(park: number[] = []) {
},
{ highWaterMark: 0 },
);
const writable = new WritableStream<Uint8Array>({ abort: (reason) => void onAbort(reason) });
const written: Uint8Array[] = [];
const writable = new WritableStream<Uint8Array>({
write: (chunk) => void written.push(chunk),
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 });
streams.push({ inbound, reading, aborted, written, release });
return opened;
},
} as unknown as WebTransport;
Expand Down Expand Up @@ -882,6 +889,86 @@ test.each([
await stuck.aborted;
});

// A setup that outlived its deadline is over: its TRACK stream is reset, even one still waiting for a
// slot, and the subscription is neither registered again nor sent as a SUBSCRIBE.
test.each([
["lite-05 subscribe waiting on the TRACK_INFO", Version.DRAFT_05, false],
["lite-06 subscribe waiting on the TRACK_INFO", Version.DRAFT_06, false],
["lite-07 subscribe waiting on the TRACK_INFO", Version.DRAFT_07, false],
["lite-05 subscribe waiting on a stream slot for the TRACK", Version.DRAFT_05, true],
] as const)("a %s that times out leaves nothing behind", async (_, version, parked) => {
jest.useFakeTimers();
const warn = spyOn(console, "warn").mockImplementation(() => {});
const { quic, streams } = fakeSession(parked ? [0] : []);
const subscriber = new Subscriber(quic, version, HopSchema.parse(1n));
try {
const track = subscriber.consume(Path.from("room")).track("video").subscribe();
await drainUntil(() => streams.length === 1);
if (!parked) await streams[0].reading;

jest.advanceTimersByTime(SUBSCRIBE_SETUP_TIMEOUT_MS);
await drainUntil(() => track.closed.peek() !== undefined);

let aborted = false;
void streams[0].aborted.then(() => {
aborted = true;
});
if (parked) streams[0].release();
await drainUntil(() => aborted);
expect(aborted).toBe(true);
expect(streams.length).toBe(1);

// A GROUP for a forgotten id is ignored without touching its stream.
const touched: PropertyKey[] = [];
const reader = new Proxy({} as Reader, {
get: (_, key) => {
touched.push(key);
return () => {};
},
});
await subscriber.runGroup(new Group({ subscribe: 0n, sequence: 0 }), reader);
expect(touched).toEqual([]);
} finally {
subscriber.close();
warn.mockRestore();
jest.useRealTimers();
}
});

// The SUBSCRIBE stream can open after the deadline too; it is reset without carrying a SUBSCRIBE.
test("a lite subscribe that times out waiting on a stream slot for the SUBSCRIBE sends nothing on it", async () => {
jest.useFakeTimers();
const warn = spyOn(console, "warn").mockImplementation(() => {});
const { quic, streams } = fakeSession([1]);
const subscriber = new Subscriber(quic, Version.DRAFT_05, HopSchema.parse(1n));
try {
const track = subscriber.consume(Path.from("room")).track("video").subscribe();
await drainUntil(() => streams.length === 1);
await streams[0].reading;
// TRACK_INFO lands halfway, so the SUBSCRIBE open's own deadline is still ahead when the
// setup deadline fires.
jest.advanceTimersByTime(SUBSCRIBE_SETUP_TIMEOUT_MS / 2);
await answerTrackInfo(streams[0]);
await drainUntil(() => streams.length === 2);

jest.advanceTimersByTime(SUBSCRIBE_SETUP_TIMEOUT_MS / 2);
await drainUntil(() => track.closed.peek() !== undefined);

let aborted = false;
void streams[1].aborted.then(() => {
aborted = true;
});
streams[1].release();
await drainUntil(() => aborted);
expect(aborted).toBe(true);
expect(streams[1].written).toEqual([]);
} finally {
subscriber.close();
warn.mockRestore();
jest.useRealTimers();
}
});

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));
Expand Down
54 changes: 36 additions & 18 deletions js/net/src/lite/subscriber.ts
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,7 @@ import {
// drafts, or TRACK_INFO on lite-05+) may take. Browsers cap concurrent QUIC streams
// (Chrome ~100) and we open with waitUntilAvailable, so past the cap the open blocks
// until the peer frees a slot. The timeout turns a stall into a clear error.
const SUBSCRIBE_SETUP_TIMEOUT_MS = 10_000;
export const SUBSCRIBE_SETUP_TIMEOUT_MS = 10_000;

// The TRACK stream and implicit SUBSCRIBE acceptance are lite-05+.
function supportsTrackStream(version: Version): boolean {
Expand Down Expand Up @@ -540,8 +540,10 @@ export class Subscriber {
});

// Open the stream under a timeout. The stream handle flows back via `state`
// so the timeout path can abort it if it finishes opening after the deadline.
const state: { stream?: Stream } = {};
// so the timeout path can abort it if it finishes opening after the deadline,
// and `cancel` ends a setup still running once the deadline passed, resetting
// its TRACK stream so a peer that never answers can't hold one per attempt.
const state: { stream?: Stream; cancel: AbortController } = { cancel: new AbortController() };
const setup = this.#openSubscribe(state, msg, request, id, timescale);

let opened: { stream: Stream; entry: SubscribeEntry };
Expand All @@ -556,6 +558,7 @@ export class Subscriber {
// The setup outlived its deadline waiting for the first response: a control
// timeout, not content that arrived late.
const e = err instanceof TimeoutError ? controlTimeout(err) : await sessionCause(this.#quic, err);
state.cancel.abort(e);
request.reject(e);
this.#subscribes.delete(id);
console.warn(`subscribe error: id=${id} broadcast=${broadcast} track=${request.name} error=${reason(e)}`);
Expand Down Expand Up @@ -640,7 +643,7 @@ export class Subscriber {
// SUBSCRIBE is accepted implicitly (no SUBSCRIBE_OK). Older drafts carry no
// per-track properties, so they resolve to defaults and just drain SUBSCRIBE_OK.
async #openSubscribe(
state: { stream?: Stream },
state: { stream?: Stream; cancel: AbortController },
msg: Subscribe,
request: track.Request,
id: bigint,
Expand All @@ -651,7 +654,10 @@ export class Subscriber {

if (supportsTrackStream(this.version)) {
// Fetch the immutable properties once via the TRACK stream.
const info = await this.#trackInfo(msg.broadcast, msg.epoch, msg.track);
const info = await this.#trackInfo(msg.broadcast, msg.epoch, msg.track, state.cancel.signal);
// The deadline passed as TRACK_INFO landed: the request is already rejected, so don't
// register it again or send its SUBSCRIBE.
state.cancel.signal.throwIfAborted();
Comment thread
coderabbitai[bot] marked this conversation as resolved.
producer = request.accept(this.#toModelInfo(info));
timescale.set(info.timescale);
} else {
Expand Down Expand Up @@ -679,6 +685,9 @@ export class Subscriber {
this.#subscribes.set(id, entry);

state.stream = await Stream.open(this.#quic, { version: this.version });
// The deadline passed while the open waited: the late-setup handler resets the stream, so
// don't send a SUBSCRIBE on it first.
state.cancel.signal.throwIfAborted();
await state.stream.writer.u53(StreamId.Subscribe);
await msg.encode(state.stream.writer, this.version);

Expand All @@ -694,22 +703,31 @@ export class Subscriber {
}

// Opens a TRACK stream, reads the single TRACK_INFO, and FINs. Lite-05+ only.
async #trackInfo(broadcast: Path.Valid, epoch: Epoch.Valid | undefined, track: string): Promise<TrackInfo> {
return this.#exchange({ version: this.version }, async (stream) => {
await stream.writer.u53(StreamId.Track);
await new TrackMessage(broadcast, track, epoch).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;
});
async #trackInfo(
broadcast: Path.Valid,
epoch: Epoch.Valid | undefined,
track: string,
signal?: AbortSignal,
): Promise<TrackInfo> {
return this.#exchange(
{ version: this.version },
async (stream) => {
await stream.writer.u53(StreamId.Track);
await new TrackMessage(broadcast, track, epoch).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;
},
signal,
);
}

// 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<T>(options: OpenOptions, run: (stream: Stream) => Promise<T>): Promise<T> {
const closed = this.#closed.signal;
// fails. Subscriber.close() or `signal` also resets it while `run` is pending, so a peer that
// never answers cannot hold it open, and a stream that opens after either is reset at once.
async #exchange<T>(options: OpenOptions, run: (stream: Stream) => Promise<T>, signal?: AbortSignal): Promise<T> {
const closed = signal ? AbortSignal.any([this.#closed.signal, signal]) : this.#closed.signal;
closed.throwIfAborted();
const stream = await Stream.open(this.#quic, options);
const abort = () => stream.abort(error(closed.reason));
Expand Down
Loading