Repository navigation
Akka.Streams: fix a possible lost wakeup in MergeHub (missing memory fence) - #8665
Merged
Aaronontheweb merged 1 commit intoSep 29, 2026
Merged
Conversation
…fence) MergeHub's consumer stores _needWakeup = true and then re-polls the ConcurrentQueue. The queue's empty path only does volatile reads, so the store could be reordered after the re-poll (store-load reordering, legal on x86 and ARM). A producer could then enqueue, read _needWakeup == false and skip the wakeup while the consumer saw an empty queue. The hub stalled until the next enqueue, which never comes when perProducerBufferSize is 1. JVM Akka marks the flag @volatile, which is sequentially consistent. C# volatile is not, so add a full fence after the store, and the same after _shuttingDown = true in PostStop, which has the same shape. Mark both fields volatile to match upstream. Un-skip MergeHub_must_work_with_long_streams_when_buffer_size_is_1.
This was referenced Oct 1, 2026
Aaronontheweb
added a commit
to Aaronontheweb/akka.net
that referenced
this pull request
Oct 2, 2026
…fence) (akkadotnet#8665) MergeHub's consumer stored needWakeup = true and then re-polled an empty ConcurrentQueue with no full fence between them, so the store could be reordered after the poll and a producer's wakeup lost (reproduced as stream hangs with buffer size 1). Adds the fence JVM @volatile provided, the same for shuttingDown in PostStop, and un-skips the buffer-size-1 long-stream spec. (cherry picked from commit 3d89d29)
Aaronontheweb
added a commit
to Aaronontheweb/akka.net
that referenced
this pull request
Oct 3, 2026
…fence) (akkadotnet#8665) MergeHub's consumer stored needWakeup = true and then re-polled an empty ConcurrentQueue with no full fence between them, so the store could be reordered after the poll and a producer's wakeup lost (reproduced as stream hangs with buffer size 1). Adds the fence JVM @volatile provided, the same for shuttingDown in PostStop, and un-skips the buffer-size-1 long-stream spec. (cherry picked from commit 3d89d29)
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 race
MergeHub.HubLogiccoordinates one consumer (the hub's stage) with many producers through aConcurrentQueueand a_needWakeupflag:TryProcessNext): on an empty queue it stores_needWakeup = true, re-polls the queue once, then stops.Enqueue):_queue.Enqueue(e), thenif (_needWakeup) { _needWakeup = false; _wakeupCallback(); }.This is a Dekker-style store-then-load handshake on both sides, so it needs sequential consistency. JVM Akka gets that from
@volatile. In .NET,_needWakeupwas a plainbool, and nothing between the store and the re-poll acts as a full fence. I decompiledConcurrentQueueSegment<T>.TryDequeuefrom the .NET 10 CoreLib: the empty path doesVolatile.Read(Head),Volatile.Read(slot.SequenceNumber),Volatile.Read(Tail)and returnsfalse, with noInterlockedoperation. Acquire reads don't stop an earlier store from moving past them, and x86-TSO allows exactly this store-load reordering (ARM64 allows more). That makes this interleaving legal:_needWakeup = true(still in its store buffer); re-poll readsTail == Head, so the queue looks empty and the consumer stops.TailCAS publishes the element (full fence); it reads_needWakeup == falseand sends no wakeup.trueand an element is queued, but nobody is scheduled to read it.The hub stalls until another producer enqueues. With
perProducerBufferSize = 1, every producer is waiting on demand from the hub, so the stall is permanent.PostStophas the same shape:_shuttingDown = true, then drain the queue. A producer inSinkLogic.PreStartenqueuesRegisterand re-checksIsShuttingDown. If the store and the drain reorder, the producer can miss both the drain and the flag and wait forever.C#
volatilealone would not fix this. A volatile store followed by a volatile load can still reorder, so we need a full fence.Evidence
A local-only stress harness (not committed) ran repeated
MergeHub.Source<int>(1)streams for 45 s per case, with a 5 s hang timeout per iteration, on 8-core x86-64 with .NET 10:dev(3 runs)That is 9 hangs in about 26,900 iterations before, and 0 in about 41,000 after. At the old rate we would expect about 14 hangs in the post-fix runs.
Fix (Hub.cs only)
Interlocked.MemoryBarrier()right after_needWakeup = true, before the re-poll.Interlocked.MemoryBarrier()right after_shuttingDown = trueinPostStop._needWakeupand_shuttingDownare nowvolatile, which matches upstream and also safely publishes_wakeupCallback(set inPreStart) to producers that see_needWakeup == true.ConcurrentQueue.Enqueuealready acts as a fence.The fence only runs when the queue is empty, not on the per-element path.
Tests
MergeHub_must_work_with_long_streams_when_buffer_size_is_1(it was skipped as "Very racy") and gave it an explicitWaitAsynctimeout. 30/30 passes locally with the fix. It also passed 30/30 ondev, so it guards against gross regressions; it is not a reliable repro of this race. The window is too small for a single run.HubSpecpassed 5/5 runs.