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
114 changes: 113 additions & 1 deletion js/net/src/lite/subscriber.test.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,15 @@
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";
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 () => {
Expand Down Expand Up @@ -691,3 +692,114 @@ test("a draft-05 duplicate start is still a restart", async () => {
announced.close();
subscriber.close();
});

interface FakeStream {
inbound: ReadableStreamDefaultController<Uint8Array>;
// Resolves once the subscriber waits on a read the test has not answered.
reading: Promise<void>;
aborted: Promise<unknown>;
// 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<Uint8Array>;
let onRead!: () => void;
let onAbort!: (reason: unknown) => void;
let release!: () => void;
const reading = new Promise<void>((resolve) => (onRead = resolve));
const aborted = new Promise<unknown>((resolve) => (onAbort = resolve));
// No high water mark, so pull() means the subscriber is blocked on a read.
const readable = new ReadableStream<Uint8Array>(
{
start: (controller) => {
inbound = controller;
},
pull: () => onRead(),
},
{ highWaterMark: 0 },
);
const writable = new WritableStream<Uint8Array>({ 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<void> {
const chunks: Uint8Array[] = [];
const writer = new Writer(
new WritableStream<Uint8Array>({ 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);
});
86 changes: 68 additions & 18 deletions js/net/src/lite/subscriber.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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<TrackInfo> {
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<T>(options: OpenOptions | undefined, run: (stream: Stream) => Promise<T>): Promise<T> {
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);
}
}

Expand Down Expand Up @@ -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<void> {
Expand Down Expand Up @@ -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);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
}

// 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<T>(group: netGroup.Producer, step: Promise<T>): Promise<T> {
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
Expand Down
Loading