Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 5 additions & 3 deletions doc/bin/relay/config.md
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
2 changes: 1 addition & 1 deletion doc/lib/js/net.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<ReadonlyMap<Path.Valid, Route>>` of the same covered prefixes for UIs that need the current set. A borrowed `Connection.origin` also exposes `dynamic(prefix, route)` for serving paths on demand.
- **Subscriptions** carry a priority, a `Time.Milli` max age, and optional `groups` bounds. Groups arrive out of order and are read frame by frame, with `Error.TooFarBehind` when a reader asks for a frame the group never held and `Error.GroupTooLarge` when a write exceeds the cache budget and aborts the group.
Expand Down
3 changes: 2 additions & 1 deletion js/net/src/connection/forward.test.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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;
});
Expand Down
110 changes: 110 additions & 0 deletions js/net/src/connection/goaway.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
224 changes: 224 additions & 0 deletions js/net/src/connection/goaway.ts
Original file line number Diff line number Diff line change
@@ -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;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Reject redirect schemes the JavaScript client cannot dial

An HTTPS session can accept a same-host redirect such as moqt://relay.example/ because moqt: has the same scheme tier, but the JavaScript connector only has WebTransport/WebSocket dial paths and cannot establish moqt:, moql:, or tcp: URLs. The accepted assignment then becomes sticky and the pooled/private loops, which use an unlimited retry window, back off forever instead of reporting the documented terminal RefusedRedirect. Validate the redirect against the schemes connect() actually supports before returning it. (Written by GPT-5.6 Sol)

AGENTS.md reference: AGENTS.md:L16-L18

Useful? React with 馃憤聽/ 馃憥.

}

/** 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<string>(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
);
}
Loading
Loading