Skip to content

De-flake ClusterSingletonProxySpec, ClusterSingletonRestartSpec, and DistributedPubSubRestartSpec - #8498

Merged
Aaronontheweb merged 2 commits into
akkadotnet:devfrom
Aaronontheweb:test/wave1-deflake
Aug 27, 2026
Merged

Aaronontheweb merged 2 commits into
akkadotnet:devfrom
Aaronontheweb:test/wave1-deflake

Conversation

@Aaronontheweb

Copy link
Copy Markdown
Member

Wave-1 de-flake batch: three flaky test families, test-only changes, no product code touched. No timeout was raised, no sleep or retry backoff added anywhere — every fix removes a structural race or a blocking wait.

ClusterSingletonProxySpec (failed build 130872, Linux, ~36s in-test)

The identify test was synchronous end to end: TestProxy blocked on ExpectMsg for up to 25s per node while five in-process ActorSystems competed for the same starved thread pool, and nothing waited for the singleton to exist before sending — cluster formation, singleton startup, and identification all had to fit inside one message timeout.

  • The test is now async, with two explicit gates: all five nodes see five members Up, then each proxy publishes IdentifySingletonResult.Success. The event-stream subscription is made before the proxy actor is created, so the result cannot be missed; the wait fishes for Success because the proxy also publishes Timeout results on a zero-initial-delay timer.
  • The zero-buffering test gated on membership, then sent once into a proxy with BufferSize = 0 — which drops what it cannot forward. Membership does not imply identification, so the only message could vanish silently. It now gates on the identification result, and terminates its seed system instead of leaving it gossiping through the rest of the assembly.
  • Two waits in the timeout tests inherited single-expect-default (3s) on different TestKit instances, so their budgets did not nest and the retry loop got roughly one attempt against a queue pre-loaded with 500ms-cadence Timeout results. Both now fish with explicit 30s bounds.
  • Teardown awaits Task.WhenAll — the old Wait(30s) discarded its result, letting five cluster systems keep shutting down into the next test.

ClusterSingletonRestartSpec (failed build 130838, Windows, 15s removal window)

  • Shutdown(_sys1) waits 5s, then silently force-stops the user guardian — which does not stop /system, so remoting kept sys1's listener bound while sys3 was created on the same host:port. tcp-reuse-addr = off-for-windows makes that rebind a Windows-only failure, matching the Windows-only flake. Now await _sys1.Terminate(): the hand-over completes and the port is released.
  • Each proxy assertion created a TestProbe inside its retry loop. CreateTestProbe blocks the caller until the probe's PreStart runs and swaps the calling thread's SynchronizationContext — a blocking wait on work that itself needs a pool thread, once per ~1.1s attempt, at four sites. One probe per phase now; only the send is retried.
  • The reported 15s failure waits for sys2 to be fully Removed, which runs sys2's cluster-exiting phase — and that phase blocks on the singleton hand-over with a 10s timeout. JoinAsync only proves sys3's own view, so sys2 could still see sys3 as Joining when it left: no hand-over target, phase runs to timeout, removal misses the window. A mutual-convergence gate now requires both sys2 and sys3 to see the same two-member Up cluster before the Leave.
  • Join retries moved from 100ms to 500ms: a JoinTo arriving while the daemon is in TryingToJoin drops it back to Uninitialized and restarts the handshake, so ten retries a second was load on the very daemon the test was waiting on. The join must still be re-issued — sys3 reuses sys1's address and is refused until the old incarnation is removed.

DistributedPubSubRestartSpec (failed build 130808)

The confirmed occurrence livelocked on a reused Identify probe: association establishment flushes the buffered Identifies as a burst of ActorIdentity(null), and each retry then read one stale null forever — 80 attempts in 25s, none of them waiting.

  • The restart-kill loop now uses ActorSelection.ResolveOne inside the retry: the Identify round trip is the delivery confirmation, its temp actor is fresh per attempt (structurally immune to the stale-null livelock), and final failure raises ActorNotFoundException naming the unresolved path instead of a bare timeout on a null Subject.
  • Resolve and kill stay in one retry on purpose: splitting them lets the outbound stream drop the kill in the gap between a successful resolve and the Tell, with no resend — a local soak reproduced exactly that. The comment in the spec documents it so nobody re-splits it.
  • The 45s bound is derived, not guessed: third's restart (~15s) + one association cycle (15s associate timeout + 5s gate) + handshake and ack (~5s) ≈ 40s, nesting inside the 120s end barrier with 45s of headroom.
  • Baseline DeltaCount reads moved off TestActor (subscribed to topic1) onto probes with explicit bounds; CountAsync retries on the 500ms gossip tick with a fresh probe and a 1s bound per attempt instead of spending its 10s window on three 3s waits.

Verification

Check Result
Proxy identify, DOTNET_ThreadPool_ForceMaxWorkerThreads=3, ×10 10/10 green, tail latency gone (12.6–13.7s vs 12.8–18.6s before)
Restart spec, same 3-thread cap, ×12 12/12 green
Both singleton classes, same cap, independent re-run 5/5 green
Full Akka.Cluster.Tools.Tests 98/98 green
DistributedPubSubRestartSpec, artery, ×6 + independent re-run 7/7 green
DistributedPubSubRestartSpec, classic, ×3 3/3 green
Both projects at -warnaserror 0 warnings, 0 errors

Both specs build multi-node clusters in one process and only fail on loaded CI agents:
ClusterSingletonRestartSpec.Restarting_cluster_node_with_same_hostname_and_port_must_handover_to_next_oldest
on Windows/net10 (build 130838, at 15s) and
ClusterSingletonProxySpec.ClusterSingletonProxy_must_correctly_identify_the_singleton
on Linux/net10 (build 130872, ~36s in-test). Neither reproduces in isolation - 8/8 and 6/6 green
locally, and 30/30 and 24/24 green here under constrained thread pools and CPU pinning. So the work
is to take the load sensitivity out of the structure, not to widen anything. No timeout was raised,
no sleep and no backoff was added.

ClusterSingletonProxySpec

The identify test was synchronous throughout: TestProxy blocked on ExpectMsg for up to 25 seconds
per node, and the finally block sat in Task.WhenAll(...).Wait(30s). On a two-core agent the pool
starts at two worker threads and injects more slowly, so a blocked test thread is a large fraction
of the pool that the five ActorSystems need in order to gossip, heartbeat and deliver the very reply
being waited on. The test is now async: TestProxyAsync awaits ExpectMsgAsync and the teardown awaits
WhenAll. Awaiting the teardown also matters on its own - the discarded Wait(30s) result meant five
cluster systems could still be shutting down while the next test in the assembly ran.

The bigger correctness gap was that nothing waited for the singleton to exist. The first TestProxy
started before the cluster had formed, so cluster formation, singleton startup and proxy
identification all had to fit inside a message timeout. Two gates replace that: every node must see
all five members Up, and then each proxy must publish IdentifySingletonResult.Success. The subscription
is made in the ActorSys constructor before the proxy actor is created, so the event cannot be missed,
and the wait fishes for Success because the proxy also publishes Timeout results on a timer that
starts with no initial delay.

ClusterSingletonProxy_with_zero_buffering_should_work had the same gap with sharper teeth. It waited
for membership and then sent one message into a proxy configured with BufferSize 0, which drops
anything it cannot forward (ClusterSingletonProxy.Buffer). Membership does not imply identification,
so that send could vanish and the test would wait out its full 25 seconds for a reply that was never
coming. It now waits on the identification result. It also terminates its seed node, which was left
running for the remainder of the assembly.

Two waits in ClusterSingletonProxySingletonTimeoutTest2 read the identify result with
AwaitAssertAsync around ExpectMsgAsync, both unbounded. AwaitAssertAsync resolved
akka.test.single-expect-default (3s) and so did the inner ExpectMsgAsync on a different TestKit
instance, so the budgets did not nest and the loop got about one attempt - against a queue holding
the Timeout results the proxy emits every 500ms under that config. Both are now
FishForMessageAsync with an explicit 30s bound. The AwaitConditionAsync in
ClusterSingletonProxySingletonTimeoutTest gets an explicit bound for the same reason; with none it
inherited the 3s default for a two-node join.

ClusterSingletonRestartSpec

The spec is now async end to end. Three points of substance:

Shutdown(_sys1) waits Terminate().Wait(Dilated(5s)) and, when that expires, force-stops the user
guardian and logs a warning - verifySystemShutdown defaults to false, so nothing fails. That stops
/user but not /system, so remoting keeps sys1's listener bound; the very next statements create sys3
on that same host:port. dot-netty's tcp-reuse-addr is off-for-windows, which is the kind of
asymmetry that produces a Windows-only failure. sys1 also owns the singleton at that moment, so a
truncated shutdown cuts the hand-over to sys2 short. await _sys1.Terminate() lets CoordinatedShutdown
finish: the hand-over completes and the port is released.

Each proxy assertion created a TestProbe inside its retry loop. CreateTestProbe blocks the caller
until the probe's PreStart has run on the test-actor dispatcher (TestKitBase.cs:739) and replaces the
calling thread's SynchronizationContext (TestKitBase.cs:194) - so every attempt added a blocking wait
on work that needs a pool thread, a context swap and a leaked system actor. AwaitProxyReplyAsync
builds one probe per phase and retries only the send. Replies stay useful across attempts because the
message is a plain echo.

The 15s window that CI reported waits for sys2 to be fully Removed from sys3's view. Getting there
runs sys2's cluster-exiting CoordinatedShutdown phase, which blocks on the singleton hand-over
(ClusterSingletonManager.SetupCoordinatedShutdown) and gives up after 10s. JoinAsync only proves
sys3's own view of the cluster, so sys2 could still see sys3 as Joining when it left - no hand-over
target, phase runs to its timeout, and the removal no longer fits in 15s. A convergence gate now
requires both sys2 and sys3 to see the same two-member Up cluster before the Leave.

Smaller items: the Within wrappers are gone and each AwaitAssertAsync carries the same bound
explicitly, so no wait resolves an ambient deadline; the join assertion reads one Cluster.State
snapshot instead of two; the join retry runs at 500ms rather than 100ms, because a JoinTo arriving
while the daemon is in TryingToJoin drops it back to Uninitialized and restarts the handshake
(ClusterDaemon.cs:1273) - ten of those a second is load on the daemon the test is waiting for. The
join must still be re-issued, since sys3 reuses sys1's address and is refused until the old
incarnation is removed; a single Join would then wait out retry-unsuccessful-join-after (10s).
sys1/sys2/sys3 logs now reach the test output via InitializeLogger, as in
ClusterSingletonRestart2Spec.

Verified: both specs green in constrained-pool loops (DOTNET_ThreadPool_ForceMaxWorkerThreads=3),
full Akka.Cluster.Tools.Tests suite green, build clean with -warnaserror.
… asserting

first's restart-kill loop asserted on ActorIdentity.Subject after a 2s Identify
window. Use ActorSelection.ResolveOne instead, inside the same retry: the Identify
round trip is the delivery confirmation, its temp actor is fresh per attempt, and
a final failure raises ActorNotFoundException naming the path that never resolved
rather than a bare timeout on a null Subject.

Resolve and kill stay in one retry on purpose. Splitting them lets artery's
ordinary outbound stream drop the kill in the gap between a successful resolve
and the Tell, with no resend - a local artery soak reproduced exactly that.

Also:
- read the baseline DeltaCount on a probe with an explicit bound instead of on
  TestActor, which is subscribed to topic1
- name the same-port rebind dependency when the old system fails to stop inside
  third's WhenTerminated wait
- fresh probe and an explicit 1s bound per Count attempt, so the 10s gossip window
  retries on the 500ms tick instead of spending itself on three 3s waits
@Aaronontheweb
Aaronontheweb merged commit 9b0fc9f into akkadotnet:dev Aug 27, 2026
15 checks passed
@Aaronontheweb
Aaronontheweb deleted the test/wave1-deflake branch August 27, 2026 22:33
@Aaronontheweb Aaronontheweb added akka-cluster-tools artery Akka.Remote Artery Protocol labels Aug 28, 2026
Aaronontheweb added a commit that referenced this pull request Sep 4, 2026
… DistributedPubSubRestartSpec baseline (#8509)

* Pin the legacy allocation strategy in ClusterShardingLeavingSpec

The spec asserts that entities on surviving nodes keep the same actor
incarnation across a graceful leave. Only the legacy threshold strategy
has that property. The bounded default may move survivor shards during
the leave window, which restarts their entities under new refs - that is
the optimizer working as designed, not a regression, and the reference
implementation behaves the same way. Two CI failures carried this exact
signature: the DData variant in build 130784 and the Persistent variant
in build 130956, both an entity ref identity mismatch followed by an
after-4 barrier cascade.

Pinning rebalance-absolute-limit = 0 selects the legacy strategy, so the
identity assertion tests the algorithm that guarantees it. The spec
becomes a regression test for that strategy's leave behavior instead of
a coin flip against the bounded default.

Closes #8481.

* De-flake ReplicatorChaosSpec: rendezvous the replicators before the write burst

`ReplicatorChaosSpec` failed three times across two CI builds, on the Artery
Linux and Artery Windows lanes, with one fingerprint every time:

  [Node1:first]  Timeout (00:00:03) while expecting 15 messages. Only got 10
  [Node2:second] Timeout 00:00:00.0000044 while waiting for a message of type
                 Akka.DistributedData.GetSuccess
  [Node3..5]     barrier failed:update-during-split-verified

Root cause. Nothing makes the five replicators ready at the same time. Each node
checks `ReplicaCount(5)` against its own replicator and then walks straight into
the update phase, so `first` can start writing while another replicator has not
yet processed `first`'s MemberUp. A replicator drops a Write whose sender it
does not know:

  [.../user/replicator] Ignoring message [Write] from
    [.../user/replicator/$a] unknown node [UniqueAddress: ...]

and sends no ack. `WriteAll` cannot recover from that. Under WriteAll every node
is primary, so `SendToSecondary` re-sends to nobody, and a GCounter update
leaves `WriteAggregator._delta` null, so the delta re-send does not run either.
One dropped Write costs the whole update. Build 130934 shows `first` issuing its
first Update at 40.136 and `fifth` receiving its Welcome at 40.977; `fifth`
ignored all five Writes 12ms after that, and all five KeyC updates timed out.
Every node now proves `ReplicaCount(5)` and meets at a new `replicas-ready`
barrier before anything is written.

Second defect, and the reason the failure read as "Only got 10" instead of
naming `UpdateTimeout`. `ReceiveN(15)` sat outside any `Within` and inherited
`akka.test.single-expect-default` - 3s, the same 3s as the `WriteAll` deadline
of the updates it collects. An aggregator answers no earlier than its own
deadline, counted from when the replicator picks the Update up, which is at or
after the `ReceiveN` call. The collect can therefore never see that answer; it
is a dead heat by construction. In build 130934 the aggregators replied at
43.226, roughly 70ms too late, and the ten replies that did arrive were exactly
the WriteLocal ones, which need no aggregator. Every expect that waits on a
`WriteTo`/`WriteAll` reply now gets `WriteReplyTimeout` - the aggregator
deadline plus one single-expect-default of slack, with the arithmetic stated in
a comment on the field.

Third defect. `AssertValue` ran `Within(10s)` around an `AwaitAssert` whose
inner `ExpectMsg<GetSuccess>()` carried no bound, so each attempt could consume
the whole remaining budget and the final attempt got almost none - the logged
`00:00:00.0000044`. That reports a value which never converged as a bogus
zero-budget timeout. `AssertDeleted` and the `ReplicaCount` loop had the same
shape. All three now use a fresh probe and a 1s bound per attempt, the treatment
`ClusterShardingSpec` got in #8500.

Also swept. The `TestConductor.Blackhole`, `PassThrough` and `Exit` calls used
`.Wait(timeout)` and discarded the returned bool, so a conductor round-trip
slower than the 1s cap let the test walk on as though the partition were already
installed. They are awaited now, bounded by the conductor's own 30s
query-timeout, and a failure raises instead of vanishing. The spec is converted
to the async TestKit API throughout, which removes the remaining `.Wait()`
calls.

Verified fail-first. A throwaway build that holds one replicator's membership
events back by 4s reproduces the CI fingerprint on all five nodes exactly, with
the same five ignored Writes; the same injection passes with these changes and
produces zero ignored Writes. Bounding `ReceiveN` on its own converts the
failure into the honest assertion `Expected 'UpdateSuccess' got
'UpdateSuccess','UpdateTimeout'`, and bounding the inner expect makes all five
nodes report `Assert.Equal() Failure: Values differ` where four of five
previously reported the zero-budget timeout.

Test-only change.

* De-flake DistributedPubSubRestartSpec: read the DeltaCount baseline after gossip goes quiet

Build 130927 (artery, Windows) failed on second with "Expected value to be 3L, but
found 4L" - one extra Delta arrived between the baseline read and the post-restart
read. This is a different signature from the one PR #8498 fixed on this spec, and it
is not caused by third's restart at all.

Mechanism, confirmed by instrumenting the mediator locally and reading the counter it
actually exports:

DeltaCount counts Delta MESSAGES received, not registry changes. The mediator's Delta
handler increments _deltaCount and only then checks _nodes.Contains(bucket.Owner), so
a Delta whose payload is thrown away is still counted. While second is still waiting
for third's MemberUp, first pushes third's bucket on every 500ms gossip tick and
second discards each push - all counted. Convergence therefore arrives as a burst,
and redundant Deltas trail it, because a peer keeps re-sending until second's next
outbound gossip tells it that second caught up.

CountAsync unblocks on the first push second actually merges, so it returns while the
burst is still draining. The old code read the baseline immediately after that. Over
20 local runs the last Delta landed just 104-1010ms before that read, against a 500ms
gossip tick - one delayed tick puts a trailing Delta on the far side of the read and
the invariance assertion fails through no fault of the restart.

So the assertion was right and the baseline was wrong: "no new Deltas during third's
restart" only holds if the baseline is taken once gossip has settled. ReadStableDeltaCountAsync
now requires two samples 2s apart to agree before accepting a baseline. 2s is 4 ticks
of this spec's 500ms gossip-interval - enough for our outbound tick plus the peer's
reaction, with slack. Fresh probe per sample, matching CountAsync, so a late reply
cannot be misread as the next sample.

Verified by fault injection rather than by luck. Delivering exactly one extra Delta to
the mediator 1s after the baseline query - precisely the event the CI log implies -
reproduces the reported failure verbatim on the old code ("Expected value to be 3L,
but found 4L" on second) and passes on the new code, where the gate re-samples and
absorbs it. The post-fix margin between the last Delta and the baseline read is 3s
instead of ~150ms.

Soak: artery 8/8 under CPU load, classic 4/4. No timeout was raised, no sleep and no
backoff added; the only change is when the baseline is sampled.
Aaronontheweb added a commit that referenced this pull request Sep 12, 2026
…DistributedPubSubRestartSpec (#8498)

* De-flake ClusterSingletonProxySpec and ClusterSingletonRestartSpec

Both specs build multi-node clusters in one process and only fail on loaded CI agents:
ClusterSingletonRestartSpec.Restarting_cluster_node_with_same_hostname_and_port_must_handover_to_next_oldest
on Windows/net10 (build 130838, at 15s) and
ClusterSingletonProxySpec.ClusterSingletonProxy_must_correctly_identify_the_singleton
on Linux/net10 (build 130872, ~36s in-test). Neither reproduces in isolation - 8/8 and 6/6 green
locally, and 30/30 and 24/24 green here under constrained thread pools and CPU pinning. So the work
is to take the load sensitivity out of the structure, not to widen anything. No timeout was raised,
no sleep and no backoff was added.

ClusterSingletonProxySpec

The identify test was synchronous throughout: TestProxy blocked on ExpectMsg for up to 25 seconds
per node, and the finally block sat in Task.WhenAll(...).Wait(30s). On a two-core agent the pool
starts at two worker threads and injects more slowly, so a blocked test thread is a large fraction
of the pool that the five ActorSystems need in order to gossip, heartbeat and deliver the very reply
being waited on. The test is now async: TestProxyAsync awaits ExpectMsgAsync and the teardown awaits
WhenAll. Awaiting the teardown also matters on its own - the discarded Wait(30s) result meant five
cluster systems could still be shutting down while the next test in the assembly ran.

The bigger correctness gap was that nothing waited for the singleton to exist. The first TestProxy
started before the cluster had formed, so cluster formation, singleton startup and proxy
identification all had to fit inside a message timeout. Two gates replace that: every node must see
all five members Up, and then each proxy must publish IdentifySingletonResult.Success. The subscription
is made in the ActorSys constructor before the proxy actor is created, so the event cannot be missed,
and the wait fishes for Success because the proxy also publishes Timeout results on a timer that
starts with no initial delay.

ClusterSingletonProxy_with_zero_buffering_should_work had the same gap with sharper teeth. It waited
for membership and then sent one message into a proxy configured with BufferSize 0, which drops
anything it cannot forward (ClusterSingletonProxy.Buffer). Membership does not imply identification,
so that send could vanish and the test would wait out its full 25 seconds for a reply that was never
coming. It now waits on the identification result. It also terminates its seed node, which was left
running for the remainder of the assembly.

Two waits in ClusterSingletonProxySingletonTimeoutTest2 read the identify result with
AwaitAssertAsync around ExpectMsgAsync, both unbounded. AwaitAssertAsync resolved
akka.test.single-expect-default (3s) and so did the inner ExpectMsgAsync on a different TestKit
instance, so the budgets did not nest and the loop got about one attempt - against a queue holding
the Timeout results the proxy emits every 500ms under that config. Both are now
FishForMessageAsync with an explicit 30s bound. The AwaitConditionAsync in
ClusterSingletonProxySingletonTimeoutTest gets an explicit bound for the same reason; with none it
inherited the 3s default for a two-node join.

ClusterSingletonRestartSpec

The spec is now async end to end. Three points of substance:

Shutdown(_sys1) waits Terminate().Wait(Dilated(5s)) and, when that expires, force-stops the user
guardian and logs a warning - verifySystemShutdown defaults to false, so nothing fails. That stops
/user but not /system, so remoting keeps sys1's listener bound; the very next statements create sys3
on that same host:port. dot-netty's tcp-reuse-addr is off-for-windows, which is the kind of
asymmetry that produces a Windows-only failure. sys1 also owns the singleton at that moment, so a
truncated shutdown cuts the hand-over to sys2 short. await _sys1.Terminate() lets CoordinatedShutdown
finish: the hand-over completes and the port is released.

Each proxy assertion created a TestProbe inside its retry loop. CreateTestProbe blocks the caller
until the probe's PreStart has run on the test-actor dispatcher (TestKitBase.cs:739) and replaces the
calling thread's SynchronizationContext (TestKitBase.cs:194) - so every attempt added a blocking wait
on work that needs a pool thread, a context swap and a leaked system actor. AwaitProxyReplyAsync
builds one probe per phase and retries only the send. Replies stay useful across attempts because the
message is a plain echo.

The 15s window that CI reported waits for sys2 to be fully Removed from sys3's view. Getting there
runs sys2's cluster-exiting CoordinatedShutdown phase, which blocks on the singleton hand-over
(ClusterSingletonManager.SetupCoordinatedShutdown) and gives up after 10s. JoinAsync only proves
sys3's own view of the cluster, so sys2 could still see sys3 as Joining when it left - no hand-over
target, phase runs to its timeout, and the removal no longer fits in 15s. A convergence gate now
requires both sys2 and sys3 to see the same two-member Up cluster before the Leave.

Smaller items: the Within wrappers are gone and each AwaitAssertAsync carries the same bound
explicitly, so no wait resolves an ambient deadline; the join assertion reads one Cluster.State
snapshot instead of two; the join retry runs at 500ms rather than 100ms, because a JoinTo arriving
while the daemon is in TryingToJoin drops it back to Uninitialized and restarts the handshake
(ClusterDaemon.cs:1273) - ten of those a second is load on the daemon the test is waiting for. The
join must still be re-issued, since sys3 reuses sys1's address and is refused until the old
incarnation is removed; a single Join would then wait out retry-unsuccessful-join-after (10s).
sys1/sys2/sys3 logs now reach the test output via InitializeLogger, as in
ClusterSingletonRestart2Spec.

Verified: both specs green in constrained-pool loops (DOTNET_ThreadPool_ForceMaxWorkerThreads=3),
full Akka.Cluster.Tools.Tests suite green, build clean with -warnaserror.

* De-flake DistributedPubSubRestartSpec: resolve the association before asserting

first's restart-kill loop asserted on ActorIdentity.Subject after a 2s Identify
window. Use ActorSelection.ResolveOne instead, inside the same retry: the Identify
round trip is the delivery confirmation, its temp actor is fresh per attempt, and
a final failure raises ActorNotFoundException naming the path that never resolved
rather than a bare timeout on a null Subject.

Resolve and kill stay in one retry on purpose. Splitting them lets artery's
ordinary outbound stream drop the kill in the gap between a successful resolve
and the Tell, with no resend - a local artery soak reproduced exactly that.

Also:
- read the baseline DeltaCount on a probe with an explicit bound instead of on
  TestActor, which is subscribed to topic1
- name the same-port rebind dependency when the old system fails to stop inside
  third's WhenTerminated wait
- fresh probe and an explicit 1s bound per Count attempt, so the 10s gossip window
  retries on the 500ms tick instead of spending itself on three 3s waits

(cherry picked from commit 9b0fc9f)
Aaronontheweb added a commit that referenced this pull request Sep 12, 2026
… DistributedPubSubRestartSpec baseline (#8509)

* Pin the legacy allocation strategy in ClusterShardingLeavingSpec

The spec asserts that entities on surviving nodes keep the same actor
incarnation across a graceful leave. Only the legacy threshold strategy
has that property. The bounded default may move survivor shards during
the leave window, which restarts their entities under new refs - that is
the optimizer working as designed, not a regression, and the reference
implementation behaves the same way. Two CI failures carried this exact
signature: the DData variant in build 130784 and the Persistent variant
in build 130956, both an entity ref identity mismatch followed by an
after-4 barrier cascade.

Pinning rebalance-absolute-limit = 0 selects the legacy strategy, so the
identity assertion tests the algorithm that guarantees it. The spec
becomes a regression test for that strategy's leave behavior instead of
a coin flip against the bounded default.

Closes #8481.

* De-flake ReplicatorChaosSpec: rendezvous the replicators before the write burst

`ReplicatorChaosSpec` failed three times across two CI builds, on the Artery
Linux and Artery Windows lanes, with one fingerprint every time:

  [Node1:first]  Timeout (00:00:03) while expecting 15 messages. Only got 10
  [Node2:second] Timeout 00:00:00.0000044 while waiting for a message of type
                 Akka.DistributedData.GetSuccess
  [Node3..5]     barrier failed:update-during-split-verified

Root cause. Nothing makes the five replicators ready at the same time. Each node
checks `ReplicaCount(5)` against its own replicator and then walks straight into
the update phase, so `first` can start writing while another replicator has not
yet processed `first`'s MemberUp. A replicator drops a Write whose sender it
does not know:

  [.../user/replicator] Ignoring message [Write] from
    [.../user/replicator/$a] unknown node [UniqueAddress: ...]

and sends no ack. `WriteAll` cannot recover from that. Under WriteAll every node
is primary, so `SendToSecondary` re-sends to nobody, and a GCounter update
leaves `WriteAggregator._delta` null, so the delta re-send does not run either.
One dropped Write costs the whole update. Build 130934 shows `first` issuing its
first Update at 40.136 and `fifth` receiving its Welcome at 40.977; `fifth`
ignored all five Writes 12ms after that, and all five KeyC updates timed out.
Every node now proves `ReplicaCount(5)` and meets at a new `replicas-ready`
barrier before anything is written.

Second defect, and the reason the failure read as "Only got 10" instead of
naming `UpdateTimeout`. `ReceiveN(15)` sat outside any `Within` and inherited
`akka.test.single-expect-default` - 3s, the same 3s as the `WriteAll` deadline
of the updates it collects. An aggregator answers no earlier than its own
deadline, counted from when the replicator picks the Update up, which is at or
after the `ReceiveN` call. The collect can therefore never see that answer; it
is a dead heat by construction. In build 130934 the aggregators replied at
43.226, roughly 70ms too late, and the ten replies that did arrive were exactly
the WriteLocal ones, which need no aggregator. Every expect that waits on a
`WriteTo`/`WriteAll` reply now gets `WriteReplyTimeout` - the aggregator
deadline plus one single-expect-default of slack, with the arithmetic stated in
a comment on the field.

Third defect. `AssertValue` ran `Within(10s)` around an `AwaitAssert` whose
inner `ExpectMsg<GetSuccess>()` carried no bound, so each attempt could consume
the whole remaining budget and the final attempt got almost none - the logged
`00:00:00.0000044`. That reports a value which never converged as a bogus
zero-budget timeout. `AssertDeleted` and the `ReplicaCount` loop had the same
shape. All three now use a fresh probe and a 1s bound per attempt, the treatment
`ClusterShardingSpec` got in #8500.

Also swept. The `TestConductor.Blackhole`, `PassThrough` and `Exit` calls used
`.Wait(timeout)` and discarded the returned bool, so a conductor round-trip
slower than the 1s cap let the test walk on as though the partition were already
installed. They are awaited now, bounded by the conductor's own 30s
query-timeout, and a failure raises instead of vanishing. The spec is converted
to the async TestKit API throughout, which removes the remaining `.Wait()`
calls.

Verified fail-first. A throwaway build that holds one replicator's membership
events back by 4s reproduces the CI fingerprint on all five nodes exactly, with
the same five ignored Writes; the same injection passes with these changes and
produces zero ignored Writes. Bounding `ReceiveN` on its own converts the
failure into the honest assertion `Expected 'UpdateSuccess' got
'UpdateSuccess','UpdateTimeout'`, and bounding the inner expect makes all five
nodes report `Assert.Equal() Failure: Values differ` where four of five
previously reported the zero-budget timeout.

Test-only change.

* De-flake DistributedPubSubRestartSpec: read the DeltaCount baseline after gossip goes quiet

Build 130927 (artery, Windows) failed on second with "Expected value to be 3L, but
found 4L" - one extra Delta arrived between the baseline read and the post-restart
read. This is a different signature from the one PR #8498 fixed on this spec, and it
is not caused by third's restart at all.

Mechanism, confirmed by instrumenting the mediator locally and reading the counter it
actually exports:

DeltaCount counts Delta MESSAGES received, not registry changes. The mediator's Delta
handler increments _deltaCount and only then checks _nodes.Contains(bucket.Owner), so
a Delta whose payload is thrown away is still counted. While second is still waiting
for third's MemberUp, first pushes third's bucket on every 500ms gossip tick and
second discards each push - all counted. Convergence therefore arrives as a burst,
and redundant Deltas trail it, because a peer keeps re-sending until second's next
outbound gossip tells it that second caught up.

CountAsync unblocks on the first push second actually merges, so it returns while the
burst is still draining. The old code read the baseline immediately after that. Over
20 local runs the last Delta landed just 104-1010ms before that read, against a 500ms
gossip tick - one delayed tick puts a trailing Delta on the far side of the read and
the invariance assertion fails through no fault of the restart.

So the assertion was right and the baseline was wrong: "no new Deltas during third's
restart" only holds if the baseline is taken once gossip has settled. ReadStableDeltaCountAsync
now requires two samples 2s apart to agree before accepting a baseline. 2s is 4 ticks
of this spec's 500ms gossip-interval - enough for our outbound tick plus the peer's
reaction, with slack. Fresh probe per sample, matching CountAsync, so a late reply
cannot be misread as the next sample.

Verified by fault injection rather than by luck. Delivering exactly one extra Delta to
the mediator 1s after the baseline query - precisely the event the CI log implies -
reproduces the reported failure verbatim on the old code ("Expected value to be 3L,
but found 4L" on second) and passes on the new code, where the gate re-samples and
absorbs it. The post-fix margin between the last Delta and the baseline read is 3s
instead of ~150ms.

Soak: artery 8/8 under CPU load, classic 4/4. No timeout was raised, no sleep and no
backoff added; the only change is when the baseline is sampled.

(cherry picked from commit e340146)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant