Repository navigation
Artery: honour the reconnect backoff, return the held control element on stage stop, publish the inbound context before bind, host the materializer under /system - #8554
Conversation
49050d0 to
dbe3dbf
Compare
16c53fb to
193c13c
Compare
193c13c to
c9b211e
Compare
a8bf46e to
0713b2e
Compare
0713b2e to
55164b0
Compare
3bcf5c4 to
827e65d
Compare
Aaronontheweb
left a comment
There was a problem hiding this comment.
Largely LGTM - with some nitpicks
| { | ||
| // Idempotent: DefaultMaterializer (ActorMaterializer.cs) does the same injection for the | ||
| // /user-hosted materializer, and InjectTopLevelFallback is safe to call more than once. | ||
| System.Settings.InjectTopLevelFallback(ActorMaterializer.DefaultConfig()); |
There was a problem hiding this comment.
Technically not necessary.
| System.Settings.InjectTopLevelFallback(ActorMaterializer.DefaultConfig()); | ||
| var settings = ActorMaterializerSettings.Create(System); | ||
| var haveShutDown = new AtomicBoolean(); | ||
| var supervisor = System.SystemActorOf( |
There was a problem hiding this comment.
It's preferable that we have a static, named actor here serving as the materializer for all artery streams for the reasons elicited in this PR, but I wonder a wee bit about starting a StreamSupervisor directly and making that the materializer context, since I thought those actors were typically created downstream of the materialization root?
I suppose it doesn't necessarily matter, but I f find that a little strange. The idea in principle is still right though. We need an actor to act as the owner of all the Akka.Remote streams used for Artery and it has to be a top-level actor under the /system hierarchy so we don't accidentally terminate in-flight streams before the /user hierarchy has had a chance to complete terminating.
| if (envelope.Message is SystemMessageEnvelope) | ||
| return; | ||
|
|
||
| var requeued = streamId switch |
| // and fresh associations formed DURING a graceful cluster leave, which runs for tens of | ||
| // seconds with CoordinatedShutdown.ShutdownReason already set and this transport fully | ||
| // alive. | ||
| if (IsTransportTerminating()) |
| .WatchTermination(Keep.Both) | ||
| .Via(Flow.FromGraph(new LaneWriteBatchStage(LaneWriteBatchMaxBytes))) | ||
| .Via(Flow.FromGraph(new LaneWriteBatchStage(LaneWriteBatchMaxBytes, | ||
| onDropped: bytes => |
There was a problem hiding this comment.
Logs un-delivered outbound messages when the transport is terminated
| /// </summary> | ||
| private bool IsTransportTerminating() => | ||
| _isShutdown | ||
| || _materializer is null || _materializer.IsShutdown |
There was a problem hiding this comment.
Can't the _materializer be null during startup too?
Overall I think this needs a bit of simplification - also, doesn't it overlap quite a bit with private bool IsActorSystemTerminating() ?
e84da73 to
2e2fa76
Compare
2e2fa76 to
d7c9ee8
Compare
… on stage stop, publish the inbound context before bind Fixes three independent Artery transport defects surfaced by the DistributedPubSubRestartSpec concurrency analysis (control stream reconnecting at 5.15/s against a configured 1/s in build 131174, attempts 10-230 in 42.7s): - ScheduleOutboundRestart opened each stream's materialize-once gate (control/ordinary/large) immediately, a full backoff before its own ScheduleOnce callback ran, so an on-demand enqueue landing during the backoff window re-materialized the stream early and bypassed outbound-restart-backoff entirely. The gate is now reset only inside the post-backoff callback. - OutboundHandshakeStage's PostStop only unsubscribed the handshake listener, silently discarding the one element (HandshakeReq, ArteryHeartbeat, ArteryQuarantined, ClearSystemMessageDelivery, or an unwrapped DaemonMsgCreate) it may have been holding when a materialization stopped before completing the handshake. PostStop now returns it to the association's own channel via a new IOutboundContext.ReturnUndelivered seam, publishing Dropped instead of discarding it silently when there is no room. - ArteryRemoting.Start() published _localUniqueAddress/_inboundContext several statements after the bound port became known, widening the window in which an already-accepting listener could dispatch a connection through HandleIncomingConnection before the inbound context existed. These two fields are now published first, ahead of every other post-bind statement. Each fix ships with a focused unit test (P1: ArteryOutboundRestartBackoffSpec, reflection-driven gate-timing checks; P2: a new OutboundHandshakeStage case in ArteryHandshakeSpec; P5: ArteryInboundContextPublishSpec, using a new ArteryTransportSetup.OnBoundPortKnown test hook) verified to fail against the pre-fix code and pass against the fix. Does not implement the ordinary-lane handshake-via-control-queue coupling or ActorSystemTerminating (separate issues).
…rvive until the shutdown flush; report LaneWriteBatchStage's held bytes on stop ArteryRemoting.Start materialized every stream on ActorMaterializer.Create(System), whose StreamSupervisor is a /user actor. The system stops /user before the remoting terminator runs, so every Artery stream aborted with AbruptStageTerminationException before the shutdown flush, and flush-wait-on-shutdown never applied. The materializer's supervisor is now a /system actor, and the outbound materialize guards also check for a terminating system. LaneWriteBatchStage reports the bytes of a retained batch it drops on stop through an onDropped callback that publishes Dropped.
…make the inbound-context and backoff tests discriminate; do not return sequenced system messages - ScheduleOutboundRestart's ORDINARY and LARGE post-backoff callbacks now release the materialize-once gate (mirroring the existing pre-schedule release) when the restart is refused because the peer was quarantined during the backoff window. Quarantine is not permanent, and returning with the gate still latched wedged the stream forever, blocking even the ActorSelectionMessage and new-incarnation sends that must pierce a quarantine. Added ArteryOutboundRestartQuarantineGateSpec, which fails on the old code and passes on the fix, for both the ordinary and large streams. - Start() now fires the P5 ArteryTransportSetup.OnBoundPortKnown hook after _defaultAddress/_addresses are published instead of one statement after _inboundContext is assigned, so the hook's argument is no longer trivially true; updated ArteryInboundContextPublishSpec's remarks to describe the ordering it now actually guards. - ReturnUndeliveredOutboundElement no longer re-offers a SystemMessageEnvelope -- that traffic already has its own resend buffer, and re-offering it would hand a duplicate to a fresh SystemMessageDeliveryStage instance. Documented on IOutboundContext.ReturnUndelivered and the method itself that the re-offer lands at the channel's tail (no head-insert available on Channel<T>), and what that means for an unwrapped DaemonMsgCreate that precedes a Watch. - ArteryOutboundRestartBackoffSpec's fourth fact now counts materialize-callback invocations instead of asserting on IsControlOutboundMaterialized alone, since that flag was true on both the buggy and fixed code for opposite reasons and could not tell them apart; renamed its DisplayName to match. - Added ArteryReturnUndeliveredOutboundElementSpec, exercising the production ReturnUndeliveredOutboundElement wiring (rather than a hand-wired delegate) for the ordinary lane's routing and the full-channel Dropped path.
…act asserts on the restart bookkeeping, which the winning caller's callback could not show
…dinatedShutdown; refused materialization no longer latches the gate; the ack-race spec reaches its race again The up-front guards in MaterializeOutboundStream/MaterializeOrdinaryOutboundWithLanes used IsActorSystemTerminating(), which also trips on CoordinatedShutdown.ShutdownReason -- set as the first statement of Run(), before phase one. That refused every new outbound stream (fresh associations, reconnects after backoff, handshake replies) for the whole of a graceful cluster leave, while /user was still alive and the transport was not shutting down. Split the predicate: a new IsTransportTerminating() (the transport's own state, deliberately without ShutdownReason) gates the two materialize guards; IsActorSystemTerminating() stays as-is for the exception filters, where the breadth is harmless. The refused path also used to return quietly from inside MaterializeOnceGate.EnsureStarted, which flips the gate before invoking the callback and only resets on throw -- leaving a permanent wedge (gate latched, no stream, no restart scheduled). Both guards now reset the gate on refusal, matching the pattern the surrounding catch blocks already use. Added a regression test proving a system inside CoordinatedShutdown but before the transport shuts down can still materialize a fresh outbound stream (fails on the old broad guard, passes with the narrow one -- checked by hand both ways). Restructured ArteryShutdownSystemMessageAckRaceSpec, which the broad guard had made vacuous: it now stops the StreamSupervisor directly instead of racing a full ActorSystem.Terminate() (which, now that the supervisor lives under /system, no longer reaches the race at all), and reliably reaches and discriminates the InvalidOperationException catch again (verified by temporarily removing the catch and confirming the spec fails, then restoring it). Also: corrected two stale "/user"-hosted-supervisor comments; gave LaneWriteBatchStage's OnPush append-failure path an onDropped callback instead of a silent discard, and noted it (like the whole lane-batching stage) is unreachable at the outbound-lanes=1 default; pinned ArteryGracefulTerminateFlushSpec's zero-count EventFilter to an explicit 3s window and noted the attempt==1 cadence gate; amended BREAKING_CHANGES_V1.6.md row 47 to say the flush behavior it promised only actually takes effect starting with this change, plus two consequences of the /system move it didn't previously cover.
…ods; drop the redundant fallback injection and the null clause CreateSystemMaterializer now names its StreamSupervisor "artery-stream-supervisor" instead of one of StreamSupervisor.NextName()'s generated names, and drops the InjectTopLevelFallback call that ActorMaterializerSettings.Create already performs internally. IsTransportTerminating, IsActorSystemTerminating, and IsStreamSupervisorTerminating are rewritten as plain methods with early returns instead of chained || expressions (dropping the now-provably-dead materializer-null clause), the two inline copies of the old catch-filter condition now call IsTransportTerminating(), and comments that only restated the removed clauses are trimmed to match.
d7c9ee8 to
bf1b8a2
Compare
Stacked on #8570.
What changes
Three small defects in Artery's outbound path, each with a test that fails on the old code and passes on the new.
ScheduleOutboundRestartreset the stream's materialize gate before scheduling the backoff, so any control enqueue during the wait re-materialized the stream at once. The reset now happens inside the post-backoff callback, for the control, ordinary, and large streams. Measured on build 131174 before the fix: five control reconnects per second against a configured one, attempts 10 to 230 in 42.7 s.OutboundHandshakeStagetakes one element out of its channel while it gates on the handshake. When the stage stopped, that element was lost. It could be aHandshakeReq, anArteryHeartbeat, or a remote-deploy message.PostStopnow hands it back through a new internalReturnUndeliveredon the outbound context, which re-enqueues it into the right channel or publishesDroppedif the channel is full. The ordinary-lane path now gives each lane its own context so the element goes back to its own lane.No public API changes: the new members are on internal types, and the approval files are unchanged. These are bug fixes that restore intended behavior, so they are not recorded in the 1.6 breaking-changes ledger.
Why
How it was checked
Akka.Remotebuilds with warnings as errors.Akka.Remote.Tests, Artery filter: 278 passed.Akka.API.Tests: 18 passed, no approval diff. Each new test was run with its fix reverted by hand and failed, then with the fix restored and passed:ArteryOutboundRestartBackoffSpec(three of four facts fail without the fix), the new case inArteryHandshakeSpec, andArteryInboundContextPublishSpec.Second commit: the materializer moves under /system, and LaneWriteBatchStage reports held bytes on stop
ArteryRemoting.Startmaterialized every stream on a materializer whose stream supervisor is a user-guardian actor. The system stops the user guardian before the remoting terminator runs, so on every graceful terminate every Artery stream aborted withAbruptStageTerminationExceptionbefore the shutdown flush, andflush-wait-on-shutdownnever applied. The supervisor is now a system actor, built from public Akka.Streams API.LaneWriteBatchStagereports the bytes of a retained batch it drops on stop through a callback that publishesDropped. UnlikeOutboundHandshakeStageit cannot hand the data back, since it sits after encode and the lane merge with no single envelope left, so it makes the loss visible rather than recoverable. That stage is only materialized atoutbound-lanes > 1, so this has no effect at the shipping default of 1.This is the product half of the
RemoteNodeRestartDeathWatchSpecanalysis; the test half is #8557. Pekko hosts the transport's materializer under its system materializer for the same reason, with two differences that remain: Pekko builds a separate control materializer, and it reads the materializer settings frompekko.remote.artery.advanced.materializer, a key Akka.NET does not have, so Artery's streams here are wired to the globalakka.stream.materializersettings.What a user sees on shutdown after this change: a graceful terminate with a peer that cannot be reached now waits up to
flush-wait-on-shutdown, 2 s by default, where before the streams had already aborted and the wait cost nothing.Abort()pays that flush as well. Inbound streams also survive the user-guardian teardown, so a terminating node drains and closes its accepted connections and dead-letters late arrivals instead of resetting them. The unbind wait inShutdown()is bounded by the TCP stage's own timer; #8570 removes the idle case, and a connection accepted but not yet materialized still holds it for up tosubscription-timeout.New spec
ArteryGracefulTerminateFlushSpec: with the old placement, terminating a system with an open outbound stream logs the abrupt-termination warning; with the new one it logs none, because the shutdown's own graceful completion runs first. Verified failing with the materializer change reverted and passing with it restored. TheLaneWriteBatchStagetest was verified the same way. Artery filter: 280 passed.Akka.API.Tests: 18 passed, no approval diff; both changed classes are internal. Ledger: the existingfix/artery-shutdown-flushrow already promised this flush and has not been true on dev until now; the fifth commit amends that row rather than adding one.One pre-existing quirk noted, not changed: an inbound Artery connection that never decoded its stream-type preamble can leave a graph interpreter that only clears on the forced-shutdown timeout, independent of materializer placement.
Third commit: what the adversarial review found
ArteryOutboundRestartQuarantineGateSpecquarantines the peer during the backoff window and asserts the gate reopens and a later send materializes, for the ordinary and large streams. Both facts fail on the second commit and pass now.ArteryInboundContextPublishSpecfail; restoring them makes it pass.SystemMessageEnvelope, which the delivery stage's resend buffer retransmits until acked. Returning it only handed a duplicate to the next materialization. The doc comment said as much and the code now agrees.DaemonMsgCreatethat can reopen the create-before-watch ordering for that one element; the inbound side tolerates the race, and the comment onReturnUndeliveredsays so.ArteryReturnUndeliveredOutboundElementSpecgoes throughReturnUndeliveredOutboundElementitself, for lane routing and for theDroppedevent when the channel is full, instead of a hand-wired delegate.Artery filter: 284 passed, run twice, before and after the by-hand reversions.
Akka.API.Tests: 18 passed, no approval diff. No ledger entry.What this PR's own CI run showed, and the restack on #8570
Build 131333: the Windows unit-test job was cancelled at the 20-minute per-project cap inside
Akka.Remote.Tests, and the Linux run of the same assembly took 19 min 46 s against 8 min on other PRs. Per-test timings showed the new tests adding about a minute and existing Artery tests slowing by 12 to 54 s each, all in the peer-restart and graceful-shutdown family. A bisect across the three commits put all of it on the second one, a flat 5 s per ArteryActorSystemtermination.The cause is not in this PR. Hosting the materializer under
/systemmeans Artery's listener now performs a real graceful unbind at shutdown instead of dying with/userfirst, and Akka.Streams' TCP unbind has a latent bug: its counter of connections awaiting initialization is seeded at -1, so the idle fast path never fires and everyUnbind()waits the 5 s subscription timeout. That is fixed in #8570, one token plus a regression test, and this PR is now stacked on it. With the fix, the Artery test namespace runs in 2 min 10 s instead of 13 min 25 s, all 284 passing, and the three-cycle create-and-terminate test takes 395 ms instead of 15 s. The three commits here are unchanged in content, verified by patch-id after the rebase.Fourth commit
The fourth backoff fact still passed on the buggy code: the materialize-once gate runs the winning caller's callback, and after the buggy synchronous reset the winner is the enqueue path's own lambda, so a callback counter cannot tell the two worlds apart. The fact now asserts on the association's restart bookkeeping, which latches the instant the gate is reset. Verified by moving the reset back in front of the backoff by hand: the fact fails, and passes again with the product code restored.
Fifth commit: what the review of the second commit found
The second commit had also widened the outbound materialize guards to refuse whenever the system was "terminating", and that predicate included the coordinated-shutdown reason, which is set before phase one runs. From the moment a graceful shutdown or a cluster leave began, an Artery node would not materialize a new outbound stream, the control stream and reconnects after backoff included, for the whole of the sharding hand-off, leave, and exiting phases, and the refused stream's gate stayed latched with no restart armed. Pekko checks only the transport's own shutdown flag.
ArteryOutboundMaterializeDuringCoordinatedShutdownSpecholds coordinated shutdown at its first phase and asserts a fresh outbound stream still materializes. It fails on the fourth commit and passes on this one.ArteryShutdownSystemMessageAckRaceSpechad become vacuous, because the broad guard returned before the race it exists to cover. It now stops the system-hosted supervisor directly and reaches the race again; with its guarding catch neutered by hand it fails five of five runs, restored it passes./user; corrected. The graceful-terminate spec's zero-count filter takes an explicit 3 s window.LaneWriteBatchStage's class doc says it is unreachable at the default lane count, and its append-failure path now folds into the singleDroppedreport on stop instead of discarding silently.Artery filter after this commit: 285 passed in 2 min 7 s.
Akka.API.Tests: 18 passed, no approval diff.Sixth commit: the maintainer's review nits
The stream supervisor is now a named top-level system actor,
/system/artery-stream-supervisor, instead of a randomly named one. It stays aStreamSupervisor: the materializer attaches interpreter actors straight to its cell, but it still needs that actor'sMaterializeprotocol for a supervisor that has not started yet, and itsPostStopis what sets the flag behindIsShutdown. The explicit top-level config injection is gone, sinceActorMaterializerSettings.Createalready performs it. The three terminating predicates are plain methods with early returns; the null-materializer clause is gone, because the materializer is assigned as the first statement ofStart()and nothing can materialize before that, so a null materializer means not started rather than terminating; and the two exception filters that spelled the clauses out inline call the predicate instead. The batch-stage comment says what the callback does in two sentences. Artery filter: 285 passed in 2 min 8 s. Approvals unchanged.