feat(net): spread a prefix's pool across paths, stitch on the reply's origin - #4279
Conversation
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>
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 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".
| .writer | ||
| .buffer(&lite::SubscribeResponse::Start(lite::SubscribeStart { | ||
| group: sequence, | ||
| origin: self.ctx.served.origin(), |
There was a problem hiding this comment.
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 👍 / 👎.
There was a problem hiding this comment.
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)
| /// 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>, |
There was a problem hiding this comment.
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 👍 / 👎.
There was a problem hiding this comment.
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)
| requested: () => this.#requested(), | ||
| origin: () => origin(this.#state), | ||
| name: (named) => { | ||
| this.#state.origin = named; |
There was a problem hiding this comment.
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 👍 / 👎.
There was a problem hiding this comment.
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)
Decisions
(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>
There was a problem hiding this comment.
💡 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".
| 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()), |
There was a problem hiding this comment.
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 👍 / 👎.
There was a problem hiding this comment.
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>
|
Merge pass: the one open review item (Spread Hash missing from (Written by Claude Opus 5.5) |
There was a problem hiding this comment.
💡 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".
| self.start(stream, datagram.sequence)?; | ||
| self.ctx.serve_datagram(datagram); |
There was a problem hiding this comment.
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 👍 / 👎.
| start(sequence: number): Promise<boolean> { | ||
| this.#started ??= this.#write(() => this.#start(sequence)); | ||
| return this.#started; |
There was a problem hiding this comment.
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 👍 / 👎.
| // 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(()); |
There was a problem hiding this comment.
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 👍 / 👎.
| best = candidates | ||
| .filter(|entry| entry.serves(path)) | ||
| .min_by_key(|entry| route_order(&entry.prefix, entry)); | ||
| .min_by_key(|entry| route_order(path, entry)); |
There was a problem hiding this comment.
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 👍 / 👎.
Completes the
quest/m1/wildcard/spread.mdquest, 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_orderhashed 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
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.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.Pin::Stay): which member the new route leads to is unknown until it replies.@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.Origin(hold groups / drop datagrams until SUBSCRIBE_OK),Originon SUBSCRIBE_OK, the new FETCH_OK message, changelog entries.doc/bin/relay/cluster.mdfailover and tie-break paragraphs.spread.md, sizes the line README[S], and adds Cluster origin reply (moq-transport downstreams) and Verified route upgrade to m1.Impact
moq-lite-07-wip) SUBSCRIBE_OK appendsOrigin (i), and a new FETCH_OK{ Origin (i) }precedes a fetch's frames. TRACK_INFO is unchanged. No change to published versions.track::Provenanceandtrack::Info::names_originare crate-private;lite::SubscribeStartgainsorigin: Hopandlite::FetchOkis new.lite.SubscribeStartgainsorigin(defaults to unknown) andlite.FetchOkis new; the internal broadcast wire gainsorigin()/name().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 twinan equal-cost pool spreads its paths the same way on every nodedoes 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.datagram_only_track_names_its_origin_on_lite07(Rust, mock sessions) andmoq-lite-07-wip datagram delivery(JS integration) fail without SUBSCRIBE_OK before the first datagram;datagram_waits_for_the_admitted_originpins the Rust subscriber drop.Decisions
Settled by the maintainer:
Originrides 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
Known gaps
@moq/netrecords 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.Checks
just check,just drafts check, andjust test interop --allpass.(Written by Claude Opus 5.5)
🤖 Generated with Claude Code