Draft-20 compliance: relay backlog - #114
Merged
Merged
Conversation
… §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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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
0bc25d5INCLUDE_PROPERTIES=0 survives a REQUEST_UPDATE (§10.9, §10.2.21).82bb1b9A SUBSCRIBE_TRACKS holder isn't forwarded a track its own SUBSCRIBE is establishing (§6.1, §10.20).5e853adA PUBLISH_SKIPPED holds across a prefix update that moves away and back (§6.1).023d79eTwo forwards racing for one track send one PUBLISH.58e40d4A client can SUBSCRIBE and FETCH a track it publishes (§5.1), including through late-publisher and missing-publisher SUBSCRIBEs (§9.5).a97626eRENDEZVOUS_TIMEOUT holds a SUBSCRIBE for a publisher (§10.2.6).FETCH and fill
b1128f1FETCH_OK End Of Track is set at the END_OF_TRACK Object (§10.14).ada7018A cancelled FETCH resets its request and data streams (§5.2).4de6761Objects stitched from an upstream FETCH are bounded by its MAX_CACHE_DURATION (§12.3).2b72d55A malformed track cancels its FETCHes and resets their fetch streams (§2.4.2).8e8085bA fill ends at the Largest Object the subscriber was told (§5.1.3).2288a1fA refused FETCH's request stream is reset with its data stream's code (§3.3.3).Data plane
6e99b0bReplay streams carry the first Object's delivery-timeout override (§8, §12.1, §12.2).9518daaA merged Subgroup that a clean survivor does not cover is reset (§11.4.3).23c67e9A merged Subgroup's FIN is judged from the ledger's lowest Object (§11.4.3).Errors, tokens, GOAWAY, shutdown
5350f87Malformed Track Properties are refused with INTERNAL_ERROR, not MALFORMED_TRACK, outside FETCH (§10.6.2).18b4708A REQUEST_UPDATE's tokens are verified on SUBSCRIBE, FETCH and PUBLISH_NAMESPACE (§10.2.2, §10.9.1).6480ab3TRACK_STATUS is answered for any track SUBSCRIBE would accept (§10.15).e3a70d4No request is initiated on a session with a GOAWAY in either direction (§10.4), and GOING_AWAY is the refusal code.9ac663fShutdown closes with GOAWAY_TIMEOUT only after a GOAWAY ran out, NO_ERROR otherwise (§3.5).Housekeeping
9d31287The 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.90cba9aanda5b043dare test-only: the close-code test's flake fix and a line wrap.Decisions and interpretations
Self-subscription
RENDEZVOUS_TIMEOUT
Config.MaxRendezvousTimeout, default 30s. A negative value holds nothing.Error codes
Other defaults
API changes
relay.Config.MaxRendezvousTimeout(new).session.Session.GoawaySent()(new).Known gaps (in STATUS.md)
Start's ctx ends each session's handler without closing the session, and drops it from the setStopcloses.Stopcan then wait without bound on an upstream SUBSCRIBE's stream reader.Start's doc says the cancel "terminates live sessions".t.Contextends to avoid it.drainStragglerignoringStop's ctx, and a draining pool relay yielding DOES_NOT_EXIST are noted in the backlog.Verification
go test ./...andgolangci-lint runpass.go test -race ./pkg/relay/...passes; the final rendezvous change was run five times.make bench-quickshows allocs/op matching the refreshed baseline.🤖 Generated with Claude Code