Skip to content

Draft-20 compliance: relay backlog - #114

Merged
floatdrop merged 23 commits into
draft-20from
draft20-backlog-relay
Sep 27, 2026
Merged

floatdrop merged 23 commits into
draft-20from
draft20-backlog-relay

Conversation

@floatdrop

Copy link
Copy Markdown
Owner

The relay part of the draft-20 compliance backlog in STATUS.md, one commit per item. Every fix comes with a regression test that was seen failing before the fix, or failing under a mutant that removes the behaviour.

Commits

Subscribe, publish and forwarding

  • 0bc25d5 INCLUDE_PROPERTIES=0 survives a REQUEST_UPDATE (§10.9, §10.2.21).
  • 82bb1b9 A SUBSCRIBE_TRACKS holder isn't forwarded a track its own SUBSCRIBE is establishing (§6.1, §10.20).
  • 5e853ad A PUBLISH_SKIPPED holds across a prefix update that moves away and back (§6.1).
  • 023d79e Two forwards racing for one track send one PUBLISH.
  • 58e40d4 A client can SUBSCRIBE and FETCH a track it publishes (§5.1), including through late-publisher and missing-publisher SUBSCRIBEs (§9.5).
  • a97626e RENDEZVOUS_TIMEOUT holds a SUBSCRIBE for a publisher (§10.2.6).

FETCH and fill

  • b1128f1 FETCH_OK End Of Track is set at the END_OF_TRACK Object (§10.14).
  • ada7018 A cancelled FETCH resets its request and data streams (§5.2).
  • 4de6761 Objects stitched from an upstream FETCH are bounded by its MAX_CACHE_DURATION (§12.3).
  • 2b72d55 A malformed track cancels its FETCHes and resets their fetch streams (§2.4.2).
  • 8e8085b A fill ends at the Largest Object the subscriber was told (§5.1.3).
  • 2288a1f A refused FETCH's request stream is reset with its data stream's code (§3.3.3).

Data plane

  • 6e99b0b Replay streams carry the first Object's delivery-timeout override (§8, §12.1, §12.2).
  • 9518daa A merged Subgroup that a clean survivor does not cover is reset (§11.4.3).
  • 23c67e9 A merged Subgroup's FIN is judged from the ledger's lowest Object (§11.4.3).

Errors, tokens, GOAWAY, shutdown

  • 5350f87 Malformed Track Properties are refused with INTERNAL_ERROR, not MALFORMED_TRACK, outside FETCH (§10.6.2).
  • 18b4708 A REQUEST_UPDATE's tokens are verified on SUBSCRIBE, FETCH and PUBLISH_NAMESPACE (§10.2.2, §10.9.1).
  • 6480ab3 TRACK_STATUS is answered for any track SUBSCRIBE would accept (§10.15).
  • e3a70d4 No request is initiated on a session with a GOAWAY in either direction (§10.4), and GOING_AWAY is the refusal code.
  • 9ac663f Shutdown closes with GOAWAY_TIMEOUT only after a GOAWAY ran out, NO_ERROR otherwise (§3.5).

Housekeeping

  • 9d31287 The benchmark baseline is refreshed. FetchFromCache went from 115 to 134 allocs/op: 6 from cbbdf96, 13 from this batch. ControlRoundTrip went from 46 to 47.
  • 90cba9a and a5b043d are test-only: the close-code test's flake fix and a line wrap.

Decisions and interpretations

Self-subscription

  • The relay SUBSCRIBEs or stitch-FETCHes from the requester's own session like any other (§5.1, "identical"; FETCH follows SUBSCRIBE's matching rules, §9.5).
  • Deviation. While a SUBSCRIBE (or stitch FETCH) for a track to a session is still pending, a second one from that session gets DOES_NOT_EXIST (or its hole is marked unknown).
    • This stops the loop a relay peer creates by routing requests straight back (§6.2 has no loop protection); without the guard the tests saw 11 hops.
    • A genuine concurrent self-subscription looks the same, so it is declined too. This is recorded in STATUS.md.

RENDEZVOUS_TIMEOUT

  • The hold is capped by the new Config.MaxRendezvousTimeout, default 30s. A negative value holds nothing.
  • These count as "no publisher yet": no match, or DOES_NOT_EXIST, TIMEOUT, or a draining publisher. Any other refusal ends the hold with its error.
    • Behaviour change outside the hold: when several candidates refuse, the non-"no publisher" error now wins whatever the order; before, the last answer won.
  • The remaining budget is forwarded only to other relays (discovery remotes), not to local publishers.
  • A remote relay's hold is cut short when a publisher arrives locally, and that relay is asked again if the local one refuses.
  • A publisher that already answered isn't asked again during the hold.
  • Nothing is held on a session the relay sent GOAWAY to.

Error codes

  • A refused FETCH resets its request stream with MALFORMED_TRACK. On the §2.5.1 unknown-Mandatory-Property path it uses INTERNAL_ERROR, the same code as the data stream.
  • A token failure on a FETCH REQUEST_UPDATE resets with EXPIRED_AUTH_TOKEN or CANCELLED.

Other defaults

  • End Of Track comes only from the END_OF_TRACK Object.
  • Nothing is done after a FIN.
  • Non-consecutive Subgroups reset.
  • The delivery timeout falls back to the Track-level value.
  • A stitched MAX_CACHE_DURATION of 0 means no limit.
  • A fill for a forwarded PUBLISH uses the registration snapshot.

API changes

  • relay.Config.MaxRendezvousTimeout (new).
  • session.Session.GoawaySent() (new).

Known gaps (in STATUS.md)

  • Filters are not aggregated upstream (§6.3.1 SHOULD). Deferred.
  • Cache entries age from when an Object was read whole, not from its first byte (§12.3). Deferred.
  • Found in this batch, not fixed: cancelling Start's ctx ends each session's handler without closing the session, and drops it from the set Stop closes. Stop can then wait without bound on an upstream SUBSCRIBE's stream reader.
    • Start's doc says the cancel "terminates live sessions".
    • The new cross-relay tests stop their relays before t.Context ends to avoid it.
    • Fixing it changes a documented contract, so it's left for a decision.
  • TRACK_STATUS asymmetries, drainStraggler ignoring Stop's ctx, and a draining pool relay yielding DOES_NOT_EXIST are noted in the backlog.

Verification

  • go test ./... and golangci-lint run pass.
  • go test -race ./pkg/relay/... passes; the final rendezvous change was run five times.
  • make bench-quick shows allocs/op matching the refreshed baseline.

🤖 Generated with Claude Code

floatdrop and others added 23 commits September 27, 2026 16:33
… §10.2.21)

"If a parameter previously set on the request is not present in
REQUEST_UPDATE, its value remains unchanged" (§10.9), and
INCLUDE_PROPERTIES cannot appear in a REQUEST_UPDATE (§10.2.21).
installSubscribeParams nonetheless set it on every call, defaulting an
absent one to 1. So any REQUEST_UPDATE turned a subscription's
INCLUDE_PROPERTIES=0 back off. After that, its subgroups and datagrams
dropped the priority it had been writing inline, and the subscriber,
which had no Track Properties, resolved the wrong default Publisher
Priority (§12.4).

It is now set only when present. A new subscription's zero value
already means "include".

Verified red first: TestIncludeProperties_SurvivesUpdate.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…10.15)

"The receiver of a TRACK_STATUS message treats it identically as if it
had received a SUBSCRIBE message, except it does not create downstream
subscription state or send any Objects." A PUBLISHed track with no
Track Properties and no Objects yet is one SUBSCRIBE accepts, since it
has an established upstream. TRACK_STATUS answered it DOES_NOT_EXIST,
because it looked only for Properties or a LARGEST_OBJECT watermark.

It now also answers OK for an entry with an established upstream,
SUBSCRIBE's own condition.

Two remaining differences are recorded in the STATUS.md backlog:
  - a namespace-only answer, where SUBSCRIBE would go upstream;
  - a leftover entry that has metadata but no upstream.

Verified red first: TestTrackStatus_AnswersPublishedTrackWithNothingYet.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
GOAWAY_TIMEOUT means "the peer took too long to close the session in
response to a GOAWAY" (§3.5). Relay.Stop, and drainStraggler for a
session that registered after Stop began, force-closed every remaining
session with it. That included sessions that were never sent a GOAWAY
(GoawayTimeout 0, or a SendGoaway that failed) and those cut short by
Stop's ctx before the grace period, which the GOAWAY Timeout allows
(§10.4: "is a hint").

Both paths now close with GOAWAY_TIMEOUT only when the relay sent the
session a GOAWAY and the grace period ran out, and with NO_ERROR
otherwise. The test listener gains a wrap hook, so the tests can
observe the code the relay closes a session with.

This is the second half of the STATUS item. The first, not initiating
requests toward a peer sent GOAWAY, comes separately. A related
pre-existing gap is recorded in STATUS.md: drainStraggler ignores
Stop's ctx.

Verified red first: in TestStopClosesWithGoawayTimeoutOnlyAfterGoaway,
the no-GOAWAY and ctx-cancel cases on the unpatched code. The
TestRelay_addSessionDrainsStraggler cases were checked against a mutant
that always uses GOAWAY_TIMEOUT.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…eplay streams (§8, §12.1, §12.2)

"For each type of timeout, the publisher's value is the Object Property
when present on the first object of the subgroup, and the Track Property
otherwise" (§8); as an Object Property on the first object in a
subgroup, OBJECT/SUBGROUP_DELIVERY_TIMEOUT "overrides the Track-level
value for that subgroup" (§12.1, §12.2).

The session applies the override only on a stream with the FIRST_OBJECT
bit, which is right: a replay stream's first Object is not the
Subgroup's. But the relay handed every downstream stream the Track-level
value alone, so a replay stream kept no override: a subscriber joining
after the first Object, or a stream reopened after a §11.4.3 gap.

When the fanout claims the Subgroup's first Object, it now resolves the
Track timeouts with that Object's Properties once, stores them in the
Subgroup's writer set, and passes an immutable pointer on every later
forwarded Object. Each writer adopts it before opening a stream. When
the relay never saw the first Object, the Track-level value stays.
Visible to peers: the override now also reaches a downstream stream
without the PROPERTIES bit (INCLUDE_PROPERTIES=0), since the publisher's
value comes from the Object as received. The per-object forwarding moves
into subgroupWriterSet.forward to keep runFanout under the gocyclo limit.

Benchmarks (bench-quick, 4 runs each side): allocs/op unchanged
(Fanout1to1 5, Fanout1toN 129-130, FetchFromCache 121). The value lives
in the set, so there is no per-Subgroup allocation; fwdObject grows from
56 to 64 bytes.

Verified red first: TestFanout_ReplayStreamKeepsFirstObjectDeliveryTimeout
(joiner and gap reopen) both failed with all 6 stalled Objects delivered.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…§11.4.3)

"If a sender closes the stream before delivering all such objects to the
QUIC stream, it MUST reset the stream" (§11.4.3). When several upstream
streams feed one Subgroup (§9.3), the merged downstream streams FINned
if any contributor ended cleanly. Suppose A delivers 0 and 1, B opens a
replay stream and delivers 4, then A resets and B FINs. Objects 2 and 3
were never forwarded, but the subscriber's last stream FINned.

Each contributor now records where its clean end vouches from: 0 for a
stream with FIRST_OBJECT (§11.4.2), else its first Object's ID. A replay
stream that ended before any Object vouches for nothing. The merged
streams FIN only if the lowest such start is at or below the lowest
Object the set forwarded, or the set's forwarded run is unbroken up to
it. Otherwise a reset contributor may have held Objects nobody
forwarded, and they reset with CANCELLED. With no clean contributor the
reset code is still the reset contributor's.

Decided: a run with non-consecutive IDs (a Group split across Subgroups)
is broken, so it resets. So does one whose consecutive IDs were
forwarded out of order. Both are conservative: the relay cannot tell a
skipped ID from one that does not exist.

Assumed (deviation from the planned rule): Objects below the lowest one
the set forwarded count as before the Start Location, as they already
do for a joiner, so the set does not need to have forwarded the
Subgroup's first Object. Requiring that would reset a lone replay
upstream's FIN (TestFanout_FirstObjectBitNotInvented). The one case this
lets through: a FIRST_OBJECT contributor resets before forwarding any
Object, and a replay starting above the true first then FINs. If that
case must reset, require a forwarded first Object whenever a
FIRST_OBJECT contributor reset.

Left in the backlog: the run is per writer set, and the set is dropped
with its last contributor. A clean replay upstream that arrives only
after every earlier one left is still judged against its own set alone.
Fixing that means consulting the dedup ledger's LowestForwarded, which
would also reset a sequential failover that did continue the run.

Benchmarks (bench-quick): allocs/op unchanged (Fanout1to1 5, Fanout1toN
129-130, FetchFromCache 121).

Verified red first: TestFanout_MultiPublisher_SurvivorFINsOnlyWhatItCovers
"replay starts past undelivered Objects" and "Object IDs not
consecutive" ended with a FIN where a reset is wanted. "replay continues
the run" FINs before and after, and
TestFanout_MultiPublisher_FailoverContinuesFromSurvivor stays green.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
….14)

§10.14: End Of Track is "1 if all Objects have been published on this
Track, and the End Location is the final Object in the Track, 0 if not".
The relay never set it, although the TrackEntry ledger already records
where an END_OF_TRACK Object ended the track (§2.4.2).

TrackEntry.TrackEnd now reads that record under its lock, and
handleFetch sets End Of Track when the FETCH_OK End Location is it. The
END_OF_TRACK Object raised the watermark, so a FETCH running past the
track's end is capped to it by capFetchEndLocation and says so too.

The track end is learned only from the END_OF_TRACK Object; an upstream
FETCH_OK's End Of Track is not recorded, which the STATUS.md 10.14 row
now says.

Verified red first: TestFetch_OKEndOfTrack failed with End Of Track
false for both the whole-track FETCH and the one ending at the
END_OF_TRACK Object. TestFetch_OKEndLocationCappedToWatermark now also
asserts End Of Track stays false when End Location is merely Largest
Object; it fails against an `endLocation == largest` variant.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
§5.2: "The Publisher can remove fetch state as soon as it has received a
STOP_SENDING. It MUST reset the bidi request stream and unidirectional
data stream associated with the FETCH." handleFetch served the response
on the session context and nothing watched the request stream, so a
requester that cancelled while the relay waited on an upstream for a
hole had its data stream written and FINed once FILL_TIMEOUT ran out.

serveFetchObjects now runs streamFetchRange under a context cancelled
with errRequestCancelled when the request stream's send Context ends
(the requester's STOP_SENDING), mirroring the fill-cancel pattern in
handler_fill.go. streamFetchRange already resets the data stream with
CANCELLED for that cause, and the upstream FETCH it drives is cancelled
with it. When the response did not complete and the requester has
cancelled, the request stream is also cancelled in both directions with
CANCELLED (§3.3.3: "RESET_STREAM for a direction they are sending and
STOP_SENDING for a direction they are receiving"). A cancel after the
data stream FINed needs nothing more.

sessiontest drops reset codes, so the relay test harness gains a
recorder (pipeListener.resetsFor, resetsOn) that wraps a relay-side conn
and reports the codes the relay resets its streams with.

Verified red first: TestFetch_CancelResetsStreams saw neither the data
stream nor the request stream reset within 500ms of the cancel.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…CACHE_DURATION (§12.3)

§12.3: "If present, the relay MUST NOT start forwarding any individual
Object received through this subscription or fetch after the specified
number of milliseconds has elapsed since the beginning of the Object was
received." Objects the relay read from an upstream FETCH to fill cache
holes carried a zero ReceivedAt and no MAX_CACHE_DURATION, and
ObjectCache.Expired treated a zero ReceivedAt as never expiring, so
they were forwarded however long the upstream held its stream before
the FIN.

fetchUpstreamRange now stamps each stitched Object when it is read and
gives it the upstream FETCH_OK's MAX_CACHE_DURATION. A new exported
field, cache.CachedObject.Stitched, marks it (the struct does not grow:
CachePut stays 128 B/op, 1 alloc). Expired bounds a stitched Object by
its own positive MAX_CACHE_DURATION only, not the relay's cache TTL,
since it is passed through rather than stored. A present 0, which the
cache reads as "never serve from the cache", sets no limit on a stitched
Object, as on the live path (interpretation). An expired one becomes an
End of Unknown Range, as a cached one does.

The other half of that STATUS.md item remains: Objects age from when
they were read whole, not from their beginning. The 12.3 row now says
so and is PARTIAL.

Verified red first: TestRelay_MaxCacheDurationBoundsStitchedObject/expired
got Object {0,1} where it wants an End of Unknown Range. The "zero" case
fails against a variant that applies a present 0 as the cache does.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ck (§2.4.2)

§2.4.2: "When a subscriber detects a Malformed Track, it MUST cancel any
corresponding subscription or fetches for that Track from that publisher
[...] If a relay detects a Malformed Track, it MUST immediately
terminate downstream subscriptions with PUBLISH_DONE and reset any fetch
streams with Status Code MALFORMED_TRACK." endMalformedTrack cancelled
only upstream subscriptions. An upstream FETCH filling a cache hole kept
waiting until FILL_TIMEOUT, and a downstream fetch or fill fetch stream
served from the cache kept being written.

TrackEntry gains AddFetch and CancelFetches. The FETCH response
(serveFetchObjects) and fill fetch stream (maybeServeFill) register the
cancel func of the context they already own. endMalformedTrack cancels
them all with a malformedTrackCause that wraps the detection error and
names the source session. ctxResetCode maps that cause to
MALFORMED_TRACK, so each downstream fetch stream is reset with it.

Every upstream FETCH runs under a downstream fetch stream's context, so
this cancels it too. Its data stream gets STOP_SENDING MALFORMED_TRACK
when the upstream is the source, else CANCELLED; fr.Close cancels its
request stream with CANCELLED, as on the self-detected path. This
registers the downstream streams rather than each upstream FETCH, as
first planned: one registration then covers the cache-only and fill
streams too.

The relay test harness's reset recorder now tells stream kinds apart by
the FETCH_HEADER type byte, and also reports STOP_SENDING on fetch
streams the relay reads.

Left open, recorded in STATUS.md: a FETCH whose data stream is reset
with MALFORMED_TRACK leaves its request stream open and unread (already
so on the §2.5.1 refusal path). Which code to reset it with is a
wire-visible choice.

Cost: BenchmarkFetchFromCache goes from 132 to 134 allocs/op, per FETCH
request.

Verified red first, against the pre-fix production code:
TestRelay_MalformedTrackCancelsUpstreamFetch (the upstream FETCH was not
cancelled within 1s), TestRelay_MalformedTrackResetsCachedFetch and
TestRelay_MalformedTrackResetsFillStream (no fetch stream reset within
1s). The first also fails if the upstream data stream is cancelled with
CANCELLED.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…PUBLISH_NAMESPACE (§10.2.2, §10.9.1)

§10.2.2: the AUTHORIZATION TOKEN parameter "MAY appear in a PUBLISH,
SUBSCRIBE, REQUEST_UPDATE, ... PUBLISH_NAMESPACE, TRACK_STATUS or FETCH
message. This parameter conveys information to authorize the sender to
perform the operation carrying the parameter."

The relay resolved a REQUEST_UPDATE's tokens through the alias cache on
every request but ran them through the TokenVerifier only on
SUBSCRIBE_NAMESPACE and SUBSCRIBE_TRACKS. On SUBSCRIBE (including the
subscription a forwarded PUBLISH opens), FETCH and PUBLISH_NAMESPACE a
token the verifier would deny was accepted with REQUEST_OK.

Each of those now verifies the update's tokens, and a denial answers
REQUEST_ERROR with the verifier's code, then does what §10.9.1 requires
of a failed update:

  "When a REQUEST_UPDATE is unsuccessful, the publisher MUST also
  terminate the subscription by sending a PUBLISH_DONE with error code
  UPDATE_FAILED. When a REQUEST_UPDATE fails for a FETCH, the publisher
  MUST reset the FETCH data stream. When a REQUEST_UPDATE fails for a
  SUBSCRIBE_NAMESPACE, SUBSCRIBE_TRACKS or PUBLISH_NAMESPACE, the
  responder MUST close the bidi stream"

- SUBSCRIBE: PUBLISH_DONE UPDATE_FAILED, as an invalid Range Filter
  already did.
- FETCH: the request ends (REQUEST_ERROR and FIN) and the fetch data
  stream is reset, so streamFetchRange now returns the stream. The relay
  FINs that stream before it reads any update, so the reset only aborts
  what the requester has not acknowledged. Its code is EXPIRED_AUTH_TOKEN
  when the verifier said so, else CANCELLED (§3.3.4 has no reset code
  for other denials). This is a judgement call.
- PUBLISH_NAMESPACE: REQUEST_ERROR and FIN close the stream, which
  withdraws the namespace. It now passes its own update func, so
  serveNamespaceFollowups loses its write parameter and enqueueReply.

The shared check is refuseUpdateTokens, which refuseUpdate now uses too.

Verified red first: all four TestRequestUpdate_TokenVerified_* tests
failed on the unpatched relay with "want *RequestRejectedError, got
<nil>" (REQUEST_OK answered the denied update). Not covered: the FETCH
data-stream reset. The in-process transport keeps a FIN over a later
reset, and FaultFunc has no CancelWrite op, so no test fails without the
out.Cancel.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…RROR, not MALFORMED_TRACK (§10.6.2)

§10.6.2 defines the code for one request type only:

  "MALFORMED_TRACK: In response to a FETCH, a relay publisher detected
  the track was malformed (see Section 2.4.2)."

The session and relay nevertheless answered REQUEST_ERROR MALFORMED_TRACK
outside a FETCH:

- a PUBLISH whose Track Properties do not parse (§12.7 makes the track
  malformed), through session.TrackPropertiesRejectCode in
  Request.AcceptPublish and the relay's handlePublish;
- a downstream SUBSCRIBE whose upstream SUBSCRIBE_OK carried such Track
  Properties (upstreamRejection);
- a downstream SUBSCRIBE whose upstream answered REQUEST_ERROR
  MALFORMED_TRACK, which upstreamRejection passed through.

All three now answer INTERNAL_ERROR ("An implementation specific or
generic error occurred"). TrackPropertiesRejectCode, which is exported,
returns RequestInternalError where it returned RequestMalformedTrack, and
MALFORMED_TRACK moves into upstreamRejection's INTERNAL_ERROR case list.
UNSUPPORTED_EXTENSION for an unknown Mandatory Track Property (§2.5.1) is
unchanged. MALFORMED_TRACK is kept where it is defined: PUBLISH_DONE
(§10.12) and the stream reset code (§3.3.4), which is how the relay
reports a malformed track after FETCH_OK. The doc comments that called
MALFORMED_TRACK this package's interpretation are updated.

Verified red first. I flipped the four existing expectations to
INTERNAL_ERROR and each failed on the unpatched code with code 0x12:
TestAcceptPublishTrackPropertiesRejected/malformed (session),
TestRelay_PublishTrackPropertiesRejected/malformed,
TestRelay_UpstreamSubscribeOKTrackPropertiesRejected/malformed and
TestSubscribe_UpstreamRejects_PropagatesRejection/0x12.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…away and back (§6.1)

§6.1: "The Publisher MUST NOT send a PUBLISH for a Track for a given
SUBSCRIBE_TRACKS after PUBLISH_SKIPPED has been sent, scoped to a single
PUBLISH. If, for example, the publisher disconnects from a relay and
later reconnects and sends a new PUBLISH, the relay MAY send the new
PUBLISH downstream."

Nothing recorded the skip. When a SUBSCRIBE_TRACKS's
TRACK_NAMESPACE_PREFIX moved away from a skipped track and back,
tracksUpdate saw the track newly match and offered it again. That meant
a second PUBLISH_SKIPPED, or a PUBLISH once credit returned, for the
same upstream PUBLISH.

Each track entry now carries an upstream epoch, set on every AddUpstream
from a process-wide counter, so it only grows and never repeats on an
entry recreated for the same track. When a PUBLISH_SKIPPED is queued,
the SubscriberEntry records the epoch it was sent at, and ClaimForward
refuses the track while the epoch forwardTrack read is not newer. A
forward decided on a stale read is therefore covered by a later skip. A
new upstream PUBLISH or SUBSCRIBE advances the epoch, so its PUBLISH may
be offered (and skipped) again, as TestPublishSkipped_NotStickyAcrossRePublish
still requires. The skip is recorded before the claim is released, so a
concurrent forward of the track cannot slip in between.

STATUS.md rows 10.20 and 10.21 now describe the scope of the skip.

Verified red first: TestPublishSkipped_StickyAcrossPrefixMoveAndBack
failed on the unpatched relay with "the prefix moved back to the skipped
track: unexpected *message.PublishSkipped &{TrackNamespaceSuffix:/cam7
TrackName:rtp}".

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…n SUBSCRIBE is establishing (§6.1, §10.20)

forwardTrack sends no PUBLISH for a track the subscriber "publishes or
already receives": §6.1 excludes "tracks published by the subscriber",
and one it SUBSCRIBEd to arrives on that subscription. The receives
half checked HasDownstreamOn(h.sess), which stays false between a
SUBSCRIBE establishing the track's upstream (AddUpstream) and
registering its downstream (AddDownstreamSnapshotLargest). A forward
in that window sent the session a PUBLISH for the track it had just
SUBSCRIBEd to. That is the load-only flake in
TestSubscribeTracks_ForwardsTrackGainedBySubscribe: handleSubscribeTracks
replies REQUEST_OK before it forwards the existing tracks, and the
forward found the upstream the concurrent SUBSCRIBE had made.

Each sessionHandler now counts, per track key, its SUBSCRIBEs in
flight. handleSubscribe marks the key before it looks for or
establishes an upstream, and clears it right after the registration
loop (also on every early return). forwardTrack checks the set before
HasDownstreamOn: the SUBSCRIBE registers its downstream before it
clears the mark, so a forward sees one or the other. The check lives in
forwardTrack, so it covers every path that forwards to this session:
the SUBSCRIBE_TRACKS existing-tracks loop, a TRACK_NAMESPACE_PREFIX
update (tracksUpdate), and a PUBLISH or an upstream-creating SUBSCRIBE
from another session (forwardToTrackSubscribers).

A forward held back this way is recorded. When the last SUBSCRIBE for
the key ends, it is offered again. forwardTrack declines it when a
SUBSCRIBE registered its downstream, and sends it when the SUBSCRIBE
failed, since the session then receives the track no other way
(§10.20). Otherwise the fix would have traded a duplicate PUBLISH for a
lost one.

Test hook testHookBeforeDownstreamRegistered (with its setter in
export_test.go, as for testHookAfterAliasRegistered) parks a SUBSCRIBE
in the window.

Verified red first. Before the fix, each of
TestSubscribeTracks_OwnSubscribeInFlight_{ExistingTracks,PrefixUpdate,NewPublish}
failed with "... while the session's SUBSCRIBE is in flight: forwarded
PUBLISH for own-window-...". TestSubscribeTracks_OwnSubscribeFails_HeldForwardSent,
added for the replay after review, failed with "no PUBLISH forwarded to
the SUBSCRIBE_TRACKS holder" before the replay existed.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Same machine and toolchain as the run it replaces (Apple M3 Max,
darwin/arm64, go1.27), so the two are directly comparable. The file
also gains BenchmarkCheckObjectProperties, which it lacked.

Two benchmarks moved on allocs/op, both on per-request paths. Every
per-object benchmark (fanout, subgroup codec and throughput, cache
put/get) is unchanged.

    FetchFromCache    115 -> 134
    ControlRoundTrip   46 ->  47

FetchFromCache, per FETCH request:
  - 115 -> 121, cbbdf96 (§5.1.1, §5.1.3.1): streams are reset when their
    request is cancelled. The fill path gained a cancel-with-cause
    context watching the request stream, found by bisecting a5fad8f..
  - 121 -> 132, the cancelled-FETCH reset (§5.2: the publisher "MUST
    reset the bidi request stream and unidirectional data stream"):
    serveFetchObjects runs under a context.WithCancelCause, with a
    context.AfterFunc on the request stream's context. The AfterFuncs
    and derived contexts under it now register as children of that new
    context, with its own children map, rather than of the session's.
  - 132 -> 134, registering each fetch stream on its track so a
    malformed track can reset it (§2.4.2).

ControlRoundTrip: +1 from releasing a Subscription's Track Alias on
Close (§11.1). Each iteration now registers the alias anew, which
allocates the channel that wakes awaitInboundTrack.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…n either direction (§10.4)

§10.4 has two rules. An endpoint that received a GOAWAY "SHOULD NOT
initiate new requests to the peer". The sender of one "SHOULD avoid
initiating requests unless required by migration". The relay honoured
only the first. During its own drain it went on SUBSCRIBEing, FETCHing
and PUBLISHing to peers it had sent GOAWAY to.

New exported Session.GoawaySent. The relay's check becomes goingAway,
true for a GOAWAY either way, at every request the relay initiates:
upstream SUBSCRIBE, the FETCH upstream pick, forwarded PUBLISH and the
upstream pool. The relay never needs such a request for migration, and
it does not follow a New Session URI.

A downstream SUBSCRIBE refused because its publisher is going away now
gets REQUEST_ERROR GOING_AWAY with a jittered ~1s Retry Interval, not
DOES_NOT_EXIST. GOING_AWAY is "The endpoint has received a GOAWAY and
MAY reject new requests" (§10.6.2), and during the relay's drain it
"has sent or received a GOAWAY" (§3.3.4).

excessiveLoadRetryInterval generalizes to retryIntervalAfter. A related
case is recorded in STATUS.md: a draining upstream relay found through
the pool still yields DOES_NOT_EXIST.

Verified red first: TestRelay_NoRequestsToPeerSentGoaway (the relay
sent SUBSCRIBE to a publisher it had sent GOAWAY) and
TestRelay_NoUpstreamSubscribeToGoingAwayPublisher (DOES_NOT_EXIST, not
GOING_AWAY).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…(§5.1.3)

A Next Object filter "coupled with an open-ended fill range, which the
publisher will end at Largest Object" delivers each Object exactly once,
and the subscriber "learns the Largest Object from the LARGEST_OBJECT
parameter in SUBSCRIBE_OK or REQUEST_UPDATE_OK" (§5.1.3).
maybeServeFill re-read the track's Largest Object when it evaluated the
fill, after that response was sent. An Object arriving in between went
out both live and on the fill. When SUBSCRIBE_OK had reported no Largest
at all, a fill opened anyway.

maybeServeFill now takes the Largest Object from its caller:
  - SUBSCRIBE: the registration snapshot it put in SUBSCRIBE_OK;
  - REQUEST_UPDATE: the value it put in REQUEST_UPDATE_OK;
  - a forwarded PUBLISH (§10.20.1): the registration snapshot, which the
    live filter is anchored on. The PUBLISH's own LARGEST_OBJECT was read
    earlier and can be older; ending there would leave the Objects in
    between neither live nor filled.

testHookBeforeFill lets a test publish an Object in the window.

Verified red first: in TestFill_EndsAtTheLargestObjectReported, the two
SUBSCRIBE_OK cases on the unpatched code. All three cases, REQUEST_UPDATE
included, fail against a mutant that re-reads Largest inside
maybeServeFill.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…so two forwards send one PUBLISH

forwardTrack sends "At most one PUBLISH per track" to a SUBSCRIBE_TRACKS
holder. That is relay policy: §6.1 only excludes "tracks published by
the subscriber". It checked HasDownstreamOn(h.sess) before
SubscriberEntry.ClaimForward. So two forwards of one track to one
holder, say the existing-tracks loop and a second publisher's
forwardToTrackSubscribers, could both pass the check. One then claimed,
opened the PUBLISH, registered its downstream and released the claim.
The other claimed next without looking again and sent a second PUBLISH.

HasDownstreamOn now runs after the claim, and the claim is released if
the check fails. serveForwardedPublish releases the claim only after
AddDownstreamSnapshotLargest, so whichever forward claims second sees
the first one's downstream.

The in-flight SUBSCRIBE check (holdForward) stays before the claim. A
held forward records its entry there. If it held the claim at that
point, the replay run when the SUBSCRIBE ends could find the claim
taken and drop the PUBLISH. The PUBLISH_SKIPPED epoch check in
ClaimForward is unchanged.

Test hook testHookBeforeForwardClaim (setter in export_test.go) parks a
forward just before its claim.

Verified red first: TestSubscribeTracks_ConcurrentForwardsSendOnePublish
failed 3/3 before the reorder with "the held forward of a track already
forwarded: forwarded PUBLISH for concurrent-forwards". The test holds
the existing-tracks forward, lets a second publisher's forward complete,
then releases the first.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ject (§11.4.3)

"If a sender closes the stream before delivering all such objects to the
QUIC stream, it MUST reset the stream" (§11.4.3). outcome measured a
clean replay contributor against its writer set's own run of forwarded
Object IDs, and that set is dropped when its last contributor leaves.
When contributors took turns, Objects went missing unnoticed. A
delivers 0 and 1 and resets, which releases the set. B then arrives with
a replay from 4 and FINs. The new set's run starts at 4, so the
subscriber's stream FINned although 2 and 3 were never delivered.

outcome now measures coverage from subgroupWriterSet.lowest. That value
is seeded from the dedup ledger's LowestForwarded, which outlives the
set. The streams FIN only if the clean start is at or below it, or this
set's unbroken run starts at it and reaches the clean start. It can only
reset more than before, never FIN more, since lowest <= runLo.

Trade-off, decided by the user ("catch it, reset more"): the earlier
set's run is forgotten, so a replay after A left that does continue the
run (A 0,1 resets, then B 2,3) now also resets. So does a replay that
overlaps the departed range without starting at its lowest Object (A
0..5 resets, B 3..7). Both are labelled a Deviation from §11.4.3's "MUST
close the stream with a FIN" in outcome's comment. A replay covering
from the lowest Object (B 0..3, with 0 and 1 redundant) still FINs.
Checking each ID below the clean start against the ledger's per-Group
Object set would avoid the trade-off; it is not done here.

No existing test expectation flipped. The table in
TestFanout_MultiPublisher_SurvivorFINsOnlyWhatItCovers gains an
aLeavesFirst dimension. Its payloads now derive from the Object ID, so a
redundant copy matches the original (§9.1).

Benchmarks (bench-quick): allocs/op unchanged against the refreshed
baseline (Fanout1to1 5, Fanout1toN 130, FetchFromCache 134, session
2/47/2).

Verified red first: "replay after A left starts past undelivered
Objects" and "replay after A left continues the run" both FINned where
a reset is wanted. "replay after A left covers from the lowest Object"
is the positive control. It passes before and after, and it fails under
a mutation that makes a rebuilt set always reset.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
TestStopClosesWithGoawayTimeoutOnlyAfterGoaway called Stop as soon as
the client's handshake returned. The relay can finish the handshake
before it registers the session, so Stop's snapshot sometimes missed
it. The session then became a straggler, drained on its own grace
period and ignoring Stop's ctx. The ctx-cancel case then closed with
GOAWAY_TIMEOUT after the full 5s, which failed about one run in ten
under load.

A TRACK_STATUS round trip first guarantees registration, because the
relay registers a session before serving its requests. The test now
passes 40 runs in a row. The straggler's own ctx gap is already in the
STATUS.md backlog.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
"An endpoint MAY SUBSCRIBE to a Track it is publishing ... Such
self-subscriptions are identical to subscriptions initiated by other
endpoints" (§5.1). The relay skipped the requester's own session as an
upstream everywhere: subscribeUpstream, the §9.5 missing- and late-publisher
SUBSCRIBEs (which also skipped a publisher receiving the track), and the
stitch FETCH (§9.5: FETCH uses SUBSCRIBE's matching rules). All skips go.

With them goes the relay's only protection against a relay peer that routes
the relay's SUBSCRIBE or stitch FETCH straight back to it on the same session
(§6.2: no loop protection), which would recurse without bound. Two guards,
both scoped to the requester's own session so concurrent subscribers on
other sessions are unaffected:

- subscribeUpstream skips the requester while its upstream claim on it is
  refused, i.e. a SUBSCRIBE for the track to it is pending;
- pickFetchUpstream skips the requester while a stitch FETCH for the track to
  it is in flight (new TrackRegistry.BeginFetch/FetchPending), so the hole is
  marked unknown (§10.13).

Deviation, recorded in STATUS.md: a genuine concurrent self-subscription is
indistinguishable from a routed-back one and is declined the same way
(DOES_NOT_EXIST / unknown range); waiting on the pending claim would deadlock
the loop case.

Tests: six in self_subscribe_test.go and the flipped late-publisher test, all
verified red at 90cba9a; the two loop tests were also verified red against
this change with only the guards disabled (11 SUBSCRIBEs / 11 stitch FETCHes
instead of 1).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…am's code (§3.3.3)

§2.4.2: a relay that detects a malformed track MUST "reset any fetch
streams with Status Code MALFORMED_TRACK", and §2.5.1: for an unknown
Mandatory Track Property in FETCH_OK, "If the relay has already
forwarded data on a fetch stream, it MUST reset the stream." The relay
reset only the data stream. The FETCH request stream was left open and
unread: the relay neither ended it nor answered a REQUEST_UPDATE on it,
and both ends held a half-open stream until the session closed.

The maintainer decided the request stream is then cancelled too, with
the data stream's code. §3.3.3: "Implementations cancel a request by
abruptly terminating any directions of the stream that are still open,
using RESET_STREAM for a direction they are sending and STOP_SENDING
for a direction they are receiving."

streamFetchRange now reports whether the track was refused, and with
which code it reset the data stream:
  - MALFORMED_TRACK for a malformed track found live (§2.4.2), in an
    upstream FETCH response (a Mandatory Track Property as an Object
    Property, §2.5.1), or in unparseable upstream FETCH_OK Track
    Properties (§12.7);
  - INTERNAL_ERROR for an upstream FETCH_OK with an unknown Mandatory
    Track Property (§2.5.1).
serveFetchObjects then cancels the request stream in both directions
with that code. A requester's own cancel still gets CANCELLED. A write
that fails because ctx's reset already hit the stream now keeps ctx's
code instead of reporting INTERNAL_ERROR.

In the tests, awaitResets waits for several stream kinds at once. The
fake upstream in refusedFetchResetsStream sends its bad FETCH_OK Track
Properties only once armed, so no refusal can happen before arming.

Verified red first: the request-stream reset was missing in
TestRelay_MalformedTrackCancelsUpstreamFetch,
TestRelay_MalformedTrackResetsCachedFetch,
TestRelay_UpstreamFetchMalformedObjectResetsStream and both subcases of
TestRelay_UpstreamFetchOKUnknownMandatoryPropertyResetsStream. The
"unknown Mandatory Track Property" subcase also fails against a variant
that cancels the request stream only on MALFORMED_TRACK.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…T (§10.2.6)

"If the RENDEZVOUS_TIMEOUT is present, the relay SHOULD hold the
subscription and wait for a publisher to appear, up to the specified
duration ... If the timeout expires without a publisher, the relay SHOULD
respond with REQUEST_ERROR with error code TIMEOUT" (§10.2.6). The relay
ignored the parameter and answered DOES_NOT_EXIST at once.

The hold (acquireUpstream):
- lasts min(requested, Config.MaxRendezvousTimeout), a new field defaulting
  to 30s ("The relay MAY use a shorter timeout"); negative holds none. 0 or
  absent still answers DOES_NOT_EXIST at once, as §10.2.6 requires;
- wakes when a publisher arrives for the track: a PUBLISH of it, or a
  covering namespace newly published here or advertised through Discovery.
  The registries signal only the waiters an arrival matches
  (TrackRegistry.AwaitUpstream, NamespaceRegistry.AwaitPublisher);
- treats DOES_NOT_EXIST, TIMEOUT and a draining publisher as "no publisher
  yet"; any other refusal outranks them in subscribeUpstream, whatever the
  candidate order, and ends the hold with its error;
- asks no publisher twice, unless its answer was cut short;
- forwards what is left on upstream SUBSCRIBEs to other relays (not to local
  publishers), and cuts such a SUBSCRIBE short when a publisher arrives here
  or the subscriber cancels (STOP_SENDING), re-asking it after a refusal;
- holds nothing on a session the relay sent GOAWAY (§10.4).

subscribeMissingPublishers now gets the publisher Seq from the last look,
not from before the hold.

Tests: rendezvous_test.go (eleven relay-level cases, four across two
relays) and registry arrivals_test.go. Each behaviour was verified red with
a mutant removing it: no hold, no budget forwarding, cancel ignored, asked
set not kept, no cap, remote not cut short, last-error-wins, GOAWAY not
awaitable, remote namespace event not waking, cut-short remote not
re-asked. TestRendezvous_NoHold/zero passes before this change by design:
it guards the existing MUST.

Found on the way, recorded in STATUS.md, not fixed: cancelling Start's ctx
leaves sessions open and unregistered, so Stop can wait without bound on an
upstream stream reader. The cross-relay tests here stop their relays before
t.Context ends for that reason.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@floatdrop
floatdrop merged commit d8adc37 into draft-20 Sep 27, 2026
11 checks passed
@floatdrop
floatdrop deleted the draft20-backlog-relay branch September 27, 2026 14:10
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