Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
39 commits
Select commit Hold shift + click to select a range
5d06fb8
quest: plan the 14 open issues without a quest (#4388)
kixelated Sep 28, 2026
a2d501d
chore: bump quest to 8590d2a (#4413)
kixelated Sep 28, 2026
61f642d
refactor(watch): share hang's catalog containment check (#4415)
kixelated Sep 29, 2026
9048b33
fix(net): an origin::Dynamic keeps its origin alive (#4417)
kixelated Sep 29, 2026
0ded2ca
fix(cli): refuse every MoQ-side flag a verb never reads (#4419)
kixelated Sep 29, 2026
de4db07
fix(net): skip announce updates the peer cannot tell apart (#4423)
kixelated Sep 29, 2026
138a051
fix: stop logging a stream reset before its header as UnknownSession …
kixelated Sep 29, 2026
9537167
doc(uring): noq paces inside poll_transmit, not ignored (#4400)
kixelated Sep 29, 2026
4e9c9ed
test: prove stopped relays and worker groups closed their sockets ins…
kixelated Sep 29, 2026
ea148fe
feat(net): drain queued stream data before a graceful close (#4430)
kixelated Sep 29, 2026
bd28408
chore(quest): JS IETF reprices a namespace in place (#4424)
kixelated Sep 29, 2026
785ee4c
chore(quest): drop the duplicate Opus mapping family on dev (#4427)
kixelated Sep 29, 2026
5b36481
fix(tokio): handle IPv6 literals in TLS server names (#4322)
shermerL Sep 29, 2026
ecf90c9
docs(quest): drop suffix-based routing from the plans (#4382)
kixelated Sep 29, 2026
1ee10df
fix(kio): keep a lost waiter's recorded lists idempotent (#4422)
kixelated Sep 29, 2026
a0929c9
fix(moq-gst): moqsrc waits for its session to end on stop (#4416)
kixelated Sep 29, 2026
8e13f46
quest(export-linger): tell a finished TS track from a removed one (#4…
kixelated Sep 29, 2026
a31afc4
quest: IETF peers without MoQ Hidden see hidden namespaces (#4394)
kixelated Sep 29, 2026
42ea3eb
chore(quest): graceful close in bindings and the moq fetch flake caus…
kixelated Sep 29, 2026
22f6d68
fix(sock): resolve an ephemeral reuseport group's port with a plain b…
kixelated Sep 29, 2026
bf59c93
fix(video): re-anchor native capture above its last timestamp (#4418)
kixelated Sep 29, 2026
539850d
docs(quest): record the routing simulator's findings in cluster routi…
kixelated Sep 29, 2026
2f4b785
fix(uring): publish a local close only once its CONNECTION_CLOSE is s…
kixelated Sep 29, 2026
ed56cab
chore(quest): plan the PR-merge follow-ups (#4435)
kixelated Sep 29, 2026
75d93c3
fix(tokio): drain a GOAWAY predecessor on Connection::close (#4436)
kixelated Sep 29, 2026
57eaf69
feat(publish): announce once every captured rendition resolves (#4414)
kixelated Sep 29, 2026
7eebe7c
chore(quest): cover the inline path and first-hop change in JS IETF r…
kixelated Sep 29, 2026
da93877
test(tokio): match a worker group's full address when counting its so…
kixelated Sep 29, 2026
b14a61f
fix(tokio): keep a WebTransport session's H3 streams open while it cl…
mmcc Sep 29, 2026
8ddae84
feat(net): BigInt-free varint codec, internal U64, and checked Varint…
kixelated Sep 29, 2026
9d2a4f6
chore(quest): settle the PR-merge session's follow-ups (#4468)
kixelated Sep 29, 2026
93b7888
chore(quest): retire flate-binary on main and track the JS closed-tra…
kixelated Sep 29, 2026
62ef9a8
Merge remote-tracking branch 'origin/main' into sync-main-into-dev
kixelated Sep 29, 2026
68f5bc6
chore(quest): a cluster publisher change updates in place (#4467)
kixelated Sep 29, 2026
e9ef2bf
chore(quest): fix the WebTransport close capsule upstream (#4497)
kixelated Sep 29, 2026
4d08ba6
Merge remote-tracking branch 'origin/main' into sync-main-into-dev
kixelated Sep 29, 2026
4aa0b04
test(cli): wait for the publisher's route, not dev's first announce e…
kixelated Sep 29, 2026
b8647f1
test(tokio): main's close-drain tests skip dev's Live marker
kixelated Sep 29, 2026
33d5dea
test(go): an ended empty origin reports Live before the end
kixelated Sep 29, 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
6 changes: 6 additions & 0 deletions .github/workflows/nightly.yml
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,12 @@ jobs:
if: ${{ !cancelled() }}
run: nix develop --command bun js/net/bench/frames.ts

# Times one varint encode and decode per wire format. No threshold: it only
# has to keep running, and its numbers are the baseline for the generated codec.
- name: JS varint benchmark
if: ${{ !cancelled() }}
run: nix develop --command bun js/net/bench/varint.ts

# Fails if publishing a group costs more as the track retains more groups,
# which means the latency guard or the cache went back to scanning them all.
- name: JS track retention benchmark
Expand Down
2 changes: 2 additions & 0 deletions Cargo.lock

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

6 changes: 3 additions & 3 deletions doc/bin/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -225,9 +225,9 @@ zero-based `frame` and padded standard base64.

`<track>` is the literal track name. `/fetch` splits its path on the last `/`,
so the two agree only for names without one. Fetch only dials `--connect`, and
refuses a listener or cluster flag. It gives up after 30 seconds, as `/fetch`
does, and exits non-zero when the broadcast or group is not found (before
writing anything), the relay refuses, or the deadline passes.
refuses any listener, cluster, auth, or `--hop` flag. It gives up after 30
seconds, as `/fetch` does, and exits non-zero when the broadcast or group is not
found (before writing anything), the relay refuses, or the deadline passes.

## Multiple stages

Expand Down
12 changes: 6 additions & 6 deletions doc/bin/relay/cluster.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,18 +49,18 @@ link costs 1, which reproduces plain hop counting. Each relay adds the price of
the link an announcement arrived on before forwarding it, so a route's cost is
the sum of what it crossed.

Wildcard advertisements are forwarded and costed the same way as an exact-path
Prefix advertisements are forwarded and costed the same way as an exact-path
route: each hop appends its identity, adds the link price, and passes the
claim on. An advertisement must be contained by one of the publisher's granted
prefixes (`grant/**`); an over-wide pattern is refused rather than clamped.
claim on. An advertised prefix must overlap the publisher's grant, or it is
refused. A prefix wider than the grant is accepted, but it only routes requests
for paths the grant covers.

Routing prefers the most specific pattern, then a fully identified hop list
Routing prefers the longest covering prefix, then a fully identified hop list
over one that holds a 0 (an anonymous hop) at any depth, then the lowest cost,
then the shortest hop list, breaking any remaining tie toward the newest
announcement so a reconnecting publisher isn't outranked by the session it
replaced. An assigned identity for an anonymous peer is local selection state
and is never written into the hop list. Resolving a non-prefix pattern into a
subscription is not implemented yet.
and is never written into the hop list.

```toml
[cluster]
Expand Down
4 changes: 2 additions & 2 deletions doc/bin/relay/config.md
Original file line number Diff line number Diff line change
Expand Up @@ -77,8 +77,8 @@ io_uring = false # Drive them with io_uring instead of tokio
```

Packets are steered by connection ID, so a client that migrates stays with its
worker. The group shares one port, including an ephemeral (zero) port: the
first worker binds it and the rest join that port. Use an explicit port unless
worker. The group shares one port, including an ephemeral (zero) port, which is
resolved once and joined by every worker. Use an explicit port unless
something reads the bound address at startup. `workers` needs the `noq`
feature and real certificate files rather than `tls.generate`. A build without
QUIC rejects `workers` instead of
Expand Down
8 changes: 4 additions & 4 deletions doc/lib/go/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -95,10 +95,10 @@ Paths with a `.`-prefixed segment below the prefix are [hidden](/concept/moq-lit
`Hidden: true`.

An `OriginProducer` from `moq.NewOriginProducer` has no `Close`: its origin
ends when the garbage collector reaches the last producer, and every consumer
and `OriginDynamic` made from it then fails with `moq.ErrClosed`. Keep the
producer reachable (a field on a long-lived struct, or `runtime.KeepAlive`) for
as long as the origin should serve.
ends when the garbage collector reaches the last owner (each producer, published
broadcast, and `OriginDynamic`), and every consumer made from it then fails with
`moq.ErrClosed`. Keep an owner reachable (a field on a long-lived struct, or
`runtime.KeepAlive`) for as long as the origin should serve.

Every call that can block takes a `context.Context` first. Cancelling it
returns `ctx.Err()` promptly and tears the in-flight native work down, so a
Expand Down
3 changes: 2 additions & 1 deletion doc/lib/js/hang.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,8 @@ import * as Container from "@moq/hang/container";
`Catalog.watch(broadcast)` iterates validated catalog roots. It throws
`Catalog.TooManyRenditions` for an update above the 64 rendition limit, and
`Catalog.EscapingBroadcast` for a `broadcast` reference that walks above the
handle's `path`.
handle's `path`. Run the same checks on a catalog from another source with
`Catalog.checkRenditions(root)` and `Catalog.checkResolvable(root, base)`.
`Hang.Timeline.Consumer.subscribe(broadcast, root.archive)` reads segment
`push`, `pop`, and `skip` events when a root advertises an archive.

Expand Down
2 changes: 1 addition & 1 deletion doc/lib/js/publish.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ WebCodecs, writes the catalog, and publishes a hang broadcast.
| `source` | `camera`, `screen`, or `file`. |
| `muted`, `invisible` | Disable audio or video capture. |
| `preview` | What the nested element shows: the raw `source` (default), a decoded copy of the `encoded` stream to see what viewers get, or `none`. |
| `announce` | When to advertise: once a `source` is live (default), `always`, or `never`. A camera source waits for every enabled track. The broadcast is created while connected either way, but nobody can see or subscribe to it until it is announced. |
| `announce` | When to advertise: once a `source` is live (default), `always`, or `never`. A camera source waits for every enabled track, and `source` waits until each captured track's config resolves or fails, so the first catalog lists every rendition. The broadcast is created while connected either way, but nobody can see or subscribe to it until it is announced. |

A nested `<video>` gets the raw capture stream; a `<canvas>` is drawn by the
element. `<moq-publish-support>` shows what the browser can encode.
Expand Down
4 changes: 4 additions & 0 deletions doc/lib/rs/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,10 @@ let mut broadcast = origin.publish("my-stream.hang", Default::default())?;
// each requested path for the application to accept or reject.
```

Before exiting, `session.close().await` delivers what was already queued, such as
the tracks you just finished, within one second. Then `Client::close` (on a clone of
the client) sends the QUIC close before the runtime stops.

The examples run the session and the origin work concurrently (`tokio::select!` or
`spawn`), since the announcement loop is live. Runnable examples:
[`rs/hang/examples/video.rs`](https://github.com/moq-dev/moq/blob/main/rs/hang/examples/video.rs)
Expand Down
17 changes: 14 additions & 3 deletions doc/lib/rs/moq-net.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,16 +57,27 @@ on external activity, `Ok(None)` only on external activity, and `Err` is the
terminal error (`Error::Closed` for a clean finish): stop polling. Tests drive
the same interface with explicitly advanced instants.

Dropping the last session handle requests closure on the next poll. Dropping
the driver cancels the session. `moq-tokio` and `moq-wasm` drive sessions for
their callers.
Dropping the last session handle requests closure on the next poll, and
`session.abort(err)` closes with `err`'s code. Either discards stream data the
peer has not acknowledged yet. `session.close().await` first waits, up to one
second, for finished tracks to deliver their last groups and FIN, returning
`Error::Timeout` if it gave up. Finish or abort live tracks before calling it.
moq-transport (IETF) sessions close without waiting. Dropping the driver
cancels the session. `moq-tokio` and `moq-wasm` drive sessions for their
callers.

`origin::Producer::new` returns a driver with the same `time::Driver`
interface. It calls `cache::Pool::gc(now)` after each poll and folds the next
cleanup time into its returned deadline. A standalone pool needs `gc(now)`
called by its owner, at least by the returned deadline; `None` means expiry is
disabled.

The origin driver finishes once every owner is gone: each `origin::Producer`
clone, published broadcast, and `origin::Dynamic`. Producer-side children pin
their parent where no cycle exists, so a handler can drop its producer and keep
serving. Read handles (`origin::Consumer`, announce cursors) never pin it: they
see `Closed` once the owners are gone.

Cache activity is dated lazily: reads and writes mark a group active without
reading a clock, and the next `gc` pass stamps it with the supplied instant.
Expiry is therefore approximate; a late `gc` extends retention.
Expand Down
8 changes: 4 additions & 4 deletions flake.lock

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

2 changes: 1 addition & 1 deletion flake.nix
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
# The quest CLI, which also serves the quest guide and skills the stubs in
# .claude/skills call. Bump the rev to upgrade them.
quest = {
url = "github:kixelated/quest/46d7fe89247919583632e4963aee1c9a68dfe059";
url = "github:kixelated/quest/8590d2a1ddd91c2f499adf37b78aad0d673e3228";
inputs.nixpkgs.follows = "nixpkgs";
inputs.flake-utils.follows = "flake-utils";
inputs.crane.follows = "crane";
Expand Down
4 changes: 2 additions & 2 deletions go/wrapper/moq_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,8 @@ import (
const testTimeout = 10 * time.Second

// newOrigin returns an origin that lasts the whole test. An OriginProducer has
// no Close: the collector ends its origin once nothing reaches the producer,
// even while consumers and dynamic handles made from it are still in use.
// no Close: the collector ends its origin once nothing reaches an owner, even
// while consumers made from it are still in use.
func newOrigin(t *testing.T) *moq.OriginProducer {
origin := moq.NewOriginProducer()
t.Cleanup(func() { runtime.KeepAlive(origin) })
Expand Down
7 changes: 4 additions & 3 deletions go/wrapper/origin.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,10 @@ import (
// discover them. Wire one as both a client's/server's publish source and
// consume sink for a full-duplex peer.
//
// There is no Close: the origin ends once the collector reaches every
// producer, and its consumers and dynamic handles then fail with [ErrClosed].
// Keep a producer reachable for as long as the origin should live.
// There is no Close: the origin ends once the collector reaches every owner
// (each producer, published broadcast, and [OriginDynamic]), and its consumers
// then fail with [ErrClosed]. Keep an owner reachable for as long as the origin
// should live.
type OriginProducer struct {
inner *ffi.MoqOriginProducer
}
Expand Down
71 changes: 64 additions & 7 deletions go/wrapper/origin_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,44 @@ import (
)

// An OriginProducer has no Close: the collector ends its origin by finalizing
// the last producer, even while a consumer or dynamic handle made from it is in
// use. Destroying the handle is what that finalizer does, so this pins the
// behavior the doc comment warns about without waiting on the collector.
// the last owner, even while a consumer made from it is in use. Destroying the
// handle is what that finalizer does, so this pins the behavior the doc comment
// warns about without waiting on the collector.
func TestOriginEndsWithItsProducer(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

origin := NewOriginProducer()
consumer := origin.Consume()
announced, err := consumer.Announced(AnnounceOptions{})
if err != nil {
t.Fatal(err)
}
defer announced.Cancel()

origin.inner.Destroy()

// The teardown runs on the origin's driver; the cursor ending is its signal. An
// empty origin reports Live first.
update, err := announced.Next(ctx)
if _, live := update.(AnnounceEventLive); live && err == nil {
update, err = announced.Next(ctx)
}
if update != nil || err != nil {
t.Fatalf("Next = (%v, %v), want the end of the stream", update, err)
}
if _, err := consumer.RequestBroadcast(ctx, "live"); !errors.Is(err, ErrClosed) {
t.Fatalf("RequestBroadcast err = %v, want ErrClosed", err)
}
}

// An OriginDynamic owns its origin: once the collector finalizes the last
// OriginProducer, the route still serves and a consumer made earlier still
// resolves through it.
func TestDynamicKeepsTheOrigin(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

origin := NewOriginProducer()
dynamic, err := origin.Dynamic("", Route{})
if err != nil {
Expand All @@ -25,10 +56,36 @@ func TestOriginEndsWithItsProducer(t *testing.T) {

origin.inner.Destroy()

if _, err := dynamic.RequestedBroadcast(ctx); !errors.Is(err, ErrClosed) {
t.Fatalf("RequestedBroadcast err = %v, want ErrClosed", err)
type result struct {
broadcast *BroadcastConsumer
err error
}
if _, err := consumer.RequestBroadcast(ctx, "live"); !errors.Is(err, ErrClosed) {
t.Fatalf("RequestBroadcast err = %v, want ErrClosed", err)
requested := make(chan result, 1)
go func() {
broadcast, err := consumer.RequestBroadcast(ctx, "live")
requested <- result{broadcast: broadcast, err: err}
}()

request, err := dynamic.RequestedBroadcast(ctx)
if err != nil {
t.Fatalf("RequestedBroadcast err = %v, want a request", err)
}
served, err := NewBroadcastProducer()
if err != nil {
t.Fatal(err)
}
defer func() { _ = served.Close() }()
if err := request.Accept(served); err != nil {
t.Fatal(err)
}

var res result
select {
case res = <-requested:
case <-ctx.Done():
t.Fatal(ctx.Err())
}
if res.err != nil {
t.Fatalf("RequestBroadcast err = %v, want the served broadcast", res.err)
}
}
28 changes: 27 additions & 1 deletion js/hang/src/catalog/consumer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,15 @@ import { expect, test } from "bun:test";
import * as Json from "@moq/json";
import * as Moq from "@moq/net";
import { TRACK } from "./format";
import { checkRenditions, EscapingBroadcast, MAX_RENDITIONS, type Root, TooManyRenditions, watch } from "./root";
import {
checkRenditions,
checkResolvable,
EscapingBroadcast,
MAX_RENDITIONS,
type Root,
TooManyRenditions,
watch,
} from "./root";

function catalog(count: number): Root {
return {
Expand All @@ -29,6 +37,24 @@ test("the shared cap accepts 64 renditions and refuses 65 with a typed error", (
).toThrow(TooManyRenditions);
});

test("the shared containment check covers every section carrying a broadcast reference", () => {
// @moq/watch runs this same check on the hangz, MSF, and manual catalogs, so a section left
// out here would exempt its tracks everywhere.
const base = Moq.Path.from("a/b");
for (const [section, key] of [
["video", "renditions"],
["audio", "renditions"],
["text", "renditions"],
["json", "tracks"],
["binary", "tracks"],
]) {
const reference = (broadcast: string) =>
({ [section]: { [key]: { entry: { broadcast: Moq.Path.normalizeRelative(broadcast) } } } }) as Root;
expect(() => checkResolvable(reference("../../x"), base)).toThrow(EscapingBroadcast);
expect(checkResolvable(reference("../x"), base)).toBeDefined();
}
});

test("watch refuses an oversized catalog update", async () => {
const broadcast = new Moq.Broadcast.Producer();
const track = broadcast.createTrack(TRACK);
Expand Down
4 changes: 2 additions & 2 deletions js/hang/src/catalog/root.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,8 @@ export class EscapingBroadcast extends Error {
}
}

// Refuse an update with a `broadcast` reference that walks above the root from `base`.
function checkResolvable(root: Root, base: Moq.Path.Valid): Root {
/** Refuse an update with a `broadcast` reference that walks above the root from `base`. */
export function checkResolvable(root: Root, base: Moq.Path.Valid): Root {
// Every section carrying a `broadcast` reference must be listed here; one left out
// silently exempts its tracks from the check.
const sections = [
Expand Down
5 changes: 5 additions & 0 deletions js/hang/src/container/consumer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1652,3 +1652,8 @@ for (const end of [
}
});
}

test("LegacyFormat rejects a timestamp past 2^53 - 1 instead of rounding", () => {
const frame = Varint.encode(2n ** 53n + 1n);
expect(() => new LegacyFormat("video").decode(frame)).toThrow(/larger than 53-bits/);
});
16 changes: 16 additions & 0 deletions js/loc/src/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -133,3 +133,19 @@ test("Format decodes the draft-03 timestamp property", () => {
expect(decoded.timestamp).toBe(4242 as Time.Micro);
expect(decoded.payload).toEqual(payload);
});

test("Format skips an unknown property whose value needs all 62 bits", () => {
const props = concat(
Varint.encode(0x02),
Varint.encode(2n ** 62n - 1n),
Varint.encode(PROP_TIMESTAMP - 0x02),
Varint.encode(1_000),
);
const [decoded] = new Format().decode(buildFrame(props, new Uint8Array([1])));
expect(decoded.timestamp).toBe(1_000 as Time.Micro);
});

test("Format rejects a timestamp past 2^53 - 1 instead of rounding", () => {
const props = concat(Varint.encode(PROP_TIMESTAMP), Varint.encode(2n ** 53n));
expect(() => new Format().decode(buildFrame(props, new Uint8Array()))).toThrow(/larger than 53-bits/);
});
Loading
Loading