Skip to content

feat(net): spread a prefix's pool across paths, stitch on the reply's origin - #4279

Merged
kixelated merged 6 commits into
quest/m1/wildcard/READMEfrom
quest/m1/wildcard/spread
Sep 27, 2026
Merged

kixelated merged 6 commits into
quest/m1/wildcard/READMEfrom
quest/m1/wildcard/spread

Conversation

@kixelated

@kixelated kixelated commented Sep 26, 2026 •

Copy link
Copy Markdown
Collaborator

Completes the quest/m1/wildcard/spread.md quest, the last child of the wildcard line (#4037).

Problem

Equal-cost advertisers of one prefix (a transcode pool claiming its service prefix) did not share its paths: route_order hashed the advertised prefix, so the first member took every path. Spreading them breaks first-hop pinning as an identity, since one pool route now labels many origins, so failover needs to know which origin actually served a subscription.

Approach

  • Spread: route_order's hash tie-break is keyed on the requested path instead of the advertised prefix. Cost still orders first, so a distant worker stays overflow. One path resolves to the same advertiser on every node holding the same routes.
  • Stitching identity from the reply (lite-07): SUBSCRIBE_OK gains Origin (i), and a fetch stream now opens with a FETCH_OK naming the origin before its frames. Each copy (one session's copy of a track) records the origin its replies name and holds its content until the relay's front admits that origin. Group streams wait for their subscription's SUBSCRIBE_OK, fetched frames wait behind FETCH_OK, and datagrams are dropped until the origin is admitted. The front takes its identity from the first source's first reply and only admits replacements naming the same origin. A mismatch, or a replacement whose wire cannot name one, aborts the tracks spliced from that source and ends the front; the subscriber re-requests. No content from another origin is ever delivered.
    • A front with a reply-named origin stays on its live route when a better one appears (Pin::Stay): which member the new route leads to is unknown until it replies.
    • A relay names its front's origin in its own SUBSCRIBE_OK and FETCH_OK: the reply-named origin, the route's first hop for sources on older wires, its own hop for content originating here, or a random one per front for content nobody identifies.
    • SUBSCRIBE_OK goes out ahead of the first group served by stream or datagram, so a datagram-only track still names its origin.
    • Older versions (lite-06 and earlier, moq-transport) carry no origin and keep first-hop pinning unchanged.
  • JS: @moq/net's route selector breaks cost and hop-count ties on the same Spread Hash, so a JS node picks the same pool member per path as Rust. It also records the origin an upstream SUBSCRIBE_OK or FETCH_OK names and proxies it when republishing; content with no upstream origin gets a random one per broadcast. A draft-07 JS subscriber holds groups until SUBSCRIBE_OK and drops datagrams before it, like Rust. The publisher writes SUBSCRIBE_OK/END through one ordered writer shared by the group and datagram loops.
  • Draft: Spread Hash spelled out in Routing, the identity rule rewritten around the reply's Origin (hold groups / drop datagrams until SUBSCRIBE_OK), Origin on SUBSCRIBE_OK, the new FETCH_OK message, changelog entries.
  • Docs: doc/bin/relay/cluster.md failover and tie-break paragraphs.
  • Quest: deletes spread.md, sizes the line README [S], and adds Cluster origin reply (moq-transport downstreams) and Verified route upgrade to m1.

Impact

  • Wire: lite-07 (moq-lite-07-wip) SUBSCRIBE_OK appends Origin (i), and a new FETCH_OK { Origin (i) } precedes a fetch's frames. TRACK_INFO is unchanged. No change to published versions.
  • Wire behavior: on lite-05+ SUBSCRIBE_OK is now also sent ahead of a subscription's first datagram (resolving the start there), not only ahead of its first group stream. Within the existing spec ("the first group that will be delivered").
  • Rust public API: none. track::Provenance and track::Info::names_origin are crate-private; lite::SubscribeStart gains origin: Hop and lite::FetchOk is new.
  • JS: lite.SubscribeStart gains origin (defaults to unknown) and lite.FetchOk is new; the internal broadcast wire gains origin() / name().
  • Behavior: route selection now spreads distinct paths across an equal-cost pool on every protocol. On lite-07 a subscription's groups wait for its SUBSCRIBE_OK (every subscriber, not just relays), and datagrams racing ahead of it are dropped.

Tests

  • equal_cost_pool_spreads_paths: 64 paths over a 4-member pool spread (each member takes at least 1/16), identical when the routes arrive in reverse order. Fails without the fix. The JS twin an equal-cost pool spreads its paths the same way on every node does the same, and one Spread Hash vector is pinned in both languages.
  • failover_resumes_on_the_origin_the_reply_names, failover_never_splices_another_pool_member, a_named_origin_stays_on_its_live_route, front state-machine unit tests, and a bounded-sequence invariant: an admitted origin is always the one the front serves.
  • tests/pool_failover.rs: end to end over lite-07 mock sessions, two workers behind a pool relay behind a downstream relay. Killing the serving worker ends the downstream subscription instead of splicing the survivor's frames, a re-request is served by the survivor, and a fetch crosses both relays through FETCH_OK.
  • Datagrams: datagram_only_track_names_its_origin_on_lite07 (Rust, mock sessions) and moq-lite-07-wip datagram delivery (JS integration) fail without SUBSCRIBE_OK before the first datagram; datagram_waits_for_the_admitted_origin pins the Rust subscriber drop.
  • Rust and JS SUBSCRIBE_OK and FETCH_OK wire round trips with cross-language bytes.

Decisions

Settled by the maintainer: Origin rides SUBSCRIBE_OK and FETCH_OK (not TRACK_INFO); reply-named fronts only fail over; every Rust subscriber holds a subscription's groups until SUBSCRIBE_OK; JS proxies the received origin or generates a random one; the line README is [S]; moq-transport downstreams are the Cluster origin reply quest; front upgrades are the Verified route upgrade quest.

Alternatives

  • Origin in TRACK_INFO (the first revision): left a window between TRACK and SUBSCRIBE where another origin's content could splice.
  • Buffering datagrams until SUBSCRIBE_OK instead of dropping them: they are best-effort and uncached, so the first RTT's worth is dropped instead.

Known gaps

  • @moq/net records the origin per consumed broadcast, not per track, and the latest reply wins. It has no failover of its own, so this only matters for a JS app republishing a broadcast whose upstream front changed origin mid-life.
  • A front with no reply-named origin that ends while a copy still holds unconfirmed content refuses that copy, so a reader spliced onto it ends rather than receiving unverified groups.
  • No benchmark sweeps pool size x requested paths for resolution, or tracks per front for the driver's per-event admission walk.

Checks

just check, just drafts check, and just test interop --all pass.

(Written by Claude Opus 5.5)

🤖 Generated with Claude Code

kixelated and others added 4 commits September 26, 2026 10:54
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… origin

Key the route_order hash on the requested path so equal-cost advertisers of
one prefix share its paths. On lite-07, TRACK_INFO names the origin serving
the track, and a relay's front splices a failover only between copies naming
the same origin, since a pool's one route labels many origins. Older wires
keep first-hop pinning.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Move lite-07's Origin out of TRACK_INFO into SUBSCRIBE_OK and a new FETCH_OK
ahead of a fetch's frames, so the origin names the content it arrives with.
Each copy records the origin its replies name and holds its groups until the
front admits that origin, so a replacement from another origin never delivers
anything. Plan the moq-transport follow-up as the cluster-origin quest.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…one exists

@moq/net records the origin an upstream SUBSCRIBE_OK or FETCH_OK names and
proxies it when republishing; content with no upstream origin gets a random
one per broadcast instead of the connection hop. A draft-07 JS subscriber
holds groups until SUBSCRIBE_OK, like Rust. Rust names a random origin per
front for content nobody identifies, instead of 0. Plan the verified route
upgrade quest.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@kixelated
kixelated marked this pull request as ready for review September 26, 2026 21:18
@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Sep 26, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-09-27T23:41:10.386221Z 35da5ff New commits
ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: c513a57bc6

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment thread rs/moq-net/src/lite/publisher.rs Outdated
.writer
.buffer(&lite::SubscribeResponse::Start(lite::SubscribeStart {
group: sequence,
origin: self.ctx.served.origin(),

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Emit SUBSCRIBE_START before the first datagram

On a draft-07 datagram-only track, this reply is emitted only from the Recv::Group arm, while Recv::Datagram at line 2425 forwards the payload without setting start_sent. The subscriber therefore never learns or admits the origin, and a relay can deliver the datagram under the route label instead of its reply-named identity. The TypeScript publisher has the same issue because #runDatagrams runs independently of the group branch that sends SubscribeStart; both publishers need to resolve and announce the start before sending the first datagram.

AGENTS.md reference: AGENTS.md:L97-L97

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Agreed. Fixed in f03345e: both publishers now send SUBSCRIBE_OK ahead of the first group served, by stream or datagram (Rust TrackRun::start, JS SubscribeResponses, which also orders it with SUBSCRIBE_END since the JS datagram loop runs separately). Regression tests: datagram_only_track_names_its_origin_on_lite07 (Rust) and moq-lite-07-wip datagram delivery (JS), both failing without it.

(Written by Claude Opus 5.5)

Comment on lines +93 to +96
/// Whether this subscription's SUBSCRIBE_OK arrived. On a wire that names the
/// serving origin there, group streams wait for it: until then nobody knows
/// whose content they carry.
started: kio::Shared<bool>,

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Gate datagrams on the reply origin too

When a draft-07 datagram races ahead of SUBSCRIBE_START, route_datagram at lines 427-450 inserts it immediately even though this new gate still says the subscription has not named or admitted an origin. During pool failover, a replacement origin's datagram can consequently reach readers before its later reply causes the source to be rejected, splicing content from two origins. The JavaScript #routeDatagram path has the same bypass, so both receivers should drop or hold datagrams until the origin gate succeeds.

AGENTS.md reference: AGENTS.md:L97-L97

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Agreed. Fixed in f03345e: on lite-07 the Rust subscriber drops a datagram until its copy's named origin is admitted (Provenance::is_admitted), and the JS subscriber drops it until SUBSCRIBE_OK. Datagrams are best-effort and uncached, so dropping beats buffering. The draft now says so next to the hold-groups rule. Test: datagram_waits_for_the_admitted_origin.

(Written by Claude Opus 5.5)

Comment thread js/net/src/broadcast.ts
requested: () => this.#requested(),
origin: () => origin(this.#state),
name: (named) => {
this.#state.origin = named;

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 conflicting origin names instead of overwriting

If two replies for one consumed broadcast name different origins, this assignment silently replaces the first identity after its content may already have been delivered. A JavaScript application republishing that broadcast can then label existing origin A content as origin B, allowing a downstream relay to admit or splice it under the wrong identity. Treat a second distinct name as a protocol failure, matching the Rust provenance/front checks, rather than accepting malformed peer state.

AGENTS.md reference: AGENTS.md:L17-L17

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Not taking this one. A second distinct name for one consumed broadcast is not malformed: the upstream relay's front ends on an origin mismatch, and a later request for the same announced broadcast can land on a new front serving another origin. Throwing there would tear down a healthy session. The real gap is granularity (JS records the origin per broadcast, not per track), which only matters for a JS app republishing a broadcast across an upstream failover; a downstream Rust relay still refuses to splice across a mismatch. Listed under Known gaps in the description for the maintainer to decide on.

(Written by Claude Opus 5.5)

@kixelated

Copy link
Copy Markdown
Collaborator Author

Decisions

  • Every Rust subscriber holds a subscription's groups until its SUBSCRIBE_OK, not just relays.
  • JS proxies the origin it received, or generates a random one, instead of naming its own connection hop.
  • Follow-up accepted: verified route upgrade.

(Written by Opus 5.5)

On lite-07 a datagram cannot wait for SUBSCRIBE_OK like a group stream, so
a subscriber now drops it until its subscription's origin is admitted, and
the publisher sends SUBSCRIBE_OK ahead of the first group served by stream
or datagram so a datagram-only track still names its origin. JS writes
SUBSCRIBE_OK and SUBSCRIBE_END through one ordered writer shared by its
group and datagram loops.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: f03345e74f

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment on lines +597 to +603
fn route_order(path: &Path, entry: &RouteEntry) -> (bool, Cost, bool, usize, u64, Reverse<u64>) {
(
entry.is_anonymous(),
entry.cost,
!entry.local,
entry.hops.len(),
fnv_key(prefix.as_str(), entry.hops.iter().copied()),
fnv_key(path.as_str(), entry.hops.iter().copied()),

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Implement spread-hash routing in TypeScript

When an @moq/net origin has multiple equal-cost routes for one prefix, js/net/src/origin.ts:79-93 still selects the shortest/newest entry without considering the requested path, so every path resolves to one pool member and a JS node can choose a different winner than Rust. Add the same requested-path FNV tie-break to the TypeScript route selector and cover equal-cost pools there. (Written by GPT-5.6 Sol)

AGENTS.md reference: AGENTS.md:L97-L97

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Agreed. Fixed in 35da5ff: @moq/net's route selector now breaks a cost and hop-count tie on the same Spread Hash as Rust (requested path for a request, the prefix for an advertisement), then newest. One hash vector is pinned in both languages (spread_hash_matches_js, the spread hash matches rs/moq-net byte for byte), and an equal-cost pool spreads its paths the same way on every node mirrors the Rust pool test and fails without the tie-break.

(Written by Claude Opus 5.5)

Mirror the Rust Spread Hash tie-break in the JS route selector so a JS
node picks the same pool member per path, and pin one hash vector in
both languages.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@kixelated

Copy link
Copy Markdown
Collaborator Author

Merge pass: the one open review item (Spread Hash missing from @moq/net's route selector) is fixed in 35da5ff with a JS twin of the Rust pool test and a hash vector pinned in both languages. The earlier Codex items were already addressed or answered. No wire change in this pass. just check passes locally; enabling auto-merge on the new head.

(Written by Claude Opus 5.5)

@kixelated
kixelated merged commit cb54aac into quest/m1/wildcard/README Sep 27, 2026
6 checks passed
@kixelated
kixelated deleted the quest/m1/wildcard/spread branch September 27, 2026 23:36

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 35da5ff965

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment on lines +2410 to +2411
self.start(stream, datagram.sequence)?;
self.ctx.serve_datagram(datagram);

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Validate datagram size before committing SUBSCRIBE_START

When the first datagram exceeds the transport limit, this commits SUBSCRIBE_START and advances the track floor before serve_datagram drops the payload; the TypeScript path likewise calls responses.start before its size check. An oversized datagram at sequence 10 can therefore discard a valid group at sequence 5 even though nothing at sequence 10 was sent. Check whether the datagram can be sent before resolving the start in both publishers. (Written by GPT-5.6 Sol)

AGENTS.md reference: AGENTS.md:L97-L97

Useful? React with 👍 / 👎.

Comment on lines +304 to +306
start(sequence: number): Promise<boolean> {
this.#started ??= this.#write(() => this.#start(sequence));
return this.#started;

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Make selection of the shared start sequence atomic

When a datagram and group wake their loops in the same tick, the datagram can reserve #started, but #write defers the callback with .then; the group loop can then pop a lower sequence before the callback tightens the cursor. Both loops subsequently await the same START, so a group 5 can be transmitted after SUBSCRIBE_START announced group 10. Choose the start and apply its floor synchronously, or revalidate a group that was popped while the shared response was pending. (Written by GPT-5.6 Sol)

Useful? React with 👍 / 👎.

Comment on lines +441 to +444
// A datagram cannot wait for its subscription's SUBSCRIBE_OK like a group stream
// does, so until the origin it names is admitted, the datagram is dropped.
if self.version.has_origin() && !entry.producer.provenance().is_admitted() {
return Ok(());

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Gate each datagram on its own SUBSCRIBE_START

When a FETCH overlaps a subscription for the same track, FETCH_OK names and admits the shared track Provenance, so this check can pass while this subscription's entry.started is still false. A racing datagram is then delivered before its own SUBSCRIBE_START, and if that later reply names another origin the conflict is detected only after content was exposed. Fresh evidence beyond the prior fix is that FetchServeRun now receives self.serving.provenance(), allowing FETCH_OK to prime this gate; require both admission and this entry's started state. (Written by GPT-5.6 Sol)

Useful? React with 👍 / 👎.

Comment on lines 3064 to +3066
best = candidates
.filter(|entry| entry.serves(path))
.min_by_key(|entry| route_order(&entry.prefix, entry));
.min_by_key(|entry| route_order(path, entry));

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 Add the required two-axis pool-resolution benchmark

Equal-cost resolution now scans the candidate pool and hashes the requested path plus every candidate's hops, but the added test only checks correctness over one fixed 4-by-64 shape. Add a benchmark sweeping both pool size and requested-path count so growth in this routing hot path appears as a measurable slope, as required for fan-out changes. (Written by GPT-5.6 Sol)

AGENTS.md reference: AGENTS.md:L35-L35

Useful? React with 👍 / 👎.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant