Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
705f6b5
quest: open the drain line
kixelated Sep 25, 2026
a40cc9a
feat(relay): hand the drain signal to embedders and drain arrivals to…
kixelated Sep 25, 2026
c895831
Merge origin/main into quest/m1/drain/README
kixelated Sep 25, 2026
4c44009
feat(net): migrate the JS connection on GOAWAY; refuse bad redirects …
kixelated Sep 25, 2026
5e04fdb
feat(relay): end a drain once every session has left (#4186)
kixelated Sep 25, 2026
cc15339
Merge remote-tracking branch 'origin/main' into quest/m1/drain/README
kixelated Sep 26, 2026
4665c13
test(relay): hold an io_uring drain open with a straggler
kixelated Sep 26, 2026
14519f4
test(drain): a JS viewer migrates off a draining relay without a drop…
kixelated Sep 26, 2026
944f864
style(net): format the merged error import
kixelated Sep 26, 2026
22b7d3c
style(relay): format the drain tests
kixelated Sep 26, 2026
5139fa8
quest(drain): keep the line open for a zero-budget JS handover
kixelated Sep 26, 2026
a6a67f5
Merge origin/main into quest/m1/drain/README
kixelated Sep 28, 2026
0863efe
Merge remote-tracking branch 'origin/main' into drain-land
kixelated Sep 30, 2026
c273150
quest(drain): finish the line; promote the JS handover quests to m1
kixelated Sep 30, 2026
cba7744
fix(relay): close the shared listener before run returns
kixelated Sep 30, 2026
7dab39c
Merge remote-tracking branch 'origin/main' into drain-land
kixelated Sep 30, 2026
b063f0c
fix(drain): address review on the drain line
kixelated Sep 30, 2026
944b4f6
Merge remote-tracking branch 'origin/main' into drain-land
kixelated Sep 30, 2026
714bb03
test(net): open GOAWAY streams with a version after the lite-07 varin…
kixelated Sep 30, 2026
18a8dc9
Merge remote-tracking branch 'origin/main' into merge-drain-4132
kixelated Oct 1, 2026
3a86732
test(drain): subscribe to the route before peeking it
kixelated Oct 1, 2026
6794296
Merge remote-tracking branch 'origin/main' into merge-drain-4132
kixelated Oct 1, 2026
13d1a2d
Merge remote-tracking branch 'origin/main' into merge-drain-4132
kixelated Oct 1, 2026
e24fdb8
quest: drop the plain-text blocker from redirect-resolve
kixelated Oct 1, 2026
af4e1b0
Merge remote-tracking branch 'origin/main' into merge-drain-4132
kixelated Oct 1, 2026
4860d9d
Merge remote-tracking branch 'origin/main' into merge-drain-4132
kixelated Oct 1, 2026
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
5 changes: 4 additions & 1 deletion .github/workflows/nightly.yml
Original file line number Diff line number Diff line change
Expand Up @@ -170,8 +170,11 @@ jobs:
matrix:
# `rs uring` is a separate feature compile (off the default set) and
# kernel-gated below 6.12. The embedding tests for io_uring live there.
# `test drain` stands up two relays and drains one under a live JS
# viewer: diff-independent, since a break can arrive through the relay,
# js/net, or the wire between them.
# `gh test` exercises the release tooling, which no pull request runs.
recipe: ["rs doctest --workspace", "rs loom", "test drill-sensitivity", "rs uring", "gh test"]
recipe: ["rs doctest --workspace", "rs loom", "test drill-sensitivity", "rs uring", "test drain", "gh test"]
steps:
- name: Free disk space
uses: jlumbroso/free-disk-space@ceedf095f4ec1a097402bc6bd80831f2e1a6fde6 # main
Expand Down
8 changes: 8 additions & 0 deletions bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

30 changes: 27 additions & 3 deletions doc/bin/relay/config.md
Original file line number Diff line number Diff line change
Expand Up @@ -179,9 +179,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 Expand Up @@ -238,6 +240,28 @@ secret = "./iroh-secret.key" # Persist the key so the endpoint id surviv

See [Transport](/concept/transport#iroh-peer-to-peer-experimental).

## Shutdown

```toml
drain_timeout = "10s" # Top-level key, as --drain-timeout / MOQ_DRAIN_TIMEOUT.
```

The first SIGTERM or SIGINT starts a drain: every session is sent a GOAWAY
asking it to reconnect, and is force-closed if it is still connected when the
window ends. A session that connects during the drain, such as a client with a
cached DNS answer, is sent a GOAWAY immediately, with only the time left in
the window. The relay exits as soon as every session has left, when the window
ends, or immediately on a second signal. `0` skips the GOAWAY and closes every
session at once.

The exit is logged with how long the drain took, as either
`drain complete: every session left` or `drain deadline force-closed sessions`
with the number `forced`. A session still in its handshake when the last one
leaves is not waited for.
Only moq-lite-04+ and moq-transport clients act on a GOAWAY; older ones are
closed when the window ends. An embedder can take over the signals and start
the drain itself; see [Embed](/bin/relay/#embed).

## \[log]

```toml
Expand Down
4 changes: 3 additions & 1 deletion doc/bin/relay/http.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,9 @@ split by `tier` and `role`, plus accept-loop counters per TCP listener. Alert
on `moq_relay_accept_failures_total{class="exhausted"}`, which means the
process ran out of a resource `accept` needs. Content dropped for drifting past
a subscriber's budget is counted separately as `moq_relay_stale_bytes_total`
and friends. Host CPU and memory belong to a node exporter.
and friends. During a [shutdown drain](/bin/relay/config#shutdown),
`moq_relay_draining_sessions` counts the sessions sent a GOAWAY that have not
left yet. Host CPU and memory belong to a node exporter.

Traffic and session counters accumulate for the node's lifetime, including
broadcasts and sessions that have ended. The stats publishing prefix (normally
Expand Down
14 changes: 10 additions & 4 deletions doc/bin/relay/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,12 +81,18 @@ TLS, and a certificate fingerprint for client pinning.

The accessors borrow and `run` consumes the relay, so clone `cluster`,
`auth`, `client`, `stats`, `shutdown`, and `shutdown_trigger` for application
tasks before calling it. `trigger.start()` drains every session with a GOAWAY
and `run` returns once the drain window elapses, with the listeners released
and the workers joined. Build routes from `web().routes()` (or
tasks before calling it. `trigger.start()` drains every session with a GOAWAY,
including any that connect afterwards, and `run` returns once every session
has left or the drain window elapses, with the listeners released and the
workers joined. `run` also starts
the drain on SIGTERM or SIGINT. An application that owns those signals, for
example to withdraw the node from DNS and wait out the TTL before draining,
calls `with_signals(false)` and fires the trigger itself. Build routes from `web().routes()` (or
`internal().routes()`): `with_web` replaces the router, so `Router::new()`
drops the built-in routes. Extra listeners (RTMP, SRT, ...) sit beside `run`
in the application's `select!`. `runtime.workers` and `runtime.io_uring` stay
in the application's `select!`. The drain reaches only MoQ sessions: an
extra listener has no GOAWAY, so the application stops it on its own deadline
and leaves the encoder to reconnect. `runtime.workers` and `runtime.io_uring` stay
inside the owner; do not split the worker group yourself. An application that
decides admissions itself leaves `[auth]` empty and answers
`relay.admissions()`; see [Authentication](/bin/relay/auth#in-process). See
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. A refusal from the most specific route ends the request: `request.closed` settles with the handler's error and no broader route is asked. 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);
});
Loading
Loading