Skip to content

De-flake ClusterShardingQueriesSpec: wait for all regions to register before triggering allocation - #8524

Merged
Aaronontheweb merged 2 commits into
akkadotnet:devfrom
Aaronontheweb:fix/deflake-cluster-sharding-queries-spec
Sep 9, 2026
Merged

Aaronontheweb merged 2 commits into
akkadotnet:devfrom
Aaronontheweb:fix/deflake-cluster-sharding-queries-spec

Conversation

@Aaronontheweb

Copy link
Copy Markdown
Member

Summary

De-flakes ClusterShardingQueriesSpec, which failed on all four nodes in PR validation build 131139 (Windows, Artery) on a PR that only moved source-generator code. Test-only change; product code is untouched.

Root cause

The spec asserts a two-two-two shard layout across the three regions: it computes timeouts as shards divided by region count, expects four stats and two failures, and later expects two shards and two failures per region. Nothing guaranteed that layout. The sharding started barrier only proves each node called StartSharding; it says nothing about which regions the coordinator has registered.

From the failing build's log, with third hosting the coordinator:

  • second and busy sent their first Register while the coordinator was still reading its initial state from DData. WaitingForInitialState drops a Register in that state rather than stashing it, and logs "ShardRegion tried to register but ShardCoordinator not initialized yet".
  • State load completed 15 ms later. A region only re-sends on its RegisterRetry timer, 250 ms and doubling.
  • third's local region won the retry race. The controller's pings arrived and shards 0, 5, 1, and 3 were all allocated to third, the only region the allocation strategy could see. second and busy registered afterward and received shards 4 and 2.

Result: a four-one-one layout. Five shards answer the stats query, one fails on busy's zero-millisecond timeout, and the assertion reads "expected 4, but found 5". Because rebalance-interval is 120 s, that layout is permanent, so the existing AwaitAssertAsync loops re-read a state that never changes. The other nodes then time out waiting on the coordinator's fan-out, and the controller fails the barrier. The race is between two remote round trips, the coordinator's majority read and the regions' singleton identification; Artery's shorter first-contact path shifts the odds toward the registers arriving first.

Fix

Before the controller sends the pings that trigger allocation, it polls GetCurrentRegions through its proxy until the coordinator reports all three regions, with a fresh probe per attempt and a dilated 30 s budget. This is the same gate ClusterShardingRolePartitioningSpec already uses. Once all three regions are registered, allocation is two-two-two by construction, because each ShardHomeAllocated update stashes the next GetShardHome, so allocation proceeds sequentially against fresh state. The comment at the fix site records the race, the log evidence, and the Artery angle.

Verification, all under AKKA_MNTR_TRANSPORT=artery

  • Before the fix, stock build: 10 of 10 passed locally. The window does not open on this machine by chance.
  • Throwaway fault injection, not committed: delay the coordinator's initial state by 500 ms so the first registers drop as in CI, and give busy and second a 3 s registration retry. Unfixed spec: 5 of 5 failed with the CI signature, all six shards on third, plus the barrier cascade. Fixed spec on the same injected build: 5 of 5 passed.
  • After the fix, clean build: 20 of 20 passed.
  • Full Akka.Cluster.Sharding.Tests.MultiNode assembly: 137 of 137 passed.

Related

ClusterDeathWatchSpec, DistributedPubSubRestartSpec, and ClusterShardingRegistrationCoordinatedShutdownSpec each failed once on Artery PR builds the same week and are not addressed here.

…e allocating shards

`ClusterShardingQueriesSpec` failed on all four nodes in the Artery Windows MNTR
lane of build 131139 (PR akkadotnet#8519, which only moved source-generator code). `third`
reported the primary failure after exhausting a full 30s converge-then-assert
budget:

  AwaitAssert failed, timeout [00:00:30] is over after [30] attempts
  Expected regions.Values.Select(i => i.Stats.Count).Sum() to be 4, but found 5.

`second` and `busy` then timed out 10s waiting for `ClusterShardingStats` (the
coordinator's fan-out waits on the region `third` had just torn down), and
`controller` failed the `received failed stats from timed out shards vs empty`
barrier.

Mechanism. The spec asserts a 2/2/2 shard layout (`timeouts = NumberOfShards /
regions.Count`, Stats == 4, Failed == 2, and later `Shards.HaveCount(2)` /
`Failed.HaveCount(2)` per region) but nothing guaranteed it. The `sharding
started` barrier only proves that each node called `StartSharding`; it says
nothing about which regions the coordinator has on its books. In the failing
run the coordinator on `third` was still reading its initial state from DData
when `second` and `busy` sent their first `Register` (20.318 and 20.339; the
state load completed at 20.354). `DDataShardCoordinator.WaitingForInitialState`
drops a `Register` rather than stashing it:

  DatatypeA: ShardRegion tried to register but ShardCoordinator not
    initialized yet: [[akka://...@localhost:54620/system/sharding/DatatypeA]]

and a region only re-sends on its `RegisterRetry` timer (250ms, doubling towards
`retry-interval`). `third`'s own region won the retry race (registered at
20.557), the controller's 20 pings reached the coordinator at 20.643, and
`LeastShardAllocationStrategy` could only pick from the regions it knew about:

  20.669  Shard [0] allocated at [.../DatatypeA#1012486927]        (third)
  20.739  Shard [5] allocated at [.../DatatypeA#1012486927]        (third)
  20.797  Shard [1] allocated at [.../DatatypeA#1012486927]        (third)
  20.833  Shard [3] allocated at [.../DatatypeA#1012486927]        (third)
  20.833  ShardRegion registered: [...@localhost:54620/...]        (second)
  20.840  ShardRegion registered: [...@localhost:54619/...]        (busy)
  20.957  Shard [4] allocated at [...@localhost:54620/...]         (second)
  20.991  Shard [2] allocated at [...@localhost:54619/...]         (busy)

4/1/1: five shards answer the stats query and only `busy`'s single shard fails
on its 0ms `shard-region-query-timeout`, hence "found 5". The spec disables
rebalancing (`rebalance-interval = 120s`), so that layout is permanent and the
`AwaitAssertAsync` loops from akkadotnet#8310 cannot converge - they re-read a state that
will never change.

The race is between two remote round trips - the coordinator's majority read of
its state and the regions' identification of the singleton - so the transport
only shifts the odds. Artery's first-contact path is shorter than DotNetty's
association handshake, which is why the remote Registers land that much earlier
relative to the state load in the Artery lane.

Fix. Before the controller sends the pings that trigger allocation, it polls
`GetCurrentRegions` through its proxy until the coordinator reports all three
regions - fresh probe and 1s bound per attempt inside a 30s dilated budget, the
same gate `ClusterShardingRolePartitioningSpec` already uses. Once the
coordinator knows all three regions the layout is 2/2/2 by construction: each
`ShardHomeAllocated` update stashes the next `GetShardHome`, so allocations are
sequential against fresh state and the least-shards ordering fills the regions
round-robin. The product is doing what it is designed to do (allocate among the
regions that have registered); the spec asserted a layout it had not waited
for.

Verified fail-first. The race does not open on its own on a Linux workstation
(10/10 green before the change, ~6s per run). A throwaway build that holds the
coordinator's initial-state result back by 500ms - so every first Register is
dropped, as in CI - and gives `busy` and `second` a 3s registration retry
reproduces the signature on the unfixed spec 5/5 times ("Expected ...
Stats.Count).Sum() to be 4, but found 6": all six shards on `third`), and
passes 5/5 with this gate on the same injected build. With the
injection removed: 20/20 green under Artery, and the full
Akka.Cluster.Sharding.Tests.MultiNode assembly passes under Artery
(137/137, 9m11s).

Test-only change.

@Aaronontheweb Aaronontheweb left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@Aaronontheweb
Aaronontheweb merged commit 78f5cdc into akkadotnet:dev Sep 9, 2026
15 checks passed
@Aaronontheweb
Aaronontheweb deleted the fix/deflake-cluster-sharding-queries-spec branch September 9, 2026 13:30
Aaronontheweb added a commit that referenced this pull request Sep 10, 2026
… before triggering allocation (#8524) (#8539)

Backport of #8524 to v1.5.
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