diff --git a/doc/bin/relay/config.md b/doc/bin/relay/config.md index 27563226ae..c2588abc45 100644 --- a/doc/bin/relay/config.md +++ b/doc/bin/relay/config.md @@ -160,9 +160,11 @@ A draining upstream may name a replacement URI. `same-host` follows it only onto the host we already dialed, so a peer moves us between ports and schemes; `follow` also lets it choose the host, which means trusting it not to point us into the local network, since a name it controls resolves wherever it likes; -`ignore` keeps the current address list. Empty, malformed, or refused redirects -also preserve caller-configured fallbacks; only an accepted redirect replaces -the list with the peer's URI. `handover` is a cap: a shorter deadline on +`ignore` keeps the current address list. An empty URI also keeps it, including +caller-configured fallbacks; an accepted redirect replaces the list with the +peer's URI. A malformed or refused redirect ends the connection with an error +rather than redialing the old address or a fallback, and so does one leaving +the host a `tls.fingerprint` pin verifies. `handover` is a cap: a shorter deadline on the received GOAWAY wins, a longer one does not extend it. ## \[cache] diff --git a/doc/lib/js/net.md b/doc/lib/js/net.md index f086c11cd9..7abdc2e542 100644 --- a/doc/lib/js/net.md +++ b/doc/lib/js/net.md @@ -48,7 +48,7 @@ for (;;) { ``` - **Origins** hold the broadcasts, not the connection: closing a session unannounces them but leaves them created for the next one. `origin.request(path)` resolves an announced local broadcast with no round trip, so a page that watches what it publishes reads its own copy, unless a cheaper route announces the same path. Create, populate, then `announce()` for an exact path; use `dynamic(prefix, route)` when the set of paths is not known: an exact-path subscribe before the tracks exist is refused, and nobody, local or remote, can see or reach a broadcast until it announces. -- **Connections** race WebTransport against WebSocket. `new Connection({ url })` pools one connection per relay URL and reconnects with backoff, which the elements use. Supplying WebTransport/WebSocket options, discovery, delay, or a caller-owned origin selects a private loop; explicit `share: true` refuses those options. `closed` settles when the handle is released (`null` on a clean close); the failure that stopped retrying the current URL is `error`, and a new URL recovers the same handle. A connection owns one send-rate sampler and one `Bandwidth.Allocator`; publishers reserve against it so their encoder targets sum to the estimate instead of each matching it. +- **Connections** race WebTransport against WebSocket. `new Connection({ url })` pools one connection per relay URL and reconnects with backoff, which the elements use. Supplying WebTransport/WebSocket options, discovery, delay, `goaway`, or a caller-owned origin selects a private loop; explicit `share: true` refuses those options. On a relay's GOAWAY the connection dials the replacement at once while the old session keeps serving its groups until it closes or the `goaway.handover` cap (default 10s, lowered to the relay's deadline) passes; the handle and its origin stay the same. An empty GOAWAY redials the same URL through a fresh DNS resolve. A redirect is followed onto the same host by default (`goaway.redirect: "follow"` lets the relay pick the host, `"ignore"` never moves), and moves the pooled entry to the new URL; a malformed or refused one, including a host change under a pinned certificate, ends the connection with `Error.RefusedRedirect` instead of reconnecting. `closed` settles when the handle is released (`null` on a clean close); the failure that stopped retrying the current URL is `error`, and a new URL recovers the same handle. A connection owns one send-rate sampler and one `Bandwidth.Allocator`; publishers reserve against it so their encoder targets sum to the estimate instead of each matching it. - **Bandwidth** (`Bandwidth.Allocator`) divides the connection's send-rate estimate by track priority, max-min fair within a tier. An idle track claims nothing. The receive side is untouched. - **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. diff --git a/js/net/src/connection/forward.test.ts b/js/net/src/connection/forward.test.ts index d9d9b74481..6039d312a8 100644 --- a/js/net/src/connection/forward.test.ts +++ b/js/net/src/connection/forward.test.ts @@ -1,4 +1,5 @@ import { expect, test } from "bun:test"; +import { Once } from "@moq/signals"; import * as Announce from "../announced.ts"; import { type Consumer as BroadcastConsumer, Producer as BroadcastProducer } from "../broadcast.ts"; import { Route } from "../hop.ts"; @@ -33,7 +34,7 @@ class FakeSession { constructor(discovery = true) { this.discovery = discovery; - registerWire(this, { consume: (path) => this.consume(path) }); + registerWire(this, { consume: (path) => this.consume(path), goaway: new Once() }); this.closed = new Promise((resolve) => { this.#die = resolve; }); diff --git a/js/net/src/connection/goaway.test.ts b/js/net/src/connection/goaway.test.ts new file mode 100644 index 0000000000..fb99bbe1c2 --- /dev/null +++ b/js/net/src/connection/goaway.test.ts @@ -0,0 +1,110 @@ +import { expect, test } from "bun:test"; +import { RefusedRedirect } from "../error.ts"; +import * as Time from "../time.ts"; +import { dialed, handover, isLocal, pinnedTransport, type Redirect, target } from "./goaway.ts"; + +const current = new URL("https://relay.example/room?jwt=secret"); + +test("no URI, or a policy that ignores it, keeps the current URL", () => { + expect(target("follow", "", current, false)).toBeUndefined(); + expect(target("same-host", "", current, false)).toBeUndefined(); + expect(target("ignore", "https://other.example/", current, false)).toBeUndefined(); +}); + +test("an explicit URI the policy will not follow is refused, not ignored", () => { + const refused: [Redirect, string][] = [ + ["same-host", "https://other.example/"], + ["follow", "not a url"], + ["ignore", "not a url"], + ["follow", "http://relay.example/"], + ["follow", "unix:///tmp/moq.sock"], + ["follow", "https://127.0.0.1/"], + ["follow", "https://[::ffff:127.0.0.1]/"], + ["follow", "https://[::ffff:7f00:1]/"], + ["follow", "moqt://169.254.169.254/"], + ]; + for (const [policy, uri] of refused) { + expect(() => target(policy, uri, current, false), `${policy}: ${uri}`).toThrow(RefusedRedirect); + } +}); + +test("a refusal never repeats the URI, which can carry credentials", () => { + try { + target("same-host", "https://other.example/?jwt=leaked", current, false); + throw new Error("not refused"); + } catch (err) { + expect(err).toBeInstanceOf(RefusedRedirect); + expect(String(err)).not.toContain("leaked"); + } +}); + +test("same-host follows a port or scheme move on the host already dialed", () => { + expect(target("same-host", "https://relay.example:5443/", current, false)?.port).toBe("5443"); + // An explicit assignment is still one when it names the current URL. + expect(target("same-host", current.href, current, false)?.href).toBe(current.href); + // An upgrade is fine; only a downgrade is refused. + const plain = new URL("http://relay.example/"); + expect(target("same-host", "https://relay.example/", plain, false)?.protocol).toBe("https:"); +}); + +test("follow lets the peer name another public host", () => { + expect(target("follow", "https://other.example/next", current, false)?.href).toBe("https://other.example/next"); +}); + +test("a local endpoint may redirect to another local one", () => { + const local = new URL("https://localhost:4443/"); + expect(target("follow", "https://127.0.0.1:9999/", local, false)?.port).toBe("9999"); +}); + +test("a certificate pin holds the host even under follow", () => { + expect(() => target("follow", "https://other.example/", current, true)).toThrow(RefusedRedirect); + expect(target("follow", "https://relay.example:5443/", current, true)?.port).toBe("5443"); +}); + +test("local literals are recognized in every spelling", () => { + const local = [ + "https://127.0.0.1/", + "https://localhost/", + "https://a.localhost/", + "https://[::1]/", + "https://[::]/", + "https://10.0.0.1/", + "https://172.16.0.1/", + "https://192.168.1.1/", + "https://169.254.1.1/", + "https://0.0.0.0/", + "https://[::ffff:10.0.0.1]/", + "https://[fe80::1]/", + "https://[fc00::1]/", + "moqt://127.0.0.1/", + "moqt://[::ffff:127.0.0.1]/", + "unix:///tmp/moq.sock", + ]; + for (const url of local) expect(isLocal(new URL(url)), url).toBe(true); + + for (const url of ["https://example.com/", "https://8.8.8.8/", "https://172.32.0.1/", "https://[2606:4700::1]/"]) { + expect(isLocal(new URL(url)), url).toBe(false); + } +}); + +test("a WebSocket fallback that connected is the host a redirect is judged against", () => { + const primary = new URL("https://relay.example/"); + const socket = new URL("wss://edge.example/moq"); + expect(dialed(primary, "websocket", socket).href).toBe(socket.href); + expect(dialed(primary, "webtransport", socket).href).toBe(primary.href); + expect(dialed(primary, "websocket").href).toBe(primary.href); +}); + +test("a certificate pin holds only the WebTransport session that used it", () => { + expect(pinnedTransport("webtransport", true)).toBe(true); + expect(pinnedTransport("websocket", true)).toBe(false); + expect(pinnedTransport("webtransport", false)).toBe(false); +}); + +test("the handover is the cap, lowered only by a positive peer deadline", () => { + const cap = Time.Milli(10_000); + expect(handover(cap)).toBe(cap); + expect(handover(cap, Time.Milli(0))).toBe(cap); + expect(handover(cap, Time.Milli(3_000))).toBe(Time.Milli(3_000)); + expect(handover(cap, Time.Milli(3_600_000))).toBe(cap); +}); diff --git a/js/net/src/connection/goaway.ts b/js/net/src/connection/goaway.ts new file mode 100644 index 0000000000..a304e45838 --- /dev/null +++ b/js/net/src/connection/goaway.ts @@ -0,0 +1,224 @@ +/** + * GOAWAY handling for the reconnect loop: the peer's drain signal and the policy for the + * redirect it may name. Mirrors `Goaway` and `Redirect` in `moq_tokio::connection`. + * + * @module + */ +import { RefusedRedirect } from "../error.ts"; +import * as Time from "../time.ts"; + +/** A peer's GOAWAY as a session received it. */ +export interface Drain { + /** Where to reconnect, including any credentials it needs. Empty means the same endpoint. */ + readonly uri: string; + /** When the peer force-closes the session. Undefined when the wire carried none. */ + readonly timeout?: Time.Milli; +} + +/** + * What to do with the URI a peer names in its GOAWAY. + * + * - `same-host` (the default) follows it only onto the host already dialed, so a peer can move + * the connection between ports or schemes but not to another host. + * - `follow` also lets the peer name the host. The name is dialed as written, so only use it + * with a relay trusted not to point the connection into the local network. + * - `ignore` reconnects to the current URL whatever the peer names. + * + * A malformed or refused URI ends the connection with {@link RefusedRedirect}. + */ +export type Redirect = "follow" | "same-host" | "ignore"; + +/** How a reconnecting connection reacts to a peer's GOAWAY. */ +export interface GoawayProps { + /** What to do with the URI the peer names (default: `same-host`). */ + redirect?: Redirect; + + /** + * How long the old session keeps serving after the GOAWAY while the replacement dials + * (default: 10000ms). A cap: a shorter deadline from the peer wins, a longer one does not + * extend it. + */ + handover?: Time.Milli; +} + +/** The handover cap when the caller names none. */ +export const DEFAULT_HANDOVER = Time.Milli(10_000); + +/** + * How long a drained session may keep serving: the cap, lowered to the peer's deadline when + * it named a positive one. Absence is not a zero-length handover. + * + * @internal + */ +export function handover(cap: Time.Milli, timeout?: Time.Milli): Time.Milli { + return timeout !== undefined && timeout > 0 ? Time.Milli(Math.min(cap, timeout)) : cap; +} + +/** + * The URL a redirect is judged against. A WebSocket fallback that won the race is the host + * we dialed; the primary URL was not, so same-host must not treat it as the current peer. + * + * @internal + */ +export function dialed(primary: URL, transport: "webtransport" | "websocket", websocket?: URL): URL { + return transport === "websocket" && websocket ? websocket : primary; +} + +/** + * Whether a certificate pin constrains this session. The pin is a WebTransport option, so it + * holds the host only when that transport is the one that connected. + * + * @internal + */ +export function pinnedTransport(transport: "webtransport" | "websocket", configured: boolean): boolean { + return configured && transport === "webtransport"; +} + +/** + * The URL a GOAWAY assigns: `undefined` keeps the current URL (the peer named none, or the + * policy ignores a URI it could parse), and a URL replaces it. Throws {@link RefusedRedirect} + * for an explicit URI the policy will not follow, including a malformed one under `ignore`. + * + * `pinned` is a certificate pin on the connection, which can only verify the host it was + * configured for, so it refuses a host change even under `follow`. + * + * @internal + */ +export function target(policy: Redirect, uri: string, current: URL, pinned: boolean): URL | undefined { + if (uri === "") return undefined; + + // The URI can carry credentials, so the error names the reason, never the URI. + // Parse before `ignore`: a malformed redirect is terminal even when the policy + // would otherwise stay on the current URL. + let next: URL; + try { + next = new URL(uri); + } catch { + throw new RefusedRedirect("the GOAWAY URI is malformed"); + } + if (policy === "ignore") return undefined; + + if (schemeTier(next.protocol) < schemeTier(current.protocol)) { + throw new RefusedRedirect("the GOAWAY redirect downgrades the scheme"); + } + + // Only as far as the URL itself says: a name is dialed, never resolved here, so this + // refuses a peer naming a local address outright, not one hiding it behind a hostname. + // That gap is why `same-host` is the default. + if (isLocal(next) && !isLocal(current)) { + throw new RefusedRedirect("the GOAWAY redirect widens reachability to a local address"); + } + + // Host only, not the full authority: the port is what a peer legitimately moves us + // across when it hands off to a sibling process on the same box. + const sameHost = next.hostname === current.hostname; + if (policy === "same-host" && !sameHost) { + throw new RefusedRedirect("the GOAWAY redirect leaves the current host"); + } + if (pinned && !sameHost) { + throw new RefusedRedirect("the GOAWAY redirect leaves the host a certificate pin verifies"); + } + + return next; +} + +/** Rank a scheme so a redirect cannot silently drop encryption; unknown schemes rank lowest. */ +function schemeTier(protocol: string): number { + switch (protocol) { + case "https:": + case "wss:": + case "moqt:": + case "moql:": + return 2; + case "http:": + case "ws:": + case "tcp:": + return 1; + default: + return 0; + } +} + +/** + * Whether a URL says it names something only reachable from this host or network. A judgement + * about the URL, not about where a dial lands: `false` means "not local on its face". + * + * @internal + */ +export function isLocal(url: URL): boolean { + const host = url.hostname.toLowerCase(); + // No host at all, e.g. a `unix:` socket path. + if (host === "") return true; + if (host === "localhost" || host.endsWith(".localhost")) return true; + + if (host.startsWith("[") && host.endsWith("]")) { + const segments = parseIpv6(host.slice(1, -1)); + return segments !== undefined && isLocalV6(segments); + } + + const v4 = parseIpv4(host); + return v4 !== undefined && isLocalV4(v4); +} + +function parseIpv4(host: string): number[] | undefined { + const parts = host.split("."); + if (parts.length !== 4) return undefined; + const octets = parts.map((part) => (/^\d{1,3}$/.test(part) ? Number(part) : Number.NaN)); + return octets.every((octet) => octet <= 255) ? octets : undefined; +} + +function isLocalV4([a, b, c, d]: number[]): boolean { + return ( + a === 127 || + a === 10 || + (a === 172 && b >= 16 && b <= 31) || + (a === 192 && b === 168) || + (a === 169 && b === 254) || + (a === 0 && b === 0 && c === 0 && d === 0) + ); +} + +/** Parse an IPv6 literal (without brackets) into its eight 16-bit segments. */ +function parseIpv6(host: string): number[] | undefined { + // Drop a zone id; it scopes the address without changing it. + let text = host.split("%")[0] ?? ""; + + // A trailing dotted quad (`::ffff:127.0.0.1`) spells the last two segments. + const lastColon = text.lastIndexOf(":"); + const quad = parseIpv4(text.slice(lastColon + 1)); + if (quad) { + const [a = 0, b = 0, c = 0, d = 0] = quad; + const word = (hi: number, lo: number) => ((hi << 8) | lo).toString(16); + text = `${text.slice(0, lastColon + 1)}${word(a, b)}:${word(c, d)}`; + } + + const halves = text.split("::"); + if (halves.length > 2) return undefined; + const words = (part: string) => (part === "" ? [] : part.split(":")); + const head = words(halves[0] ?? ""); + const rest = halves.length === 2 ? words(halves[1] ?? "") : []; + const missing = 8 - head.length - rest.length; + if (halves.length === 2 ? missing < 0 : missing !== 0) return undefined; + + const all = [...head, ...Array(missing).fill("0"), ...rest]; + const segments = all.map((word) => (/^[0-9a-f]{1,4}$/i.test(word) ? Number.parseInt(word, 16) : Number.NaN)); + return segments.some(Number.isNaN) ? undefined : segments; +} + +function isLocalV6(segments: number[]): boolean { + // An IPv4-mapped address reaches the same host as the v4 it wraps. + if (segments.slice(0, 5).every((s) => s === 0) && segments[5] === 0xffff) { + const [hi = 0, lo = 0] = segments.slice(6); + return isLocalV4([hi >> 8, hi & 0xff, lo >> 8, lo & 0xff]); + } + const first = segments[0] ?? 0; + const zeroPrefix = segments.slice(0, 7).every((s) => s === 0); + const last = segments[7] ?? 0; + return ( + // Loopback (::1) and unspecified (::). + (zeroPrefix && (last === 1 || last === 0)) || + // Unique local (fc00::/7) and link local (fe80::/10). + (first & 0xfe00) === 0xfc00 || + (first & 0xffc0) === 0xfe80 + ); +} diff --git a/js/net/src/connection/migrate.test.ts b/js/net/src/connection/migrate.test.ts new file mode 100644 index 0000000000..a4af8d032f --- /dev/null +++ b/js/net/src/connection/migrate.test.ts @@ -0,0 +1,277 @@ +import { afterEach, expect, test } from "bun:test"; +import { RefusedRedirect } from "../error.ts"; +import * as Lite from "../lite/index.ts"; +import { createMockTransportPair } from "../mock.ts"; +import { Producer as OriginProducer } from "../origin.ts"; +import * as Path from "../path.ts"; +import { Stream } from "../stream.ts"; +import * as Time from "../time.ts"; +import { accept } from "./accept.ts"; +import { Connection, resetShared } from "./pool.ts"; +import { Reload } from "./reload.ts"; + +const original = globalThis.WebTransport; + +afterEach(() => { + resetShared(); + globalThis.WebTransport = original; +}); + +async function settle() { + await new Promise((resolve) => setTimeout(resolve, 0)); +} + +// Polls until `pred` holds, so a regression fails the test instead of hanging it. +async function waitUntil(pred: () => boolean, ms = 2000): Promise { + const deadline = Date.now() + ms; + for (;;) { + if (pred()) return; + if (Date.now() > deadline) throw new Error("timed out waiting for condition"); + await settle(); + } +} + +/** One accepted session on the fake fleet: the URL it was dialed at, and its server side. */ +interface Dial { + url: string; + server: WebTransport; + closed: boolean; +} + +/** + * Stand in for a fleet of relays behind any URL: every dial gets a fresh session serving + * `origin`, so each one publishes the same broadcasts the way siblings of a fleet do. + */ +function fleet(origin = new OriginProducer()): { dials: Dial[]; origin: OriginProducer } { + const dials: Dial[] = []; + const stub = function StubWebTransport(url: string | URL) { + const pair = createMockTransportPair(Lite.ALPN_06); + const dial: Dial = { url: new URL(url).href, server: pair.server, closed: false }; + dials.push(dial); + void pair.server.closed.then( + () => { + dial.closed = true; + }, + () => { + dial.closed = true; + }, + ); + void accept({ transport: pair.server, url: new URL(url), publish: origin.consume() }); + return pair.client; + }; + globalThis.WebTransport = stub as unknown as typeof WebTransport; + return { dials, origin }; +} + +/** Send a moq-lite GOAWAY from the server side of a session. */ +async function goaway(server: WebTransport, uri: string): Promise { + const stream = await Stream.open(server); + await stream.writer.u53(Lite.StreamId.Goaway); + await new Lite.Goaway(uri).encode(stream.writer, Lite.Version.DRAFT_06); +} + +const url = new URL("https://relay.example/room"); + +// Short enough to watch, long enough that a handover visibly overlaps the replacement. +const handover = Time.Milli(200); + +test("an empty GOAWAY migrates without unrouting the path", async () => { + const { dials, origin } = fleet(); + const broadcast = origin.createBroadcast(Path.from("cam")); + broadcast.announce(); + + const consume = new OriginProducer(); + const reload = new Reload({ + url, + websocket: { enabled: false }, + consume, + goaway: { handover }, + // A session younger than this counts as redirected immediately and backs off first. + delay: { initial: Time.Milli(1) }, + }); + const watched = consume.request(Path.from("cam"), { announced: true }); + + try { + await waitUntil(() => watched.active.peek() !== undefined); + const first = reload.established.peek(); + + // Record every moment the path had no route, from here on. + let gaps = 0; + const stop = watched.active.subscribe((active) => { + if (active === undefined) gaps++; + }); + + const drained = dials[0]; + if (!drained) throw new Error("no first dial"); + await goaway(drained.server, ""); + + // The replacement dials the configured URL at once, while the old session still serves. + await waitUntil(() => dials.length === 2); + expect(dials[1]?.url).toBe(url.href); + await waitUntil(() => reload.established.peek() !== first && reload.established.peek() !== undefined); + expect(drained.closed).toBe(false); + expect(reload.status.peek()).toBe("connected"); + + // The old session closes at the handover cap, and the path never went unrouted. + await waitUntil(() => drained.closed, handover * 10); + await settle(); + expect(watched.active.peek()).not.toBeUndefined(); + expect(gaps).toBe(0); + expect(dials.length).toBe(2); + stop(); + } finally { + watched.close(); + reload.close(); + consume.close(); + broadcast.close(); + } +}); + +test("a GOAWAY without a timeout hands over at the configured cap", async () => { + const { dials } = fleet(); + const reload = new Reload({ + url, + websocket: { enabled: false }, + goaway: { handover }, + delay: { initial: Time.Milli(1) }, + }); + + try { + await waitUntil(() => reload.status.peek() === "connected"); + const drained = dials[0]; + if (!drained) throw new Error("no first dial"); + + const sent = performance.now(); + await goaway(drained.server, ""); + await waitUntil(() => dials.length === 2 && reload.status.peek() === "connected"); + + // Lite carries no deadline, which must read as "the cap", never as a zero handover. + await waitUntil(() => drained.closed, handover * 10); + expect(performance.now() - sent).toBeGreaterThanOrEqual(handover * 0.9); + } finally { + reload.close(); + } +}); + +test("a refused redirect ends the connection instead of reconnecting", async () => { + const refused = ["https://other.example/", "not a url", "http://relay.example/", "https://127.0.0.1/"]; + for (const uri of refused) { + const { dials } = fleet(); + const reload = new Reload({ url, websocket: { enabled: false }, delay: { initial: Time.Milli(1) } }); + + try { + await waitUntil(() => reload.status.peek() === "connected"); + const drained = dials[0]; + if (!drained) throw new Error("no first dial"); + + await goaway(drained.server, uri); + await waitUntil(() => reload.error.peek() !== undefined); + expect(reload.error.peek(), uri).toBeInstanceOf(RefusedRedirect); + expect(reload.status.peek()).toBe("disconnected"); + await waitUntil(() => drained.closed); + + // Nothing redials: not the original URL, not anything else. + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(dials.length, uri).toBe(1); + } finally { + reload.close(); + } + } +}); + +test("a certificate pin refuses a redirect to another host even under follow", async () => { + const { dials } = fleet(); + const reload = new Reload({ + url, + websocket: { enabled: false }, + webtransport: { serverCertificateHashes: [{ value: "00".repeat(32) }] }, + goaway: { redirect: "follow" }, + delay: { initial: Time.Milli(1) }, + }); + + try { + await waitUntil(() => reload.status.peek() === "connected"); + await goaway(dials[0]?.server as WebTransport, "https://other.example/"); + await waitUntil(() => reload.error.peek() !== undefined); + expect(reload.error.peek()).toBeInstanceOf(RefusedRedirect); + expect(dials.length).toBe(1); + } finally { + reload.close(); + } +}); + +// A same-host move to another port: what `same-host` exists to allow. +const moved = new URL("https://relay.example:5443/room"); + +test("a redirect moves the pool key while the handle keeps its origin", async () => { + const { dials } = fleet(); + + const handle = new Connection({ url }); + try { + await waitUntil(() => handle.status.peek() === "connected"); + const origin = handle.origin.peek(); + + await goaway(dials[0]?.server as WebTransport, moved.href); + await waitUntil(() => dials.length === 2); + expect(dials[1]?.url).toBe(moved.href); + expect(handle.origin.peek()).toBe(origin); + + // A caller configured with the target shares the migrated connection. + const joined = new Connection({ url: moved }); + await waitUntil(() => joined.status.peek() === "connected"); + expect(joined.origin.peek()).toBe(origin); + expect(dials.length).toBe(2); + + // One still asking for the original URL gets a fresh entry. + const fresh = new Connection({ url }); + await waitUntil(() => fresh.status.peek() === "connected"); + expect(fresh.origin.peek()).not.toBe(origin); + expect(dials.length).toBe(3); + + joined.close(); + fresh.close(); + } finally { + handle.close(); + } +}); + +test("a redirect onto an already pooled key leaves both entries to their handles", async () => { + const { dials } = fleet(); + + const redirected = new Connection({ url }); + const resident = new Connection({ url: moved }); + try { + await waitUntil(() => redirected.status.peek() === "connected" && resident.status.peek() === "connected"); + const migrated = redirected.origin.peek(); + const target = resident.origin.peek(); + expect(migrated).not.toBe(target); + + const drained = dials.find((dial) => dial.url === url.href); + await goaway(drained?.server as WebTransport, moved.href); + await waitUntil(() => dials.filter((dial) => dial.url === moved.href).length === 2); + await waitUntil(() => redirected.status.peek() === "connected"); + + // Existing handles keep their own entries. + expect(redirected.origin.peek()).toBe(migrated); + expect(resident.origin.peek()).toBe(target); + + // The target key still belongs to the entry that was there. + const joined = new Connection({ url: moved }); + await waitUntil(() => joined.origin.peek() !== undefined); + expect(joined.origin.peek()).toBe(target); + + // The original key was vacated, so a new caller dials fresh. + const before = dials.length; + const fresh = new Connection({ url }); + await waitUntil(() => fresh.status.peek() === "connected"); + expect(fresh.origin.peek()).not.toBe(migrated); + expect(fresh.origin.peek()).not.toBe(target); + expect(dials.length).toBe(before + 1); + + joined.close(); + fresh.close(); + } finally { + redirected.close(); + resident.close(); + } +}); diff --git a/js/net/src/connection/pool.ts b/js/net/src/connection/pool.ts index 959697ca09..b7da0472bf 100644 --- a/js/net/src/connection/pool.ts +++ b/js/net/src/connection/pool.ts @@ -20,6 +20,7 @@ import { type WebTransportProps as WebTransportPropsType, } from "./connect.ts"; import type { Established as EstablishedType } from "./established.ts"; +import type { GoawayProps, Redirect as RedirectType } from "./goaway.ts"; import { Reload, type ReloadDelay, type ReloadStatus } from "./reload.ts"; import type { Probe as ProbeType, Stats as StatsType } from "./stats.ts"; import type { Transport as TransportType } from "./transport.ts"; @@ -69,6 +70,13 @@ export interface ConnectionProps { /** Backoff settings for the reconnect loop; an unset field uses its default. */ delay?: ReloadDelay; + /** + * How to react to the relay's GOAWAY; an unset field uses its default. Every connection + * migrates on GOAWAY, dialing the replacement while the old session finishes its groups. + * This only tunes where it may go and for how long the old session serves. + */ + goaway?: GoawayProps; + /** A Connection owns the abort signal for each connection attempt. */ signal?: never; @@ -93,7 +101,11 @@ export interface ConnectionProps { * connection so the next handle dials fresh; a new URL on this handle starts another * sequence. {@link closed} settles only when this handle is released. * - * Options the pool cannot honor (transport options, discovery, delay, a pinned + * A relay's GOAWAY migrates rather than drops: the replacement dials at once while the old + * session finishes its groups, and the handle and its origin carry across. A redirect the + * {@link ConnectionProps.goaway} policy refuses stops the loop like an auth rejection does. + * + * Options the pool cannot honor (transport options, discovery, delay, goaway, a pinned * certificate, caller-owned origins) use a private loop. An explicit `share: true` * refuses them. A supplied transport cannot reconnect at all; pass it to {@link Connection.connect} * instead. @@ -246,6 +258,7 @@ export class Connection { // A handle nobody watches wants unlimited retries; an auth rejection still // stops this URL, and a new one starts another sequence. delay: { timeout: Time.Milli(0), ...props.delay }, + goaway: props.goaway, }); this.#signals.cleanup(() => loop.close()); @@ -339,6 +352,10 @@ export namespace Connection { export type AcceptProps = AcceptPropsType; /** Backoff settings for a private reconnect loop. */ export type Backoff = ReloadDelay; + /** How a connection reacts to the relay's GOAWAY. */ + export type Goaway = GoawayProps; + /** What to do with the URI a relay names in its GOAWAY. */ + export type Redirect = RedirectType; /** Current state of a {@link Connection}. */ export type Status = ReloadStatus; /** The current connection's PROBE estimates. */ @@ -389,6 +406,9 @@ function refuse(props?: ConnectionProps): void { if (props.delay !== undefined) { throw new Error("delay cannot be shared; pass share: false"); } + if (props.goaway !== undefined) { + throw new Error("goaway cannot be shared; pass share: false"); + } } /** Options tied to one handle cannot be represented by a URL-keyed shared entry. */ @@ -398,6 +418,7 @@ function requiresPrivate(props: ConnectionProps): boolean { props.websocket !== undefined || props.discovery !== undefined || props.delay !== undefined || + props.goaway !== undefined || props.publish !== undefined || props.consume !== undefined ); @@ -405,6 +426,8 @@ function requiresPrivate(props: ConnectionProps): boolean { /** One shared connection and the handles keeping it alive. */ interface Entry { + /** The pool key: the dialed URL, which a GOAWAY redirect moves. */ + key: string; origin: Origin.Producer; connection: Reload; refs: number; @@ -431,14 +454,25 @@ function acquire(key: string, linger?: Time.Milli): Entry & { release: () => voi delay: { timeout: Time.Milli(0) }, }); - const created: Entry = { origin, connection, refs: 0, linger: linger ?? LINGER_MS }; + const created: Entry = { key, origin, connection, refs: 0, linger: linger ?? LINGER_MS }; - // The loop only stops on a peer saying these credentials will never work. Drop the - // entry so a later handle dials fresh rather than joining a loop that has stopped; - // handles already on it keep it until they release, since a redial would be refused - // the same way. + // The loop only stops on a peer saying these credentials will never work, or on a + // GOAWAY redirect it refused. Drop the entry so a later handle dials fresh rather + // than joining a loop that has stopped; handles already on it keep it until they + // release. connection.error.subscribe((err) => { - if (err !== undefined && pool.get(key) === created) pool.delete(key); + if (err !== undefined && pool.get(created.key) === created) pool.delete(created.key); + }); + + // An accepted redirect moves the entry to the URL it now dials, so a later handle + // configured with that URL shares it, and one asking for the old URL dials fresh. + // A live entry already at the target wins: this one leaves the pool and serves only + // the handles it has until they release it. + connection.redirect.subscribe((redirect) => { + if (!redirect || redirect.href === created.key) return; + if (pool.get(created.key) === created) pool.delete(created.key); + created.key = redirect.href; + if (!pool.has(created.key)) pool.set(created.key, created); }); entry = created; @@ -464,7 +498,7 @@ function acquire(key: string, linger?: Time.Milli): Entry & { release: () => voi if (taken.refs > 0) return; taken.timer = setTimeout(() => { - if (pool.get(key) === taken) pool.delete(key); + if (pool.get(taken.key) === taken) pool.delete(taken.key); taken.connection.close(); taken.origin.close(); }, taken.linger); diff --git a/js/net/src/connection/reload.ts b/js/net/src/connection/reload.ts index b7b77bcaef..a555ff9bf3 100644 --- a/js/net/src/connection/reload.ts +++ b/js/net/src/connection/reload.ts @@ -8,6 +8,7 @@ import * as Time from "../time.ts"; import { wireOf } from "../wire.ts"; import { type ConnectProps, connect, type WebSocketProps, type WebTransportProps } from "./connect.ts"; import type { Established } from "./established.ts"; +import { DEFAULT_HANDOVER, type Drain, dialed, type GoawayProps, handover, pinnedTransport, target } from "./goaway.ts"; import type { Probe, Stats } from "./stats.ts"; /** @@ -59,6 +60,9 @@ export type ReloadProps = Omit & { /** Backoff settings for the reconnect loop; every field falls back to its default. */ delay?: ReloadDelay; + + /** How to react to the peer's GOAWAY; every field falls back to its default. */ + goaway?: GoawayProps; }; /** @@ -162,6 +166,16 @@ export class Reload { /** Backoff settings for the reconnect loop; an unset field uses its default. */ delay: ReloadDelay; + /** How to react to the peer's GOAWAY (not reactive). */ + goaway: GoawayProps; + + /** + * The URL an accepted GOAWAY redirect assigned, or undefined while the loop dials + * {@link Reload.url}. Sticky across reconnects: a redirect is an assignment, not a + * detour. Cleared when a new URL or a disable/re-enable starts another sequence. + */ + readonly redirect: Getter; + /** The reactive effect scope driving the connect loop; closed by {@link Reload.close}. */ #signals = new Effect(); @@ -179,6 +193,16 @@ export class Reload { #closed = new Once(); #error = new Signal(undefined); + #redirect = new Signal(undefined); + // The configured href the redirect was assigned for. + #redirectHref: string | undefined; + + // A session the peer sent GOAWAY on, serving its groups in flight while the replacement + // dials. Retired when it closes, at its handover cap, or when the sequence ends. + #draining: Draining | undefined; + + // Whether a connect attempt is in flight, so a retiring predecessor reports the right status. + #dialing = false; // The current wait between attempts, doubling per failure, and when the retry window expires. // Both are undefined between sequences, so a later edit to `delay` applies to the next one. @@ -207,6 +231,8 @@ export class Reload { this.url = Signal.from(props?.url); this.enabled = Signal.from(props?.enabled ?? true); this.delay = props?.delay ?? {}; + this.goaway = props?.goaway ?? {}; + this.redirect = this.#redirect; this.webtransport = props?.webtransport; this.websocket = props?.websocket; this.discovery = props?.discovery; @@ -279,6 +305,7 @@ export class Reload { if (!enabled) { this.#givenUpHref = undefined; this.#resetSequence(); + this.#redirect.set(undefined); return; } @@ -293,6 +320,7 @@ export class Reload { if (!href) { this.#givenUpHref = undefined; this.#resetSequence(); + this.#redirect.set(undefined); return; } const url = new URL(href); @@ -303,49 +331,100 @@ export class Reload { if (this.#sequenceHref !== href) { this.#resetSequence(); + // A redirect was assigned for another URL; this one starts from itself. A page + // hide/show resumes the same URL, so it keeps the assignment. + if (this.#redirectHref !== href) this.#redirect.set(undefined); this.#sequenceHref = href; this.#error.set(undefined); this.#givenUpHref = undefined; } - effect.set(this.status, "connecting", "disconnected"); + // A drained predecessor still serves while its replacement dials. + if (!this.#draining) this.status.set("connecting"); // This run's teardown, handed to connect() so a rerun cancels the attempt in flight. const signal = effect.abort; + // The session this run serves, closed with the run. A drained one leaves this slot + // for #draining, which outlives the run. + let current: Established | undefined; + effect.cleanup(() => { + if (current) { + current.close(); + if (this.established.peek() === current) this.established.set(undefined); + current = undefined; + } + if (this.established.peek() === undefined) this.status.set("disconnected"); + }); + effect.spawn(async () => { // Set once the session is live, so #retry can tell a healthy session that // later dropped from a connect failure or a peer that flaps immediately. let connected: DOMHighResTimeStamp | undefined; try { - const connection = await connect({ - url, - websocket: this.websocket, - webtransport: this.webtransport, - discovery: this.discovery, - publish: this.publish, - consume: this.consume, - signal, - }); - - // Hand the connection to the effect, which closes it now if this run is already over. - effect.cleanup(() => connection.close()); - if (signal.aborted) return; - - effect.set(this.established, connection); - effect.set(this.status, "connected", "disconnected"); - - connected = performance.now(); + // Loops only to migrate: a GOAWAY off a healthy session dials its replacement + // straight away, with no backoff. + for (;;) { + const dialing = this.#redirect.peek() ?? url; + + this.#dialing = true; + let connection: Established; + try { + connection = await connect({ + url: dialing, + // A redirect names the relay; a fallback URL pinned for the old one does not follow. + websocket: this.#redirect.peek() ? { ...this.websocket, url: undefined } : this.websocket, + webtransport: this.webtransport, + discovery: this.discovery, + publish: this.publish, + consume: this.consume, + signal, + }); + } finally { + this.#dialing = false; + } - // A cancelled effect resolves undefined, so the sentinel tells the session - // closing (null for clean, an Error otherwise) apart from this run being - // torn down. - const closed = await effect.race(connection.closed); - if (closed === undefined) return; + // Hand the connection to the effect, which closes it now if this run is already over. + if (signal.aborted) { + connection.close(); + return; + } + current = connection; + + // The replacement serves now; a predecessor keeps draining its groups in flight. + this.established.set(connection); + this.status.set("connected"); + connected = performance.now(); + + // A cancelled effect resolves undefined, so the sentinel tells the session + // closing (null for clean, an Error otherwise) apart from this run being + // torn down. Anything else is the peer's GOAWAY. + const ended = await effect.race(connection.closed, wireOf(connection).goaway); + if (ended === undefined) return; + if (ended === null || ended instanceof Error) { + console.warn("connection closed, reconnecting"); + if (this.established.peek() === connection) this.established.set(undefined); + this.#retry(effect, connected, ended ?? undefined); + return; + } - console.warn("connection closed, reconnecting"); - this.#retry(effect, connected, closed ?? undefined); + current = undefined; + // A pinned WebSocket URL can win the race against the primary. Judge the + // redirect against that endpoint, not the primary we never reached. + const socket = this.#redirect.peek() ? undefined : this.websocket?.url; + if (!this.#migrate(connection, dialed(dialing, connection.transport, socket), ended)) return; + + // A session that outlived the initial delay was healthy, so its handover is not a + // failure. One redirected almost at once still migrates, but through the backoff, + // so two peers bouncing us between them escalate and eventually give up. + if (performance.now() - connected < this.#initial()) { + this.#retry(effect, connected, new Error("peer redirected immediately")); + return; + } + this.#delay = undefined; + this.#deadline = undefined; + } } catch (err) { // Treat teardown as cancellation, not a connection failure. if (signal.aborted) return; @@ -356,6 +435,56 @@ export class Reload { }); } + /** + * Act on the peer's GOAWAY for `connection`, which was dialed at `dialing`: resolve where + * to go next and leave the old session serving until it drains. Returns false when the + * redirect is refused, which ends the sequence rather than redialing. + */ + #migrate(connection: Established, dialing: URL, drain: Drain): boolean { + const hashes = this.webtransport?.serverCertificateHashes?.length ?? 0; + const configured = hashes > 0 || this.webtransport?.serverCertificate !== undefined; + // The pin is a WebTransport option. A WebSocket that won the race never used it. + const pinned = pinnedTransport(connection.transport, configured); + + let next: URL | undefined; + try { + next = target(this.goaway.redirect ?? "same-host", drain.uri, dialing, pinned); + } catch (err) { + // The peer is leaving and named somewhere we won't go: redialing the old address + // would ignore it, so stop here. + console.warn("GOAWAY redirect refused:", err); + connection.close(); + this.established.set(undefined); + this.status.set("disconnected"); + this.#giveUp(error(err)); + return false; + } + + // Only an accepted redirect replaces the URL; an empty one keeps it. + if (next) { + this.#redirect.set(next); + this.#redirectHref = this.#sequenceHref; + } + + console.info("GOAWAY received; migrating"); + // A newer GOAWAY retires an older predecessor rather than holding two open. + this.#draining?.retire(); + const cap = handover(this.goaway.handover ?? DEFAULT_HANDOVER, drain.timeout); + const draining = new Draining(connection, cap, () => { + if (this.#draining === draining) this.#draining = undefined; + // If nothing replaced it yet, nothing is serving. + if (this.established.peek() !== connection) return; + this.established.set(undefined); + this.status.set(this.#dialing ? "connecting" : "disconnected"); + }); + this.#draining = draining; + return true; + } + + #initial(): Time.Milli { + return this.delay?.initial ?? DEFAULT_DELAY.initial; + } + /** * Schedule the next connect attempt after the current backoff, or stop once the retry window * has expired. `connected` is when the dead session was established, if it ever was, and @@ -368,15 +497,14 @@ export class Reload { // optional value passes an explicit undefined, which a spread would take as the // answer, turning the backoff into NaN or the window into forever. const delay = this.delay ?? {}; - const initial = delay.initial ?? DEFAULT_DELAY.initial; + const initial = this.#initial(); const multiplier = delay.multiplier ?? DEFAULT_DELAY.multiplier; const max = delay.max ?? DEFAULT_DELAY.max; const timeout = delay.timeout ?? DEFAULT_DELAY.timeout; - // Any session is dead now: report disconnected during the backoff rather than - // when the retry reruns the effect. - this.established.set(undefined); - this.status.set("disconnected"); + // Report disconnected during the backoff rather than when the retry reruns the + // effect, unless a drained predecessor still serves until it retires. + if (this.established.peek() === undefined) this.status.set("disconnected"); // A session that outlived the initial delay was healthy, so clear the backoff and // start a fresh retry window: a one-off drop should reconnect promptly. Anything @@ -432,6 +560,7 @@ export class Reload { this.#delay = undefined; this.#deadline = undefined; this.#sequenceHref = undefined; + this.#draining?.retire(); } /** @@ -507,6 +636,37 @@ export class Reload { /** Stop reconnecting, close the current connection, and settle {@link Reload.closed}. Idempotent. */ close(abort?: Error) { this.#signals.close(); + this.#draining?.retire(); if (this.#closed.peek() === undefined) this.#closed.set(abort ?? null); } } + +/** + * A session the peer sent GOAWAY on, left serving so its groups in flight finish. It retires + * when it closes on its own or overstays its handover window, and `onRetire` runs once either way. + */ +class Draining { + #connection: Established; + #timer: ReturnType; + #onRetire: () => void; + #retired = false; + + constructor(connection: Established, handover: Time.Milli, onRetire: () => void) { + this.#connection = connection; + this.#onRetire = onRetire; + this.#timer = setTimeout(() => { + console.warn("old session did not drain in time; closing"); + this.retire(); + }, handover); + void connection.closed.then(() => this.retire()); + } + + /** Close the old session now, whatever remains of its window. Idempotent. */ + retire(): void { + if (this.#retired) return; + this.#retired = true; + clearTimeout(this.#timer); + this.#connection.close(); + this.#onRetire(); + } +} diff --git a/js/net/src/error.ts b/js/net/src/error.ts index 912fae64cd..dc048d19fe 100644 --- a/js/net/src/error.ts +++ b/js/net/src/error.ts @@ -246,6 +246,21 @@ export class NotFound extends Stream { } } +/** + * A peer's GOAWAY named a redirect the connection refuses, or one it could not parse. + * + * Terminal: the peer is leaving, so the connection stops rather than redialing the old + * address. Mirrors the Rust `Error::RefusedRedirect`. + * + * @public + */ +export class RefusedRedirect extends Error { + constructor(reason: string) { + super(`GOAWAY redirect refused: ${reason}`); + this.name = "RefusedRedirect"; + } +} + /** * A peer broke the protocol in a way the spec says must end the session. * diff --git a/js/net/src/errors.ts b/js/net/src/errors.ts index 09cc4dd606..d84bf02d1e 100644 --- a/js/net/src/errors.ts +++ b/js/net/src/errors.ts @@ -8,6 +8,7 @@ export { GroupTooLarge, NotFound, ProtocolViolation, + RefusedRedirect, Session, Stream, type StreamOptions, diff --git a/js/net/src/ietf/adapter.test.ts b/js/net/src/ietf/adapter.test.ts index 1945cef825..ebde36a91f 100644 --- a/js/net/src/ietf/adapter.test.ts +++ b/js/net/src/ietf/adapter.test.ts @@ -4,6 +4,7 @@ import * as Path from "../path.ts"; import { Stream } from "../stream.ts"; import { ControlStreamAdapter } from "./adapter.ts"; import { toRequestCode } from "./error.ts"; +import { GoAway } from "./goaway.ts"; import { PublishNamespace, PublishNamespaceCancel, PublishNamespaceDone } from "./publish_namespace.ts"; import { RequestError } from "./request.ts"; import { ALPN, Version } from "./version.ts"; @@ -184,3 +185,51 @@ test("withdrawals distinguish the same namespace by direction", async () => { await cancel(peer, namespace); expect(await closed(outgoing)).toBe(true); }); + +/** + * Draft-14 to -16 carry GOAWAY on the shared control stream. The adapter decodes it and keeps + * routing, so the session serves its groups in flight while the caller migrates. + */ +test("the control stream adapter decodes a GOAWAY and keeps running", async () => { + const { adapter, peer } = await connect(); + + await peer.writer.u53(GoAway.id); + await new GoAway({ newSessionUri: "https://relay.example/next" }).encode(peer.writer, VERSION); + + const drain = await adapter.goaway; + expect(drain.uri).toBe("https://relay.example/next"); + // These drafts carry no timeout: absence means the caller's cap, never a zero handover. + expect(drain.timeout).toBeUndefined(); + + // Still routing: a later announcement opens its virtual stream. + await announce(peer, 1n, Path.from("still")); + await accept(adapter); +}); + +test("a server adapter rejects a client GOAWAY that names a redirect", async () => { + const pair = createMockTransportPair(ALPN.DRAFT_15); + const control = await Stream.open(pair.server, { version: VERSION }); + const adapter = new ControlStreamAdapter(pair.server, control, VERSION, 100n, false); + const running = adapter.run(); + const peer = await Stream.accept(pair.client, VERSION); + if (!peer) throw new Error("no control stream"); + + await peer.writer.u53(GoAway.id); + await new GoAway({ newSessionUri: "https://other.example/" }).encode(peer.writer, VERSION); + await expect(running).rejects.toThrow("client GOAWAY must not name a redirect"); +}); + +test("a second GOAWAY on the control stream closes the session", async () => { + const pair = createMockTransportPair(ALPN.DRAFT_15); + const control = await Stream.open(pair.server, { version: VERSION }); + const adapter = new ControlStreamAdapter(pair.server, control, VERSION, 100n, true); + const running = adapter.run(); + const peer = await Stream.accept(pair.client, VERSION); + if (!peer) throw new Error("no control stream"); + + for (let i = 0; i < 2; i++) { + await peer.writer.u53(GoAway.id); + await new GoAway({ newSessionUri: "" }).encode(peer.writer, VERSION); + } + await expect(running).rejects.toThrow("duplicate GOAWAY"); +}); diff --git a/js/net/src/ietf/adapter.ts b/js/net/src/ietf/adapter.ts index fd16a5349d..ad2ae94e2e 100644 --- a/js/net/src/ietf/adapter.ts +++ b/js/net/src/ietf/adapter.ts @@ -1,6 +1,10 @@ +import { Once } from "@moq/signals"; import { Mutex } from "async-mutex"; +import type { Drain } from "../connection/goaway.ts"; +import { ProtocolViolation } from "../error.ts"; import { Reader, Stream, type Writer } from "../stream.ts"; import * as Varint from "../varint.ts"; +import { GoAway } from "./goaway.ts"; import * as Namespace from "./namespace.ts"; import { type IetfVersion, Version } from "./version.ts"; @@ -90,6 +94,9 @@ export class ControlStreamAdapter implements Session { #writeMutex = new Mutex(); readonly version: IetfVersion; + /** The peer's GOAWAY, which draft-14 to -16 carry on the shared control stream. */ + readonly goaway = new Once(); + // Virtual streams keyed by requestId #streams = new Map(); @@ -123,6 +130,9 @@ export class ControlStreamAdapter implements Session { #closed = false; + // Whether this side opened the session. Only a server may name a redirect. + #client: boolean; + constructor( quic: WebTransport, controlStream: Stream, @@ -138,6 +148,7 @@ export class ControlStreamAdapter implements Session { this.version = version; this.#maxRequestId = maxRequestId; this.#requestId = client ? 0n : 1n; + this.#client = client; } /** @@ -263,8 +274,15 @@ export class ControlStreamAdapter implements Session { const classified = await this.#classify(typeId, body); if (classified.route === Route.GoAway) { - console.warn("received GOAWAY on control stream"); - return; + // The session keeps serving: a GOAWAY asks us to migrate, not to stop reading. + const msg = await GoAway.decodeBody(body, this.version); + if (this.goaway.peek() !== undefined) throw new ProtocolViolation("duplicate GOAWAY"); + // A client may leave, but only the server may name where to go. + if (!this.#client && msg.newSessionUri !== "") { + throw new ProtocolViolation("client GOAWAY must not name a redirect"); + } + this.goaway.set(msg.drain()); + continue; } const { route, requestId } = classified; diff --git a/js/net/src/ietf/connection.ts b/js/net/src/ietf/connection.ts index e1d22f9bf8..7869b48188 100644 --- a/js/net/src/ietf/connection.ts +++ b/js/net/src/ietf/connection.ts @@ -1,6 +1,7 @@ -import { type Getter, Signal } from "@moq/signals"; +import { type Getter, Once, Signal } from "@moq/signals"; import type * as announce from "../announced.ts"; import type { Established } from "../connection/established.ts"; +import type { Drain } from "../connection/goaway.ts"; import { type Probe, type Stats, transportStats } from "../connection/stats.ts"; import { type Transport, transportOf } from "../connection/transport.ts"; import { error, fromClose, ProtocolViolation, StreamCode, StreamError } from "../error.ts"; @@ -45,6 +46,9 @@ export class Connection implements Established { // The established WebTransport session. #quic: WebTransport; + // Whether this side opened the session. Only a server may name a redirect. + #client: boolean; + // Session abstraction: adapter for v14-v16, native for v17. #session: Session; @@ -60,6 +64,9 @@ export class Connection implements Established { // The Hop IDs this session declared; see {@link Cluster}. #cluster?: Cluster.Hops; + // The peer's GOAWAY: read here on v17+, by the control stream adapter before that. + #goaway: Once; + // Just to avoid logging when `close()` is called. #closed = false; @@ -116,15 +123,18 @@ export class Connection implements Established { this.version = versionName(version); this.transport = transportOf(quic); this.#quic = quic; + this.#client = client; // Two-path dispatch: v14-v16 uses adapter, v17+ uses native bidi streams if (version >= Version.DRAFT_17) { this.#session = new NativeSession(quic, version, client); + this.#goaway = new Once(); // v17+: control/setup stream only carries GoAway void this.#runGoAway(control, version); } else { const adapter = new ControlStreamAdapter(quic, control, version, maxRequestId, client); this.#session = adapter; + this.#goaway = adapter.goaway; // Start the adapter read loop (routes control messages to virtual streams) void adapter.run().catch((err: unknown) => { if (!this.#closed) console.error("adapter error", err); @@ -142,7 +152,7 @@ export class Connection implements Established { this.#solicit = solicit; this.#cluster = cluster; this.#subscriber = new Subscriber({ session: this.#session, cluster, hidden }); - registerWire(this, { consume: (path) => this.#subscriber.consume(path) }); + registerWire(this, { consume: (path) => this.#subscriber.consume(path), goaway: this.#goaway }); void this.#run(); } @@ -317,18 +327,29 @@ export class Connection implements Established { /** * v17+ only: reads GoAway from the setup/control stream. + * + * The session keeps serving after a GOAWAY so its groups in flight can finish while the + * caller migrates; only the stream ending, or a second GOAWAY, closes it here. */ async #runGoAway(controlStream: Stream, version: IetfVersion) { try { - const done = await controlStream.reader.done(); - if (done) return; + for (;;) { + const done = await controlStream.reader.done(); + if (done) return; + + const typeId = await controlStream.reader.u53(); + if (typeId !== GoAway.id) { + console.warn(`unexpected message on setup stream: 0x${typeId.toString(16)}`); + return; + } - const typeId = await controlStream.reader.u53(); - if (typeId === GoAway.id) { const msg = await GoAway.decode(controlStream.reader, version); - console.warn(`received GOAWAY with redirect URI: ${msg.newSessionUri}`); - } else { - console.warn(`unexpected message on setup stream: 0x${typeId.toString(16)}`); + if (this.#goaway.peek() !== undefined) throw new ProtocolViolation("duplicate GOAWAY"); + // A client may leave, but only the server may name where to go. + if (!this.#client && msg.newSessionUri !== "") { + throw new ProtocolViolation("client GOAWAY must not name a redirect"); + } + this.#goaway.set(msg.drain()); } } catch (err) { if (!this.#closed) { diff --git a/js/net/src/ietf/goaway.test.ts b/js/net/src/ietf/goaway.test.ts new file mode 100644 index 0000000000..22c35e5eef --- /dev/null +++ b/js/net/src/ietf/goaway.test.ts @@ -0,0 +1,82 @@ +import { expect, test } from "bun:test"; +import { createMockTransportPair } from "../mock.ts"; +import { Stream } from "../stream.ts"; +import * as Time from "../time.ts"; +import { wireOf } from "../wire.ts"; +import { Connection } from "./connection.ts"; +import { GoAway } from "./goaway.ts"; +import { ALPN, Version } from "./version.ts"; + +// Draft-17+ carry GOAWAY on the setup stream, with the peer's deadline. The session keeps +// serving afterwards so its groups in flight finish while the caller migrates. +test("a draft-17 GOAWAY surfaces its URI and deadline without closing the session", async () => { + const version = Version.DRAFT_17; + const pair = createMockTransportPair(ALPN.DRAFT_17); + const control = await Stream.open(pair.client, { version }); + const peer = await Stream.accept(pair.server, version); + if (!peer) throw new Error("no setup stream"); + + const connection = new Connection({ + url: new URL("https://relay.example/"), + quic: pair.client, + control, + maxRequestId: 100n, + version, + client: true, + }); + + let closed = false; + void connection.closed.then(() => { + closed = true; + }); + + try { + await peer.writer.u53(GoAway.id); + await new GoAway({ newSessionUri: "", timeout: 5000n }).encode(peer.writer, version); + + const drain = await wireOf(connection).goaway; + expect(drain.uri).toBe(""); + expect(drain.timeout).toBe(Time.Milli(5000)); + + await new Promise((resolve) => setTimeout(resolve, 20)); + expect(closed).toBe(false); + } finally { + connection.close(); + } +}); + +test("a server rejects a client GOAWAY that names a redirect", async () => { + const version = Version.DRAFT_17; + const pair = createMockTransportPair(ALPN.DRAFT_17); + const control = await Stream.open(pair.client, { version }); + const peer = await Stream.accept(pair.server, version); + if (!peer) throw new Error("no setup stream"); + + const connection = new Connection({ + url: new URL("https://relay.example/"), + quic: pair.client, + control, + maxRequestId: 100n, + version, + client: false, + }); + + let closed = false; + void connection.closed.then(() => { + closed = true; + }); + + try { + await peer.writer.u53(GoAway.id); + await new GoAway({ newSessionUri: "https://other.example/", timeout: 0n }).encode(peer.writer, version); + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(closed).toBe(true); + } finally { + connection.close(); + } +}); + +test("a zero GOAWAY timeout reads as no deadline", async () => { + const msg = new GoAway({ newSessionUri: "https://relay.example/next", timeout: 0n }); + expect(msg.drain()).toEqual({ uri: "https://relay.example/next", timeout: undefined }); +}); diff --git a/js/net/src/ietf/goaway.ts b/js/net/src/ietf/goaway.ts index 9619b6ab4a..d77f188cfe 100644 --- a/js/net/src/ietf/goaway.ts +++ b/js/net/src/ietf/goaway.ts @@ -1,4 +1,6 @@ -import type { Reader, Writer } from "../stream.ts"; +import type { Drain } from "../connection/goaway.ts"; +import { Reader, type Writer } from "../stream.ts"; +import * as Time from "../time.ts"; import * as Message from "./message.ts"; import { type IetfVersion, Version } from "./version.ts"; @@ -35,6 +37,25 @@ export class GoAway { return Message.decode(r, (mr) => GoAway.#decode(mr, version)); } + /** + * Decode a body the caller already unframed, as the draft-14 to -16 control stream + * adapter does before routing a message. + */ + static async decodeBody(body: Uint8Array, version: IetfVersion): Promise { + const r = new Reader(undefined, body, version); + const msg = await GoAway.#decode(r, version); + if (!(await r.done())) throw new Error("GOAWAY has trailing bytes"); + return msg; + } + + /** The drain signal this message carries. A zero timeout is the wire saying "none". */ + drain(): Drain { + return { + uri: this.newSessionUri, + timeout: this.timeout > 0n ? Time.Milli(Number(this.timeout)) : undefined, + }; + } + static async #decode(r: Reader, version: IetfVersion): Promise { const newSessionUri = await r.string(); // All drafts cap the New Session URI at 8,192 bytes; a longer one is a diff --git a/js/net/src/lite/connection.test.ts b/js/net/src/lite/connection.test.ts index bf7717fefb..703a6702ae 100644 --- a/js/net/src/lite/connection.test.ts +++ b/js/net/src/lite/connection.test.ts @@ -1,7 +1,12 @@ import { expect, test } from "bun:test"; -import { probeLevel } from "./connection.ts"; +import { createMockTransportPair } from "../mock.ts"; +import { Stream } from "../stream.ts"; +import { wireOf } from "../wire.ts"; +import { Connection, probeLevel } from "./connection.ts"; +import { Goaway } from "./goaway.ts"; import { ProbeLevel } from "./setup.ts"; -import { Version } from "./version.ts"; +import { StreamId } from "./stream.ts"; +import { ALPN_04, Version } from "./version.ts"; /** A transport whose `getStats` behaves as described, or is absent entirely. */ function transport(getStats?: () => Promise): WebTransport { @@ -45,3 +50,39 @@ test("a throwing getStats advertises None rather than propagating", async () => }); expect(await probeLevel(quic, Version.DRAFT_05)).toBe(ProbeLevel.None); }); + +async function sendGoaway(server: WebTransport, uri: string): Promise { + const stream = await Stream.open(server); + await stream.writer.u53(StreamId.Goaway); + await new Goaway(uri).encode(stream.writer, Version.DRAFT_04); + stream.writer.close(); +} + +test("a lite GOAWAY keeps the session open, and a second one closes it", async () => { + const pair = createMockTransportPair(ALPN_04); + const connection = new Connection({ + url: new URL("https://relay.example/"), + quic: pair.client, + version: Version.DRAFT_04, + }); + + let closed = false; + void connection.closed.then(() => { + closed = true; + }); + + try { + await sendGoaway(pair.server, ""); + const drain = await wireOf(connection).goaway; + expect(drain.uri).toBe(""); + + await new Promise((resolve) => setTimeout(resolve, 20)); + expect(closed).toBe(false); + + await sendGoaway(pair.server, "https://other.example/"); + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(closed).toBe(true); + } finally { + connection.close(); + } +}); diff --git a/js/net/src/lite/connection.ts b/js/net/src/lite/connection.ts index 75e681d6f8..85b48684cf 100644 --- a/js/net/src/lite/connection.ts +++ b/js/net/src/lite/connection.ts @@ -1,9 +1,10 @@ -import { type Getter, Signal } from "@moq/signals"; +import { type Getter, Once, Signal } from "@moq/signals"; import type * as announce from "../announced.ts"; import type { Established } from "../connection/established.ts"; +import type { Drain } from "../connection/goaway.ts"; import { type Probe, type Stats, transportStats } from "../connection/stats.ts"; import { type Transport, transportOf } from "../connection/transport.ts"; -import { error, fromClose, StreamCode, StreamError } from "../error.ts"; +import { error, fromClose, ProtocolViolation, StreamCode, StreamError } from "../error.ts"; import { type Hop, randomHop } from "../hop.ts"; import type { Consumer as OriginConsumer } from "../origin.ts"; import type * as Path from "../path.ts"; @@ -94,6 +95,9 @@ export class Connection implements Established { // Written by the Subscriber as PROBE messages arrive. #probe = new Signal({}); + // The peer's GOAWAY. Lite carries no deadline, so only the URI is set. + #goaway = new Once(); + /** * The {@link Role} the peer advertised in its SETUP, for a server deciding whether the * peer's authorization grants the direction it intends to use. @@ -126,7 +130,7 @@ export class Connection implements Established { this.hop = randomHop(); this.#publisher = new Publisher(this.#quic, this.#version, this.hop, publish); this.#subscriber = new Subscriber(this.#quic, this.#version, this.hop, this.#probe, this.#peerSetup); - registerWire(this, { consume: (path) => this.#subscriber.consume(path) }); + registerWire(this, { consume: (path) => this.#subscriber.consume(path), goaway: this.#goaway }); void this.#run(); } @@ -218,6 +222,10 @@ export class Connection implements Established { this.#runBidi(stream) .catch((err: unknown) => { stream.writer.reset(err); + // A protocol violation on one stream is the peer breaking the session. + // Resetting that stream leaves it free to repeat the violation; a duplicate + // GOAWAY is the one this dispatcher raises. + if (err instanceof ProtocolViolation) this.close(); }) .finally(() => { stream.writer.close(); @@ -246,7 +254,9 @@ export class Connection implements Established { await this.#publisher.runProbe(stream); } else if (typ === StreamId.Goaway) { const msg = await Goaway.decode(stream.reader, this.#version); - console.info("received goaway:", msg.uri); + // A peer sends at most one; a second is a protocol violation. + if (this.#goaway.peek() !== undefined) throw new ProtocolViolation("duplicate GOAWAY"); + this.#goaway.set({ uri: msg.uri }); } else { throw new Error(`unknown stream type: ${typ.toString()}`); } diff --git a/js/net/src/origin.test.ts b/js/net/src/origin.test.ts index 47d044158e..bc7427169c 100644 --- a/js/net/src/origin.test.ts +++ b/js/net/src/origin.test.ts @@ -1300,6 +1300,52 @@ test("a rejected request is not asked of the same route again", async () => { origin.close(); }); +// A relay migration lands the replacement session's route next to the draining one's. The +// request must hand over to it, not drop to nothing while the new session answers. +test("an outranked route keeps serving until its replacement answers", async () => { + const origin = new Producer(); + const consumer = origin.consume(); + const path = Path.from("migrating"); + + const older = new BroadcastProducer(); + const keepOlder = older.consume(); + const disposeOlder = serve(origin, path, () => keepOlder.clone()); + + const request = consumer.request(path); + await settle(); + const first = request.active.peek(); + expect(first).toBeDefined(); + + let gaps = 0; + const stop = request.active.subscribe((active) => { + if (active === undefined) gaps++; + }); + + // The newer route wins the tie but has not answered yet. + const newer = wireOf(origin).receive(path); + const asked = newer.requested().next(); + await settle(); + expect(request.active.peek()).toBe(first); + + // Once it answers, the request swaps straight across. + const replacement = new BroadcastProducer(); + const { value: req } = await asked; + req?.accept(replacement.consume()); + await settle(); + expect(request.active.peek()).toBeDefined(); + expect(request.active.peek()).not.toBe(first); + expect(gaps).toBe(0); + + stop(); + request.close(); + newer.close(); + disposeOlder(); + keepOlder.close(); + older.close(); + replacement.close(); + origin.close(); +}); + test("a rejected request falls through to the next-best route", async () => { const origin = new Producer(); const consumer = origin.consume(); diff --git a/js/net/src/origin.ts b/js/net/src/origin.ts index 28dd06e62d..9d5a604623 100644 --- a/js/net/src/origin.ts +++ b/js/net/src/origin.ts @@ -321,11 +321,10 @@ class OriginState { const requests = this.requests.peek(); for (const [path, cached] of [...this.materialized]) { if (!Path.hasPrefix(prefix, path)) continue; - const refused = requests?.get(path)?.refused; - if (cached.entry !== this.bestEntry(path, (entry) => refused?.has(entry) ?? false)) { - this.materialized.delete(path); - cached.front.close(); - } + // A merely outranked provider is left to `route`, which holds it until the new one serves. + if (this.present(cached.entry)) continue; + this.materialized.delete(path); + cached.front.close(); } for (const [path, slot] of requests ?? []) { if (Path.hasPrefix(prefix, path)) slot.route.set(this.route(path, slot)); @@ -369,6 +368,14 @@ class OriginState { cached.front.close(); } + /** Whether `entry` is still in the table, rather than retracted. */ + present(entry: RouteEntry): boolean { + for (const entries of this.routes.peek()?.values() ?? []) { + if (entries.includes(entry)) return true; + } + return false; + } + /** The preferred entry on the most specific route covering `path`, ignoring skipped entries, if any. */ bestEntry(path: Path.Valid, skip?: (entry: RouteEntry) => boolean): RouteEntry | undefined { let bestPrefix: Path.Valid | undefined; @@ -405,7 +412,9 @@ class OriginState { * * Materialization is lazy and cached per path: the first request under a route opens * the providing session's subscription, repeats share it, and a provider change (the - * route retracting, a better session taking over) swaps it out. + * route retracting, a better session taking over) swaps it out. A route that was only + * outranked keeps serving until its replacement answers, so the swap never leaves the + * path unrouted in between: a relay migration hands over rather than dropping out. */ route(path: Path.Valid, slot: Pick): broadcast.Consumer | undefined { const entry = this.bestEntry(path, (candidate) => slot.refused.has(candidate)); @@ -416,23 +425,34 @@ class OriginState { return local; } - const cached = this.materialized.get(path); + let cached = this.materialized.get(path); if (cached && cached.entry === entry) { if (cached.front.closed.peek() === undefined) return cached.front; this.materialized.delete(path); - } else if (cached) { + cached = undefined; + } + // Whatever `cached` holds now belongs to another provider, and goes once this one serves. + const replace = () => { + if (!cached) return; this.materialized.delete(path); cached.front.close(); + }; + if (!entry?.server) { + replace(); + return slot.answer; } - if (!entry?.server) return slot.answer; const served = entry.server.served.get(path); if (served && served.closed.peek() === undefined) { + replace(); this.materialized.set(path, { entry, front: served }); return served; } entry.server.enqueue(path); + const standby = cached && !slot.refused.has(cached.entry) && this.present(cached.entry); + if (standby && cached?.front.closed.peek() === undefined) return cached?.front; + replace(); return undefined; } } diff --git a/js/net/src/wire.ts b/js/net/src/wire.ts index 3348b172bc..f4eb0a86d0 100644 --- a/js/net/src/wire.ts +++ b/js/net/src/wire.ts @@ -7,8 +7,9 @@ * * @module */ -import type { Dispose, Getter } from "@moq/signals"; +import type { Dispose, GetPromise, Getter } from "@moq/signals"; import type * as broadcast from "./broadcast.ts"; +import type { Drain } from "./connection/goaway.ts"; import type { Consumer as GroupConsumer } from "./group.ts"; import type { Route } from "./hop.ts"; import type * as origin from "./origin.ts"; @@ -53,9 +54,11 @@ export interface Advertised { readonly route: Route; } -/** The protocol-facing operation behind an established session. */ +/** The protocol-facing operations behind an established session. */ export interface Established { consume(path: Path.Valid): broadcast.Consumer; + /** Settles with the peer's GOAWAY; the session keeps serving until it closes. */ + readonly goaway: GetPromise; } type View = Broadcast | OriginProducer | OriginConsumer | Established; diff --git a/quest/m1/drain/README.md b/quest/m1/drain/README.md index 93428a5a8c..c8e0a45e04 100644 --- a/quest/m1/drain/README.md +++ b/quest/m1/drain/README.md @@ -11,13 +11,12 @@ node anyway (a cached resolve, or a pool alias) just gets another GOAWAY. Only after sessions drain or the stop deadline expires does the process exit and the new software boot. -GOAWAY only reaches MoQ sessions, and only the Rust client acts on it today: -`moq_tokio::Connection` migrates, while `js/net` decodes and logs the message -and at most closes the session afterwards (the IETF path does, the lite path -does not), leaving any reconnect to the ordinary close-triggered backoff -rather than migrating. The wire message is lite04+/IETF only besides. So the stop deadline is the +GOAWAY only reaches MoQ sessions. Both clients migrate on it: +`moq_tokio::Connection` and the `js/net` `Connection` dial the replacement +through a fresh resolve while the old session drains. The wire message is +lite04+/IETF only besides. So the stop deadline is the real backstop - for pre-lite04 versions, for client SDKs deployed before -client-goaway ships, and for in-process ingest gateways (RTMP/SRT/WHIP/WHEP), +the JS migration shipped, and for in-process ingest gateways (RTMP/SRT/WHIP/WHEP), which have no GOAWAY equivalent at all: their grace is the DNS-drain window stopping new arrivals plus the encoder's own reconnect. The DNS-drain-first ordering is what keeps that hard-close window small. @@ -33,18 +32,19 @@ The relay's drain hook has landed: `Relay::with_signals(false)` hands SIGTERM to the embedder, and its `shutdown_trigger` GOAWAYs every session, arrivals included, against one deadline. -**client-goaway.** The JS reconnector migrates like the Rust one, preserving -the app-visible session while resolving DNS again before dialing, and the Rust -path gains the regression test it lacks. This is a +**Clients (landed).** The JS reconnector migrates like the Rust one, +preserving the app-visible session while resolving DNS again before dialing. +Both are covered against stand-in servers. This is a scale-down prerequisite, not merely a deploy improvement. RTMP/SRT/WHIP/WHEP cannot receive MoQ GOAWAY, so their contract remains DNS withdrawal followed by the stop deadline and encoder reconnect. +**End to end.** The line's own remaining work: a JS client watching a live +track through an in-tree relay drained with the drain hook migrates to a +second relay behind the same name without a dropped group. + ## Quests -- [Client goaway](/quest/m1/drain/client-goaway.md) - the JavaScript client - migrates on GOAWAY with a handover and the guarded redirect the Rust client - already has, and the Rust drain path gets its regression test - [Drain exit](/quest/m1/drain/drain-exit.md) - a drain ends as soon as every session has left, and reports whether that or the deadline ended it diff --git a/quest/m1/drain/client-goaway.md b/quest/m1/drain/client-goaway.md deleted file mode 100644 index 5d611d75dc..0000000000 --- a/quest/m1/drain/client-goaway.md +++ /dev/null @@ -1,87 +0,0 @@ -# [L] Client goaway - -## Goal - -The JavaScript client migrates on GOAWAY the way the Rust client already -does: it dials the replacement while the old session keeps serving, follows a -redirect URI under the same guard, and keeps the app-visible handle and its -origins across the swap. The Rust client's fleet-drain path (an empty-URI -GOAWAY redialed through a fresh DNS resolve) gains the regression test it -lacks today. - -## Plan - -moq.pro's (downstream) fleet drain orchestration relies on this behavior: a -drained node is already out of DNS when GOAWAY fires, so a re-resolve is what -lands clients on a healthy relay. - -Everything below describes `dev`, which is where this quest lands: `main` -still has `moq-native`'s close-only `Reconnect`. The drain loop is written -once into `Connection` and its pool. - -### Rust policy and proof - -`moq_tokio::Connection` already handles GOAWAY: the session loop returns the -message, `Redirect::resolve` guards the URI (scheme tier never drops, the host -is pinned to the configured one by default, and `follow` is the opt-in that -lets a peer name another), the loop redials while a -`Draining` handle keeps the old session serving until it closes or overstays -the handover cap, and `Status::Migrating` is visible to callers. Every dial -resolves DNS again, since `Addrs` holds URLs and each backend resolves at dial -time, so nothing is pinned. The relay always sends an empty URI, and the only -end-to-end coverage is the cluster sibling test with a redirect. - -Add the fleet-drain regression test: a relay GOAWAYs with an empty URI and a -timeout, the client redials the configured URL through a fresh resolve (a -resolver the test can repoint), live tracks hand over at a group boundary, -and the old session closes within the handover cap. - -### JavaScript is greenfield - -`js/net` logs the lite GOAWAY URI and keeps that session open, -logs the IETF draft-17+ URI and then closes the session when its control -loop ends, and on the draft-14 to -16 shared control stream reads the message -body but returns without decoding it. Nothing migrates: `Connection` -reconnects only after `closed` fires, through its backoff, tears the old -connection down in its effect cleanup, and the pool -(`js/net/src/connection/pool.ts`) keys by URL href. - -- Surface the peer's GOAWAY on the live session as a drain signal carrying the - resolved URI and the timeout, decoded on every wire the client speaks, - including the draft-14 to -16 adapter route that currently returns without - decoding the body it has already read. -- `Connection` mirrors `Draining`: on GOAWAY it dials the target immediately, - swaps the origin wiring (`forwardAnnounced`, `publish`, `subscribe`) once - the replacement is established, and leaves the old session to close on its - own or at a handover cap: the configured cap, lowered to the peer's timeout - when the wire carried a positive one. Lite and IETF drafts 14 to 16 carry no - timeout, and the IETF decoder reads an absent one as zero, so absence means - the cap and never a zero-length handover. Groups in flight finish. A GOAWAY does not go through the backoff delay; a failed - replacement dial does. -- Port the guard: same-host by default, refuse a scheme-tier drop or a - widening to a local host, and offer the follow mode. An empty URI preserves - the current address list, including caller-selected fallbacks, and starts - normal migration. A malformed or policy-refused explicit redirect ends the - connection with a typed terminal error; it must not trigger a retry against - the original address or another configured fallback. Only an accepted - redirect replaces the list. Apply and test this policy in Rust as well as JS; - do not assume the existing Rust behavior already satisfies it. A redirect with a certificate pin (`serverCertificateHashes`) - is refused unless the host is unchanged, since the pin cannot verify another - relay; the pool already refuses to share pinned connections. -- The pool re-keys its entry to the redirect target, so a later - caller configured with that URL shares the migrated connection. The app's - handle and shared origin are unchanged; only the pool key moves, and a - caller still asking for the original URL gets a fresh entry. When the - target key already holds a live entry, that entry wins: the migrating entry - is removed from the pool without replacing the target entry. Existing - handles keep its migrated connection, but a new lookup for the original URL - dials fresh and a target lookup joins the target entry. The migrated entry - retires when its last existing handle releases it. Two connections to one - relay for that overlap is the honest cost; entry removal is identity-guarded - so neither cleanup can delete another entry. -- Tests against the in-tree relay: an empty-URI drain migrates without a - dropped group, a GOAWAY without a timeout hands over at the configured cap, - a redirect moves the pool key and origins, a redirect onto an already - pooled key keeps existing handles on both entries but makes a new caller for - the original URL dial fresh, each guard refusal closes rather than - reconnects, and the draft-14 to -16 route decodes the URI. diff --git a/quest/m1/transport-upgrade/README.md b/quest/m1/transport-upgrade/README.md index 7e88b12e6f..194c13b0bd 100644 --- a/quest/m1/transport-upgrade/README.md +++ b/quest/m1/transport-upgrade/README.md @@ -33,9 +33,8 @@ every version (a client may send one with an empty URI; only a redirect URI is forbidden to a moq-transport client). The origin's multi-route front prefers the newest of two equal routes and `resume` splices each track at a group boundary, capping the old segment so the old session's subscription ends at the -boundary on its own. The JavaScript handover is the -[client goaway](/quest/m1/drain/client-goaway.md) quest's, so the JS half -requires it. +boundary on its own. The JavaScript handover ships with the +[drain](/quest/m1/drain/README.md) line, so the JS half requires it. Shared decisions: @@ -58,7 +57,7 @@ Shared decisions: ## Quests - [Rust](/quest/m1/transport-upgrade/rust.md) - moq-tokio keeps the QUIC dial after WebSocket wins and migrates through the existing Draining path -- [JavaScript](/quest/m1/transport-upgrade/js.md) - js/net keeps the WebTransport dial after WebSocket wins and migrates through the client-goaway handover +- [JavaScript](/quest/m1/transport-upgrade/js.md) - js/net keeps the WebTransport dial after WebSocket wins and migrates through the drain line's GOAWAY handover ## Related diff --git a/quest/m1/transport-upgrade/js.md b/quest/m1/transport-upgrade/js.md index 4cb58a3c59..f98faccec7 100644 --- a/quest/m1/transport-upgrade/js.md +++ b/quest/m1/transport-upgrade/js.md @@ -11,10 +11,12 @@ nothing changes. ## Plan -Lands in `js/net`, after the [client goaway](/quest/m1/drain/client-goaway.md) -quest ships the handover it reuses: dial the replacement while the old session -keeps serving, swap the origin wiring once it is established, leave the old -session to close on its own or at the handover cap. See the +Lands in `js/net`, after the [drain](/quest/m1/drain/README.md) line ships the +GOAWAY handover it reuses (`Reload`'s migration in +`js/net/src/connection/reload.ts`): dial the replacement while the old session +keeps serving, let both feed the origin (a request holds the outranked route +until the new one answers), leave the old session to close on its own or at the +handover cap. See the [questline](/quest/m1/transport-upgrade/README.md) for the shared decisions. - `connectInner` (`js/net/src/connection/connect.ts`) currently resolves @@ -37,10 +39,10 @@ session to close on its own or at the handover cap. See the group across the upgrade, the WebSocket session closes within the cap, and the next connect to the same URL gives WebTransport the head start again; with no delay, WebTransport wins and no WebSocket session is ever opened. -- Public API: none beyond what client-goaway adds; `transportOf` already +- Public API: none beyond what the drain line adds; `transportOf` already reports the live transport. Update `doc/lib/js` where the fallback race is described. ## Required -- [Client goaway](/quest/m1/drain/client-goaway.md) - the handover this upgrade reuses +- [Graceful relay drains](/quest/m1/drain/README.md) - ships the JS GOAWAY handover this upgrade reuses diff --git a/rs/moq-ffi/src/binary.rs b/rs/moq-ffi/src/binary.rs index 8e343bd2ad..1ffa8ca881 100644 --- a/rs/moq-ffi/src/binary.rs +++ b/rs/moq-ffi/src/binary.rs @@ -48,7 +48,8 @@ impl MoqBroadcastProducer { let _guard = crate::ffi::enter(); self.with_state(|state| { let track = state.broadcast.create_track(name, None)?; - let producer = state.catalog.binary_snapshot(track, config.into())?; + let config: moq_mux::binary::Config = config.into(); + let producer = state.catalog.binary_snapshot(track, config)?; Ok(Arc::new(MoqBinarySnapshotProducer { inner: std::sync::Mutex::new(Some(producer)), })) @@ -66,7 +67,8 @@ impl MoqBroadcastProducer { let _guard = crate::ffi::enter(); self.with_state(|state| { let track = state.broadcast.create_track(name, None)?; - let producer = state.catalog.binary_stream(track, config.into())?; + let config: moq_mux::binary::Config = config.into(); + let producer = state.catalog.binary_stream(track, config)?; Ok(Arc::new(MoqBinaryStreamProducer { inner: std::sync::Mutex::new(Some(producer)), })) diff --git a/rs/moq-tokio/src/client.rs b/rs/moq-tokio/src/client.rs index e9133a8e89..bc1744e633 100644 --- a/rs/moq-tokio/src/client.rs +++ b/rs/moq-tokio/src/client.rs @@ -57,6 +57,9 @@ pub struct Client { pub(crate) reconnect: bool, pub(crate) backoff: Backoff, pub(crate) goaway: Goaway, + /// Whether the TLS config pins a certificate fingerprint, which only verifies + /// the configured host, so a GOAWAY may not redirect elsewhere. + pub(crate) pinned: bool, /// The resolved Happy Eyeballs timings, used by the `tcp://` dial here; the /// QUIC backend captures its own copy from the config. #[cfg(feature = "tcp")] @@ -135,6 +138,7 @@ impl Client { reconnect: !config.once.unwrap_or(false), backoff: config.backoff, goaway: config.goaway, + pinned: !config.tls.fingerprint.is_empty(), #[cfg(feature = "tcp")] failover_delay, #[cfg(feature = "tcp")] diff --git a/rs/moq-tokio/src/connection.rs b/rs/moq-tokio/src/connection.rs index 13f3fb4e2a..6038f285af 100644 --- a/rs/moq-tokio/src/connection.rs +++ b/rs/moq-tokio/src/connection.rs @@ -209,27 +209,42 @@ pub enum Redirect { impl Redirect { /// Resolve the URL to dial after a GOAWAY, falling back to `current` when the /// redirect is empty ("reconnect to me"), malformed, or refused by policy. + /// + /// Lenient on purpose, for a caller that only wants somewhere to dial. A + /// [`Connection`] is stricter: it ends with [`Error::RefusedRedirect`] rather + /// than redialing after a malformed or refused URI. pub fn resolve(&self, uri: &str, current: &Url) -> Url { - self.target(uri, current).unwrap_or_else(|| current.clone()) + self.target(uri, current, false) + .ok() + .flatten() + .unwrap_or_else(|| current.clone()) } - // Absence keeps the caller's address list; only an accepted URI replaces it. - fn target(&self, uri: &str, current: &Url) -> Option { - if uri.is_empty() || matches!(self, Self::Ignore) { - return None; + /// The URL a GOAWAY assigns. `Ok(None)` keeps the current address list (the + /// peer named no URI, or the policy ignores a URI it could parse), `Ok(Some)` + /// replaces it, and `Err` is an explicit URI this policy refuses. A malformed + /// URI is refused even under [`Self::Ignore`]. + /// + /// `pinned` is a certificate pin on the connection, which can only verify the + /// host it was configured for, so it refuses a host change even under + /// [`Self::Follow`]. + fn target(&self, uri: &str, current: &Url, pinned: bool) -> crate::Result> { + if uri.is_empty() { + return Ok(None); } - let Ok(target) = uri.parse::() else { - tracing::warn!(uri, "malformed GOAWAY URI; keeping the current addresses"); - return None; - }; + // The URI can carry credentials, so the error names the reason, never the URI. + // Parse before `Ignore`: a malformed redirect is terminal even when the policy + // would otherwise stay on the current address list. + let refuse = |reason: &str| Error::RefusedRedirect(reason.to_string()); + + let target = uri.parse::().map_err(|_| refuse("the GOAWAY URI is malformed"))?; + if matches!(self, Self::Ignore) { + return Ok(None); + } if scheme_tier(target.scheme()) < scheme_tier(current.scheme()) { - tracing::warn!( - uri, - "GOAWAY redirect downgrades the scheme; keeping the current addresses" - ); - return None; + return Err(refuse("the GOAWAY redirect downgrades the scheme")); } // Only as far as the URL itself says: a name is dialed, never resolved here, @@ -237,27 +252,29 @@ impl Redirect { // nothing about one that hides the same address behind a hostname. That gap // is why [`Self::SameHost`] is the default; see [`is_local`]. if is_local(&target) && !is_local(current) { - tracing::warn!( - uri, - "GOAWAY redirect widens reachability; keeping the current addresses" - ); - return None; + return Err(refuse("the GOAWAY redirect widens reachability to a local address")); } // Host only, not the full authority: the port is what a peer legitimately // moves us across when it hands off to a sibling process on the same box. - if matches!(self, Self::SameHost) && target.host_str() != current.host_str() { - tracing::warn!( - uri, - "GOAWAY redirect leaves the current host; keeping the current addresses" - ); - return None; + let same_host = target.host_str() == current.host_str(); + if matches!(self, Self::SameHost) && !same_host { + return Err(refuse("the GOAWAY redirect leaves the current host")); + } + if pinned && !same_host { + return Err(refuse("the GOAWAY redirect leaves the host a certificate pin verifies")); } - Some(target) + Ok(Some(target)) } } +/// Whether this scheme's dial installs the Rustls verifier, so a configured +/// fingerprint actually checked the peer. Plain and non-Rustls transports do not. +fn fingerprint_pins(scheme: &str) -> bool { + matches!(scheme, "https" | "wss" | "moqt" | "moql") +} + /// Rank a scheme so a peer-supplied redirect cannot silently drop encryption. /// Unknown schemes rank lowest, so a forgotten classification is refused. fn scheme_tier(scheme: &str) -> u8 { @@ -762,11 +779,20 @@ impl Connection { // ended sooner counts as a failed attempt however it ended. let healthy = connected.elapsed() >= initial; - // The connected target owns the policy, including in one-shot mode. - if let Ended::Goaway(msg) = &ended - && addr.addresses().is_some() - && goaway.redirect.target(msg.uri(), &url).is_some() - { + // The connected target owns the policy, including in one-shot mode. A + // refused redirect is terminal: the peer is leaving and named somewhere we + // won't go, so redialing the old address or a fallback would ignore it. + let assigned = match &ended { + // A fingerprint only checked the dial that installed the Rustls + // verifier. tcp, unix, iroh, and plaintext WebSocket never consult it. + Ended::Goaway(msg) => { + goaway + .redirect + .target(msg.uri(), &url, client.pinned && fingerprint_pins(url.scheme()))? + } + Ended::Closed(_) => None, + }; + if assigned.is_some() && addr.addresses().is_some() { return Err(Error::PinnedRedirect); } @@ -784,7 +810,7 @@ impl Connection { // An accepted redirect is an assignment: keep dialing it from here on, and // only it. The peer named exactly one place to go, which retires // whatever other addresses got us to this session. - let url = if let Some(target) = goaway.redirect.target(msg.uri(), &url) { + let url = if let Some(target) = assigned { addrs = Addrs::new(target.clone()); target } else { @@ -1495,8 +1521,8 @@ mod tests { assert_eq!(Redirect::Follow.resolve("https://other.example/", &plain), same); } - /// The three ways a redirect resolves to "redial what we already had": the peer - /// naming no URI, a URI we cannot parse, and a policy that ignores it outright. + /// `resolve` is the lenient form: every way a redirect can fail to assign a + /// new URL, refusals included, lands back on the current one. #[test] fn resolve_falls_back_to_the_current_url() { let current: Url = "https://relay.example/".parse().unwrap(); @@ -1506,11 +1532,7 @@ mod tests { current, "empty means 'reconnect to me'" ); - assert_eq!( - Redirect::Follow.resolve("not a url", ¤t), - current, - "a malformed URI is not a reason to stop reconnecting" - ); + assert_eq!(Redirect::Follow.resolve("not a url", ¤t), current); assert_eq!( Redirect::Ignore.resolve("https://other.example/", ¤t), current, @@ -1575,26 +1597,73 @@ mod tests { ); } + /// Only an explicit URI can assign, and one the policy will not follow is an + /// error rather than a quiet fallback: the loop ends on it instead of redialing. #[test] fn only_an_accepted_redirect_replaces_the_address_list() { let current: Url = "https://relay.example/".parse().unwrap(); + + // No URI, or a policy that ignores it: keep the current address list. + assert_eq!(Redirect::Follow.target("", ¤t, false).unwrap(), None); + assert_eq!( + Redirect::Ignore + .target("https://relay.example:5443/", ¤t, false) + .unwrap(), + None + ); + + // An explicit URI the policy will not follow is refused. for (policy, uri) in [ (Redirect::SameHost, "https://other.example/"), - (Redirect::Follow, ""), (Redirect::Follow, "not a url"), - (Redirect::Ignore, "https://relay.example:5443/"), + (Redirect::Ignore, "not a url"), (Redirect::Follow, "http://relay.example/"), (Redirect::Follow, "https://127.0.0.1/"), ] { - assert_eq!(policy.target(uri, ¤t), None, "{policy:?}: {uri}"); + assert!( + matches!(policy.target(uri, ¤t, false), Err(Error::RefusedRedirect(_))), + "{policy:?}: {uri}" + ); } + // An explicit assignment remains an assignment even if its URL is unchanged. assert_eq!( - Redirect::SameHost.target(current.as_str(), ¤t), + Redirect::SameHost.target(current.as_str(), ¤t, false).unwrap(), Some(current.clone()) ); let moved = "https://relay.example:5443/"; - assert_eq!(Redirect::SameHost.target(moved, ¤t), Some(moved.parse().unwrap())); + assert_eq!( + Redirect::SameHost.target(moved, ¤t, false).unwrap(), + Some(moved.parse().unwrap()) + ); + } + + /// A fingerprint is a Rustls check. Schemes that never install that verifier + /// must not inherit the pin. + #[test] + fn a_fingerprint_pin_only_covers_rustls_schemes() { + for scheme in ["https", "wss", "moqt", "moql"] { + assert!(fingerprint_pins(scheme), "{scheme}"); + } + for scheme in ["http", "ws", "tcp", "unix", "iroh"] { + assert!(!fingerprint_pins(scheme), "{scheme}"); + } + } + + /// A certificate pin verifies only the host it was configured for, so it + /// refuses a host change even when the policy would follow one. + #[test] + fn a_certificate_pin_holds_the_host() { + let current: Url = "https://relay.example/".parse().unwrap(); + assert!(matches!( + Redirect::Follow.target("https://other.example/", ¤t, true), + Err(Error::RefusedRedirect(_)) + )); + let moved = "https://relay.example:5443/"; + assert_eq!( + Redirect::Follow.target(moved, ¤t, true).unwrap(), + Some(moved.parse().unwrap()) + ); } /// `SameHost` lets a peer move us between ports or schemes on the endpoint we @@ -1950,4 +2019,134 @@ mod tests { assert_eq!(attempt_timeout(1, 3), Some(CONNECT_ATTEMPT), "still one more"); assert_eq!(attempt_timeout(2, 3), None, "the last of three"); } + + /// A stream-only server on a free loopback port, publishing `origin`. + /// + /// Returns its address, a receiver yielding each accepted session (so the test + /// can drain it), and the listener task. Probing for a free port races other + /// tests between the probe closing and the real bind, so this retries. + #[cfg(feature = "tcp")] + async fn serve( + origin: moq_net::origin::Producer, + ) -> ( + std::net::SocketAddr, + tokio::sync::mpsc::UnboundedReceiver, + tokio::task::JoinHandle<()>, + ) { + for _ in 0..20 { + let probe = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = probe.local_addr().unwrap(); + drop(probe); + + let mut config = crate::listen::Config::default(); + config.tcp.bind = Some(addr); + let Ok(mut server) = config.init(Default::default()).unwrap().listen().await else { + continue; + }; + + let (accepted, sessions) = tokio::sync::mpsc::unbounded_channel(); + let task = tokio::spawn(async move { + while let Some(request) = server.accept().await { + if let Ok(session) = request.with_publisher(&origin).ok().await { + let _ = accepted.send(session); + } + } + }); + return (addr, sessions, task); + } + panic!("could not bind a free TCP port after 20 attempts"); + } + + /// The fleet drain: a relay withdrawn from DNS sends an empty-URI GOAWAY with a + /// deadline, and the client lands on a healthy relay by resolving the configured + /// name again. The drained relay still accepts, so a cached resolve would land + /// right back on it. The live track hands over at a group boundary, and the old + /// session closes at our handover cap, well before the peer's own deadline. + #[cfg(feature = "tcp")] + #[tokio::test] + async fn a_fleet_drain_redials_through_a_fresh_resolve() { + const WAIT: Duration = Duration::from_secs(10); + const HANDOVER: Duration = Duration::from_millis(500); + const DEADLINE: Duration = Duration::from_secs(30); + + // Two relays of one fleet, serving the same live broadcast. + let origin = crate::origin::spawn(); + let broadcast = origin.create_broadcast("cam").unwrap(); + broadcast.announce(Default::default()).unwrap(); + let track = broadcast.create_track("video", None).unwrap(); + let (addr_a, mut accepted_a, _task_a) = serve(origin.clone()).await; + let (addr_b, mut accepted_b, _task_b) = serve(origin.clone()).await; + + // Unique to this test: the table is process-wide. + const HOST: &str = "fleet-drain.test"; + crate::resolve::hosts::point(HOST, [addr_a]); + + let subscriber = crate::origin::spawn(); + let mut config = crate::connect::Config::default(); + config.goaway.handover = HANDOVER; + // A session younger than the initial delay counts as redirected immediately and + // waits out a backoff, which this test is not about. + config.backoff.initial = MIN_BACKOFF; + let client = config + .init(Default::default()) + .unwrap() + .with_subscriber(subscriber.clone()); + let url: Url = format!("tcp://{HOST}:1/").parse().unwrap(); + let _connection = client.connect(url); + + let session_a = tokio::time::timeout(WAIT, accepted_a.recv()).await.unwrap().unwrap(); + + let consumer = subscriber.consume(); + let cam = tokio::time::timeout(WAIT, consumer.routed_broadcast("cam")) + .await + .unwrap() + .unwrap(); + let mut sub = cam.track("video").unwrap().subscribe(None).await.unwrap(); + + let mut group = track.append_group().unwrap(); + group.write_frame(moq_net::Timestamp::ZERO, b"g0".as_ref()).unwrap(); + group.finish().unwrap(); + let g0 = tokio::time::timeout(WAIT, sub.recv_group()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!(g0.sequence, 0); + + // Withdraw A from DNS, then drain it. + crate::resolve::hosts::point(HOST, [addr_b]); + let drained = tokio::time::Instant::now(); + session_a + .drain() + .send(moq_net::goaway::Goaway::new().with_timeout(DEADLINE)) + .unwrap(); + + let _session_b = tokio::time::timeout(WAIT, accepted_b.recv()) + .await + .expect("never redialed through the fresh resolve") + .unwrap(); + + let mut group = track.append_group().unwrap(); + group.write_frame(moq_net::Timestamp::ZERO, b"g1".as_ref()).unwrap(); + group.finish().unwrap(); + let mut g1 = tokio::time::timeout(WAIT, sub.recv_group()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!(g1.sequence, 1, "delivery resumes at the next group after the swap"); + assert_eq!(g1.read_frame().await.unwrap().unwrap().payload[..], b"g1"[..]); + + tokio::time::timeout(WAIT, session_a.closed()) + .await + .expect("the drained session never closed"); + assert!( + drained.elapsed() < DEADLINE, + "the old session outlived the handover cap and waited for the peer's deadline" + ); + assert!( + accepted_a.try_recv().is_err(), + "a cached resolve redialed the drained relay" + ); + } } diff --git a/rs/moq-tokio/src/error.rs b/rs/moq-tokio/src/error.rs index 209306d96a..f603fc1307 100644 --- a/rs/moq-tokio/src/error.rs +++ b/rs/moq-tokio/src/error.rs @@ -58,6 +58,11 @@ pub enum Error { #[error("peer redirect refused for a connection with fixed addresses")] PinnedRedirect, + /// A peer's GOAWAY named a redirect the connection's policy refuses, or one it + /// could not parse. Terminal: the peer is leaving, so redialing is ignoring it. + #[error("GOAWAY redirect refused: {0}")] + RefusedRedirect(String), + /// Reading or writing a socket, certificate, or key file failed. #[error(transparent)] Io(Arc), diff --git a/rs/moq-tokio/src/resolve.rs b/rs/moq-tokio/src/resolve.rs index 7410cae79a..3f58d119fa 100644 --- a/rs/moq-tokio/src/resolve.rs +++ b/rs/moq-tokio/src/resolve.rs @@ -249,6 +249,11 @@ impl Candidates { }, }; + #[cfg(test)] + if let Some(addrs) = hosts::lookup(domain) { + return Self::fixed(addrs); + } + Self { full: Query::start(domain, port, Lookup::Full), ipv4: Query::start(domain, port, Lookup::Ipv4), @@ -549,6 +554,32 @@ impl Candidates { } } +/// A name table tests can repoint, consulted before the system resolver. +/// +/// DNS is what moves a client off a drained node, so proving a redial resolves +/// afresh needs a name whose answer changes between dials. Entries carry the +/// port too, so two servers on one loopback address can stand in for two hosts. +#[cfg(test)] +pub(crate) mod hosts { + use std::collections::HashMap; + use std::net::SocketAddr; + use std::sync::Mutex; + + static HOSTS: Mutex>>> = Mutex::new(None); + + /// Answer every later lookup of `host` with `addrs`, ports included. + pub(crate) fn point(host: &str, addrs: impl IntoIterator) { + let mut hosts = HOSTS.lock().unwrap(); + hosts + .get_or_insert_with(HashMap::new) + .insert(host.to_string(), addrs.into_iter().collect()); + } + + pub(super) fn lookup(host: &str) -> Option> { + HOSTS.lock().unwrap().as_ref()?.get(host).cloned() + } +} + #[cfg(test)] impl Query { /// A lookup that answers with `addrs` after `delay`. diff --git a/rs/moq-tokio/tests/reconnect.rs b/rs/moq-tokio/tests/reconnect.rs index a3532e354c..c2ecf04bf7 100644 --- a/rs/moq-tokio/tests/reconnect.rs +++ b/rs/moq-tokio/tests/reconnect.rs @@ -133,8 +133,7 @@ async fn spawn_server() -> ( (port, sessions, handle) } -/// A client that redials fast, so a refused redirect lands back on the original -/// server inside the test's patience. +/// A client that redials fast, so a migration lands inside the test's patience. fn quick_client(redirect: moq_tokio::Redirect) -> moq_tokio::Client { let mut config = moq_tokio::connect::Config::default(); config.backoff.initial = Duration::from_millis(20); @@ -144,14 +143,15 @@ fn quick_client(redirect: moq_tokio::Redirect) -> moq_tokio::Client { config.init(Default::default()).expect("failed to init client") } -/// The default refuses peer-selected host changes and redials the configured URL. +/// The default refuses a peer-selected host change, and the refusal ends the +/// connection: the peer is leaving, so redialing the configured URL would ignore it. #[tokio::test] async fn a_redirect_to_another_host_is_refused_by_default() { let (port_a, mut sessions_a, _task_a) = spawn_server().await; let (port_b, mut sessions_b, _task_b) = spawn_server().await; let url: url::Url = format!("tcp://localhost:{port_a}/").parse().expect("parse url"); - let _connection = quick_client(Default::default()).connect(url); + let connection = quick_client(Default::default()).connect(url); let first = tokio::time::timeout(Duration::from_secs(10), sessions_a.recv()) .await @@ -163,12 +163,13 @@ async fn a_redirect_to_another_host_is_refused_by_default() { .send(moq_net::goaway::Goaway::redirect(format!("tcp://127.0.0.1:{port_b}/"))) .expect("send goaway"); - // Refused, so the redial goes back to A rather than to the host the peer named. - tokio::time::timeout(Duration::from_secs(10), sessions_a.recv()) + let err = tokio::time::timeout(Duration::from_secs(10), connection.closed()) .await - .expect("never redialed the configured URL") - .expect("server A stopped accepting"); + .expect("the refusal never ended the connection") + .expect_err("a refused redirect must end with an error"); + assert!(matches!(err, moq_tokio::Error::RefusedRedirect(_)), "ended with {err}"); + assert!(sessions_a.try_recv().is_err(), "redialed the configured URL"); assert!( sessions_b.try_recv().is_err(), "the peer moved us onto the host it named" @@ -202,21 +203,21 @@ async fn follow_still_honors_a_cross_host_redirect() { .expect("server B stopped accepting"); } -/// Refusing a peer-selected host must not discard caller-selected fallbacks. +/// A refused redirect must not fall through to a caller-selected fallback either. #[tokio::test] -async fn a_refused_redirect_preserves_configured_fallbacks() { +async fn a_refused_redirect_skips_configured_fallbacks() { let (port_a, mut sessions_a, task_a) = spawn_server().await; let (port_b, mut sessions_b, _task_b) = spawn_server().await; let primary: url::Url = format!("tcp://localhost:{port_a}/").parse().expect("primary URL"); let fallback: url::Url = format!("tcp://127.0.0.1:{port_b}/").parse().expect("fallback URL"); let addrs = moq_tokio::connect::Addrs::new(primary).or(fallback); - let _connection = quick_client(Default::default()).connect(addrs); + let connection = quick_client(Default::default()).connect(addrs); let first = tokio::time::timeout(Duration::from_secs(10), sessions_a.recv()) .await .expect("first dial timed out") .expect("server A stopped accepting"); - // Stop accepting before GOAWAY so reconnect must use the configured fallback. + // Stop accepting before GOAWAY so any redial could only land on the fallback. task_a.abort(); assert!(task_a.await.expect_err("listener was aborted").is_cancelled()); first @@ -224,6 +225,36 @@ async fn a_refused_redirect_preserves_configured_fallbacks() { .send(moq_net::goaway::Goaway::redirect("tcp://127.0.0.1:1/")) .expect("send goaway"); + let err = tokio::time::timeout(Duration::from_secs(10), connection.closed()) + .await + .expect("the refusal never ended the connection") + .expect_err("a refused redirect must end with an error"); + assert!(matches!(err, moq_tokio::Error::RefusedRedirect(_)), "ended with {err}"); + assert!( + sessions_b.try_recv().is_err(), + "fell through to the configured fallback" + ); +} + +/// An empty GOAWAY is "reconnect to me", which keeps every caller-selected fallback. +#[tokio::test] +async fn an_empty_goaway_preserves_configured_fallbacks() { + let (port_a, mut sessions_a, task_a) = spawn_server().await; + let (port_b, mut sessions_b, _task_b) = spawn_server().await; + let primary: url::Url = format!("tcp://localhost:{port_a}/").parse().expect("primary URL"); + let fallback: url::Url = format!("tcp://127.0.0.1:{port_b}/").parse().expect("fallback URL"); + let addrs = moq_tokio::connect::Addrs::new(primary).or(fallback); + let _connection = quick_client(Default::default()).connect(addrs); + let first = tokio::time::timeout(Duration::from_secs(10), sessions_a.recv()) + .await + .expect("first dial timed out") + .expect("server A stopped accepting"); + + // Stop accepting before GOAWAY so the migration must use the configured fallback. + task_a.abort(); + assert!(task_a.await.expect_err("listener was aborted").is_cancelled()); + first.drain().send(moq_net::goaway::Goaway::new()).expect("send goaway"); + tokio::time::timeout(Duration::from_secs(10), sessions_b.recv()) .await .expect("configured fallback was discarded")