Repository navigation
Artery: stop dropping the first ordinary message when the peer dialed first - #8504
Aaronontheweb merged 7 commits into
Conversation
|
Triage for the Linux Artery MNTR failure (build 130927): ambient agent-load flakes; the fix's behavior in these very logs is clean. Three specs failed, none of them
Why this is not the fix regressing, from these logs themselves:
Will retrigger for a second artery sample once this run settles. |
|
Windows Artery addendum for build 130927: same verdict as the Linux lane — ambient flakes, fix behavior clean.
A slimmed revision of this PR is coming (review feedback: reduce the diff to the essential core), which will trigger a fresh full matrix — that supersedes the retrigger promised earlier. |
… first Fixes akkadotnet#8496. `OutboundHandshakeStage.PreStart` treated "this association already has a `UniqueRemoteAddress`" as "our handshake is done" and skipped straight to `Completed`. That field is also set by the INBOUND direction - when WE handle the peer's `HandshakeReq` - so it proves we know the peer's uid, not that the peer knows ours. When a peer dialed us first, our first outbound ordinary stream therefore sent no `HandshakeReq` of its own, and our first user message raced our `HandshakeRsp` (a different TCP connection, no ordering relationship) into the peer's unknown-origin gate: Dropping inbound message [Akka.Actor.ActorSelectionMessage] from unknown origin uid [817823530] (no completed handshake for this uid yet). Ordinary messages have no resend path, so the message was gone. The receiving side documented an invariant ("the sender's gate cannot complete before we have processed its Req") that the `PreStart` shortcut quietly broke. `AssociationState.OutboundHandshakeCompleted` now records, per INCARNATION, whether the peer has answered a `HandshakeReq` of ours. Only the new `CompleteOutboundHandshake` sets it, and only `InboundHandshakeStage.HandleRsp` calls that - a Rsp is the sole event proving the peer registered our uid. The flag rides the same immutable snapshot and CAS swap as the rest of the association state: a uid change resets it, `Quarantine` clears it, so it never outlives the incarnation it describes. Two gates consult it, both for ordinary/large/lane streams only: * `PreStart`'s fast path (`CanSkipOwnHandshake`) - so a stream materialized against an inbound-only association runs the normal, well-tested handshake path instead of shortcutting. * `RefreshCompletionFromContext` - load-bearing too, since the generation counter it already checked is advanced by the peer's own inbound Req: without it, a stream that did send its Req would complete on the peer's Req *retry* landing while we wait for our Rsp, which is the same bug through the back door. The CONTROL stream deliberately keeps the old, weaker rule. Its traffic is never subject to the receiver's unknown-origin drop - every control envelope (handshake, heartbeat, quarantine notice, system messages and their Ack/Nack) is dispatched regardless of whether the origin uid is registered - and it must be able to complete on an inbound handshake: the `HandshakeRsp` we owe a peer is enqueued on that very control queue, so a strict rule there would deadlock two systems that dial each other at the same instant, each holding the Rsp the other is waiting for. `ForceReqOnStart` is unaffected: it still forces the full handshake path, and the two conditions compose - the fast path is skipped when either says so. Also, at both inbound gates (`InboundHandshakeStage` and the lane-routed copy in `ArteryInboundProcessingStage`), the unknown-origin drop is raised from DEBUG to a rate-limited WARNING: the first drop warns at once, later ones at most once per 10s carrying the count suppressed in between. This class of loss should be diagnosable in the field, not invisible. Tests. `ArteryPeerDialedFirstSpec` covers the issue's stage-level repro (an association created by a real inbound `HandshakeReq`; the ordinary stage must send its own Req and hold traffic until a real `HandshakeRsp` arrives), the peer's-Req-retry variant, fast-path preservation, the new warning, and an end-to-end test over two real Artery systems where B dials A first and A's `HandshakeRsp` is dropped - the deterministic form of the race. Four of the five fail on unfixed code; the fast-path test passes before and after, as it must. `AssociationStateSpec` gains the transition, incarnation-reset and quarantine coverage for the new flag. The invariant comment on `InboundHandshakeStage` now describes the corrected construction.
ce3f291 to
c486c79
Compare
…lling for it `OutboundHandshakeStage` learned that a handshake had completed only by re-reading `AssociationState` at events it already processed: `OnPull`, `OnPush`, and the retry tick. While the stage is gating traffic, none of the first two can happen - it holds the element, so it stops pulling, and it has pushed nothing, so downstream demand is still outstanding. The retry timer was therefore the only wake-up, and the first message on a gated association waited up to a full `handshake-retry-interval` (1s by default) after the peer had already answered. The previous commit's gate routes a new population of streams through that state: on a cluster join every peer dials the seed, so the seed's first ordinary message to each of them - the cluster `Welcome` - is now held. Measured join->welcome latency, ReplicatorChaosSpec on Artery, 3 runs x 4 peers per revision on one machine: dev base 0.062 - 0.085 s mean 0.071 s gate commit only 1.052 - 1.106 s mean 1.073 s The consequence is not academic. In the CI logs for this PR the seed reached its own `ReplicaCount(5)` and issued `Update(WriteAll)` while a peer, whose `Welcome` was still in flight, had not yet processed membership for the seed - so that peer's Replicator discarded all five writes at the application layer: [INFO] Cluster Node [33251] - Welcome from [...:36423] [INFO] .../user/replicator Ignoring message [Write] from [[...:36423] (x5) `WriteAll` has no re-send path for a `GCounter` delta, so each aggregator timed out one ack short and the spec's `ReceiveN(15)` saw 10. Nothing was lost by the transport; the writes simply arrived before the receiver was ready, and the extra second of `Welcome` latency is what put them there. Detection is now push: * `Association` carries a keyed listener registry (`SubscribeHandshakeStateChanged` / `UnsubscribeHandshakeStateChanged`, `ConcurrentDictionary`), fired from the CAS transition after the state swap and the generation increment, so a woken subscriber always reads the snapshot the notification belongs to. * It fires on EVERY completion, either direction - not only on `CompleteOutboundHandshake`. This is a "re-check" nudge, and every subscriber re-verifies the snapshot against its own rule before acting, so a wake that does not satisfy that rule (an inbound Req while an ordinary stream waits for its Rsp, or a completion belonging to a newer incarnation than the one the stage registered against) leaves the stage exactly as it was. One no-op callback buys the same latency fix for the control stream and for reconnects. * The stage registers a `GetAsyncCallback` when it enters `ReqInProgress`, then re-checks - register-THEN-check, so a completion landing in between is either seen by the check or delivered by the callback, never neither. On wake it re-verifies through the existing refresh logic and releases `_pendingMessage` through the same push/pull machinery the retry-tick release used. * The registration is dropped on completion and in `PostStop`; outbound streams restart against an association that outlives them, so a leaked entry would accumulate. Invoking a stopped stage's callback is independently harmless - the interpreter drops async input for a completed logic - so this is hygiene, not a correctness race. * Both timers keep their purpose: the retry timer retransmits the `HandshakeReq`, the timeout timer fails the stage. Neither is the detection path any more. Tests: a release-latency test pins a 30s retry interval, holds an element, delivers the peer's Rsp and requires the element within 3s - it fails on the previous commit with "no element signaled during 00:00:03" and passes here in milliseconds. A lifecycle test asserts the listener is registered while gated, gone after the stream stops, and that a completion recorded afterwards neither throws nor is lost. Verification: join->welcome back to 0.063-0.084 s (5 runs, mean 0.065-0.079) with the gate still firing 8-10 times per run; ReplicatorChaosSpec 5/5 on Artery; AttemptSysMsgRedelivery 4/4 Artery and 2/2 classic; artery unit suite 271/271; full Akka.Remote.Tests 670 passed / 5 skipped; Akka.API.Tests 18/18 with no approval delta.
Aaronontheweb
left a comment
There was a problem hiding this comment.
Found some issues that need fixing
| return Complete(association, peer, association.CompleteOutboundHandshake(peer)); | ||
| } | ||
|
|
||
| private AssociationState Complete( |
| internal sealed class AssociationState | ||
| { | ||
| private AssociationState(int incarnation, UniqueAddress? uniqueRemoteAddress, ImmutableHashSet<long> quarantinedUids) | ||
| private AssociationState( |
There was a problem hiding this comment.
why not make this a record?
There was a problem hiding this comment.
Done in 363de0e — internal sealed record, transitions rewritten with with-expressions, no-op paths still return this. One hazard the conversion surfaced and fixed: the CAS success checks compared with ==, which flips from reference to value equality on a record — a lost CAS whose winner was value-equal could report as won. Both are now explicit ReferenceEquals. Type docs note snapshots are reference-compared only (generated equality would be shallow over the ImmutableHashSet).
Per the style guide's "default to sealed classes and records". Behaviour is unchanged; this is a shape change plus one hardening fix that the conversion makes necessary. The four properties become `init` accessors and the three transitions are expressed as `with` copies, so each one now names exactly the fields it changes and the carry-over of everything else (notably `QuarantinedUids` across an incarnation bump) is structural rather than a positional argument that has to be got right. Two guardrails, since the CAS machinery leans on identity rather than equality: * Every no-op transition still returns `this`, never a `with` copy - `with` always allocates, and `Association`'s CAS loops plus the registry's reverse-index maintenance detect "nothing to publish" by reference. * The two `Interlocked.CompareExchange(...) == current` comparisons are now explicit `ReferenceEquals` calls. As a class that `==` was reference equality; as a record it would have become compiler-generated value equality, under which a lost CAS could masquerade as a won one. Nothing compares snapshots with `==`/`Equals` anywhere - the specs assert identity with `BeSameAs`/`NotBeSameAs` - and the generated value equality would be shallow over `ImmutableHashSet` in any case, so the type documents that snapshots are to be compared by reference only.
Adopts the injection pattern from akkadotnet#7314 (which obsoleted DateTimeOffsetNowTimeProvider in favour of passing an ITimeProvider) for the two inbound gates' unknown-origin drop-warning rate limit, so the 10s window can be driven by a virtual clock in tests instead of by waiting. Both stages take an optional `ITimeProvider? timeProvider = null`, typed as the BASE interface rather than IScheduler or IDateTimeOffsetNowTimeProvider so a TestScheduler satisfies it. ArteryRemoting passes the actor system's scheduler at all three materialization sites; a stage constructed without one resolves the materializing system's scheduler at PreStart, so the 13 existing stage constructions in the specs needed no changes and still get a real clock. The throttle arithmetic moves from `DateTime.UtcNow` to `timeProvider.Now`, with the "never warned" sentinel becoming a nullable DateTimeOffset rather than MinValue. `Now` specifically, not MonotonicClock: TestScheduler virtualizes only Now, and virtualization is the whole point of the change. ArteryUnknownOriginDropWarningSpec proves it end to end on TestConfigs. TestSchedulerConfig: a first drop warns with "[0] further drop(s) suppressed", a second inside the window is silent, `((TestScheduler)Sys.Scheduler).Advance(11s)` crosses the window with no real waiting, and the third drop warns again reporting "[1] further drop(s) suppressed". Reverting the stage to a wall clock fails that test on the third assertion, so it genuinely pins the injected clock. A second test hands the stage a frozen provider and advances the system scheduler five virtual minutes to confirm the explicit argument wins over the materializer-resolved default. The spec is its own class because virtualizing the scheduler also freezes the timers the sibling handshake specs depend on.
|
Triage for the classic Linux MNTR failure (build 130956): the #8481 family, second occurrence — unrelated to this PR.
Meanwhile the Artery Linux lane passed for the second consecutive run — the |
The dropped test verified that an explicitly passed ITimeProvider wins over the materializer-resolved default - which is testing a null-coalescing operator. The remaining test in the spec carries the real proof: the suppression window crossed with TestScheduler.Advance in virtual time, ablation-checked against a wall-clock implementation.
Aaronontheweb
left a comment
There was a problem hiding this comment.
One DateTime.UtcNow call that needs to be cleaned up but otherwise LGTM
| Pull(_stage.In); | ||
| } | ||
|
|
||
| /// <summary> |
There was a problem hiding this comment.
So this only happens because Artery can open multiple parallel connections when multiple outbound lanes are supported. If another outbound lane has already handled the handshake, we can and must skip it.
There was a problem hiding this comment.
Right, with one addition: the fast path fires even at outbound-lanes = 1, because control, ordinary, and large are separate streams sharing one association — and re-materializations of any of them hit it too. Multi-lane is one producer of the situation, not the only one. The comment now leads with your sentence and covers the single-lane case in one more.
| private bool CanSkipOwnHandshake(AssociationState state) => | ||
| state.UniqueRemoteAddress is { } already && | ||
| Equals(already.Address, _stage.Context.RemoteAddress) && | ||
| (_stage.IsControlStream || state.OutboundHandshakeCompleted); |
There was a problem hiding this comment.
If the state determines that we've already completed this outbound handshake OR we're the control stream.
If we are the control stream, we are sending and handling the handshake. Therefore, it would deadlock us if we were to wait on on it. Thus, we skip it in this instance.
There was a problem hiding this comment.
Exactly — and your wording is now the comment (with one clause of mechanism kept: the Rsp we owe the peer sits in the same control queue the gate would close, so two simultaneous dialers would each hold the other's answer).
| // the generation counter below is advanced by the peer's own inbound Req as well, | ||
| // so without this an ordinary stream would complete on the peer's Req retry that | ||
| // happens to land while we wait for our Rsp -- issue #8496 through the back door. | ||
| if (!_stage.IsControlStream && !state.OutboundHandshakeCompleted) |
There was a problem hiding this comment.
If we are not the control stream, we cannot continue until the handshake is done.
| if (_state != State.Completed) | ||
| return; | ||
|
|
||
| if (_pendingMessage is { } held && IsAvailable(_stage.Out)) |
There was a problem hiding this comment.
If we have a pending user message that hasn't been sent yet, now we can flush it now that the handshake's done. This gets invoked as part of the async callback. This is what prevents the message loss
There was a problem hiding this comment.
Correct — your three sentences are the comment now, verbatim.
| /// </para> | ||
| /// </summary> | ||
| internal sealed class AssociationState | ||
| internal sealed record AssociationState |
There was a problem hiding this comment.
LGTM - prefer record as the standard going forward
| private AssociationState( | ||
| int incarnation, | ||
| UniqueAddress? uniqueRemoteAddress, | ||
| bool outboundHandshakeCompleted, |
There was a problem hiding this comment.
new property being tracked
The outbound stage still timed its HandshakeReq injections with DateTime.UtcNow. It now takes the same optional ITimeProvider the inbound stages take, and the whole _lastInject mechanism moves to that clock in one piece. _lastInject changes from DateTime to a nullable DateTimeOffset, so the "never injected" sentinel is null instead of MinValue. All five write sites and both read sites go through the provider: the PreStart fast-path stamp, the completion stamp in RefreshCompletionFromContext, the two re-injection stamps in OnPush, TryInjectReq, and ShouldReinjectForLiveness. The field is either all provider time or all wall clock. A mixed field would compare a UtcNow stamp against a provider reading and misjudge both intervals as soon as anyone virtualizes the clock. The provider resolves at the top of PreStart, before any code stamps _lastInject: the explicit constructor argument first, then the materializing system's scheduler. ArteryRemoting passes System.Scheduler at both outbound materialization sites, which matches what the inbound sites already do. Comments in this file are rewritten in plainer English. Same facts, shorter sentences, no storytelling. The fast-path docs now lead with the reason a stream can skip the handshake - another outbound stream on the association already handled it - and then cover the single-lane case where that other stream is the control, ordinary or large stream of the same association. The control-stream paragraph states the deadlock directly: we send and handle the handshake, waiting on it would deadlock us, so we skip it, and skipping is safe because the receiver only drops ordinary envelopes as unknown-origin.
|
Final-run triage (build 131079): the one non-blocking failure is |
Fixes #8496.
Two commits: the correctness gate (revised per review feedback to its essential core), and a push-based completion wake-up that replaces the stage's timer-poll detection — added after the gate exposed a pre-existing ~1s detection latency (see below).
Problem
Artery's receiver drops any ordinary envelope from an unregistered origin uid — safe only if every sender holds ordinary traffic behind its own completed handshake.
OutboundHandshakeStage.PreStartbroke that invariant by treating "this association has aUniqueRemoteAddress" as "our handshake is done." That field is also set by the inbound direction when we process the peer'sHandshakeReq, so when the peer dialed us first, our first outbound ordinary stream skipped its ownHandshakeReq. Our first user message then raced ourHandshakeRsp— a different TCP connection, no ordering between them — into the peer's unknown-origin gate. Lost races were silent, permanent message loss (ordinary messages have no resend). Three artery MNTR lane failures carried the exact signature; classic is immune becauseEndpointWriterbuffers until association.Fix
One new fact is tracked:
AssociationState.OutboundHandshakeCompleted— per-incarnation, in the same immutable snapshot asUniqueRemoteAddress, set only via the newIInboundContext.CompleteOutboundHandshake, which onlyHandleRspcalls (aHandshakeRspis only ever sent after the responder registers the requester's uid, so it is proof the peer knows us). Reset on incarnation change, cleared by quarantine.Consulted at the two places "handshake done" is decided for ordinary/large/lane streams:
PreStart's fast path now requires it. An association born from the peer's inbound handshake falls through to the normal handshake flow — send our Req, hold traffic, existing retry/timeout machinery. Cost: one round trip before the first ordinary message on a peer-opened association.RefreshCompletionFromContext(the in-progress completion check) gets the same 2-line gate. This closes a third door: while our stream holds traffic awaiting its Rsp, the peer retrying its ownHandshakeReqadvances the handshake generation and would falsely complete us. An ablation run proves this gate is load-bearing — with it commented out, exactly one test fails (the message releases after 110ms with no Rsp).HandshakeRspthe other is waiting for. Rationale documented at the decision site.Also: the unknown-origin drop is raised from silent DEBUG to a rate-limited WARNING at both gate sites (inline fields, no new type — first drop immediate, then at most once per 10s with a suppressed count), and the
InboundHandshakeStageinvariant comment now states the real construction. A compact DEBUG line names the #8496 ordering when the corrected path fires, which is what makes the fix observable in soak logs.No wire-format change, no config change, no public API delta.
ForceReqOnStartcomposes unchanged.Fail-first
Four of the five tests fail on unfixed dev:
HandshakeReqs sentVerification (slim version)
-warnaserror: Akka.Remote + both test projectsArteryPeerDialedFirstSpecAkka.Remote.TestsAkka.API.TestsAttemptSysMsgRedeliverySpecartery/classicCaveat as before: fast loopback lets pre-fix code pass the MNTR soak (the Rsp wins the race locally), so the soak is no-regression evidence plus proof the corrected path fires; the deterministic proof is the unit/E2E suite. Across the earlier full-matrix run of the fatter revision, ~2 artery MNTR suites produced zero unknown-origin drops.
The second commit: push-based completion detection
The stage detected handshake completion only by polling on the retry timer (default 1s). The
HandshakeRsptypically lands in ~1ms, but a held first message waited for the next tick — a pre-existing latency the gate newly exposed on peer-opened associations, where the first held message is often the clusterWelcomeon a star join. That ~1.1sWelcomedelay is what madeReplicatorChaosSpecfail on this PR's earlier CI runs: peers' replicators silently ignoreWriteAllwrites that arrive before membership lands, andWriteAllnever re-sends.Detection is now push, matching the reference stage design: the association keeps a keyed listener registry fired from the CAS transition helper after every state swap; the stage registers a
GetAsyncCallbackon enteringReqInProgress(register-then-check, so nothing lands unseen), re-verifies through the existing refresh logic on wake, and releases the held message through the same path a retry-tick release used. Listeners are deregistered on completion and inPostStop(hygiene — post-stop invocation is independently safe via the interpreter's stage-completed guard). The retry timer now only retransmits Reqs; the timeout timer still bounds the whole handshake.Reviewer note on one breadth choice: the notification fires on every completion (both directions), not just outbound — every subscriber re-verifies its own condition on wake, so a non-qualifying nudge is a no-op, and the broad trigger also removes the same 1s latency for control-stream and reconnect waits. Narrowing it to outbound-only is a two-line change if preferred.
Fail-first: a release-latency test pins a 30s retry interval so polling cannot pass — it fails on the gate-only commit (
no element signaled during 00:00:03) and passes with the wake-up. Welcome-latency A/B (Join→Welcome, ReplicatorChaosSpec formation, artery):ReplicatorChaosSpecartery: 5/5 locally on the final branch. Full battery re-run green (ArteryPeerDialedFirstSpec7/7, artery unit suite 271/271,Akka.Remote.Tests670/670, API tests no delta, MNTR AttemptSysMsgRedelivery 4/4 artery + 2/2 classic,-warnaserrorclean). Branch total vs dev: 349 product lines, 540 test lines.Split out
The restart-path variant of the same conflation (
ForceReqOnStartre-materializations can complete off a generation bump caused by the peer's inbound Req) is real but not harmful today — the flag being true guarantees the peer has our uid, so nothing is dropped. Filed separately with a repro sketch rather than carried here.