request: fix deadlock between Pool.Close and BatchStore.Fetch (#843) - #1211
Conversation
849dddf to
c91739f
Compare
87c3851 to
a55f277
Compare
tock-ibm
left a comment
There was a problem hiding this comment.
Review: fix deadlock between Pool.Close and BatchStore.Fetch
The core fix looks correct. I traced the interleavings:
- Pre-canceled ctx: the new loop condition
ctx.Err() == nilis false, soWait()is skipped andFetchreturns promptly — the deadlock is gone. - ctx cancels in the window between the condition check and parking in
Wait(): no lost signal, becauseFetchholds the write lock continuously across that window and the ctx-watch goroutine cannotLock()toSignal()untilWait()has parked (and released the lock). - Insert signal / spurious wakeup: guarded by the same lock and the condition-recheck loop.
- Inner goroutine lifetime: always terminates via
<-finished(closed by defer) or<-ctx.Done(); no leak.
The inline comments below are latent/quality issues, not blockers.
One minor, pre-existing nit (it sits outside this diff, so noting it here rather than inline): the defer func() { close(finished) }() in Fetch could be the simpler defer close(finished) — equivalent, and closer to the repo's "favor explicit over clever" guidance. Optional, and not introduced by this PR.
🤖 Reviewed with Claude Code
| // lost, leaving Fetch blocked forever while holding bs.lock and thus | ||
| // deadlocking Pool.Close. | ||
| bs.lock.Lock() | ||
| bs.signal.Signal() |
There was a problem hiding this comment.
The ctx-cancel wakeup uses sync.Cond.Signal(), which wakes only one arbitrary waiter. This is safe today only because there is one NextRequests/Fetch caller per BatchStore (the batcher primary loop, sequential). But NextRequests takes only rp.lock.RLock, so nothing structurally prevents two concurrent Fetch calls on one store.
If that ever happens: both park in signal.Wait(); the ctx of Fetch-A is canceled; A's goroutine Signal()s, but the single wakeup can land on Fetch-B, which re-checks its own (non-canceled) ctx and re-parks — leaving Fetch-A blocked forever and re-introducing the very Pool.Close deadlock this PR fixes. Broadcast() would be robust to that; at minimum it's worth a comment documenting the one-caller-per-store invariant this relies on.
| @@ -157,21 +157,27 @@ func (bs *BatchStore) Fetch(ctx context.Context) ([]interface{}, []interface{}) | |||
| go func() { | |||
There was a problem hiding this comment.
The ctx-watch goroutine is spawned on every Fetch, before checking readyBatches. On the hot path — sustained load, where readyBatches is usually non-empty on entry — Fetch could return immediately, but instead pays a goroutine create + teardown (via close(finished)) once per batch.
Consider hoisting an early return for len(bs.readyBatches) > 0 || ctx.Err() != nil above the go func(), so the goroutine is spawned only when Fetch is actually about to Wait().
| func TestFetchWithCanceledContext(t *testing.T) { | ||
| sugaredLogger := testutil.CreateLogger(t, 0) | ||
|
|
||
| for i := 0; i < 2000; i++ { |
There was a problem hiding this comment.
Two small test-quality points:
- The
2000iteration count is a magic number with no comment tying it to the race window it exercises, which makes the stress bound arbitrary to tune later. A brief comment (or a named const) would help. - On the failure path,
t.Fatalffires from the main test goroutine while the blockedFetchgoroutine (and its inner goroutine) are left running for the rest of the process; and a genuine deadlock makes the suite hang up to the full 10s per iteration before reporting. Neither is fatal for a regression test, but both are worth noting.
9d4ebbd to
9f7cc5a
Compare
2de67a2 to
e64991a
Compare
tock-ibm
left a comment
There was a problem hiding this comment.
Thanks for tracking this down. The fix itself looks correct: we found no deadlock, race or lost wakeup in the new Fetch. The inline comments cover the regression test (it never exercises the fixed path), a behavior change after Prune, and a few inaccurate comments.
Not in the diff, but related:
NextRequestsholdsrp.lock.RLock()across a blockingFetch(pool.go:216), soPool.Closeonly works because callers cancel the ctx first (cancelBatch()beforeMemPool.Close()in batcher_role.go). That ordering isn't documented or enforced anywhere.- The same lost wakeup as #843 exists in the assembler.
PartitionPrefetchIndex.Stop(node/assembler/partition_prefetch_index.go:416) callscancelContextFunc(); stateCond.Broadcast()without holdingstateCond.L, whilePopOrWait/Putcheck the ctx under L and thenWait(). We reproduced it in 3/3 runs:Collator.Stophangs, soAssembler.Stop/SoftStophangs while holdinga.lock. The fix is the same as this PR's (broadcast under L). Probably worth its own issue. - Pre-existing:
batch.Prune(request/batch.go:48) only deletes from the batch's map. It never callsonDelete(which releases the semaphore permit and decrements the size) and never clearskeys2Batches. After a reconfig prune in batching mode, pruned IDs are rejected as duplicates, and capacity leaks untilSubmittimes out. Reproduced in a scratch test.
🤖 Reviewed with Claude Code
| bs := NewBatchStore(100, 100*8, func(string) {}, sugaredLogger) | ||
|
|
||
| ctx, cancel := context.WithCancel(context.Background()) | ||
| cancel() // context is already canceled before Fetch is called |
There was a problem hiding this comment.
The test doesn't reach the code the fix changed. ctx is cancelled before Fetch is called, so the new len(bs.readyBatches) == 0 && ctx.Err() == nil guard (line 155) is false. The watcher goroutine, the Lock-before-Broadcast and the wait loop never run.
We tried two mutations under -race:
- Restoring the old unlocked
bs.signal.Signal()in the watcher: this test and the whole./requestpackage still pass. - Deleting the watcher entirely: this test still passes. Only
TestBatchStorecatches it, by hanging until thego testtimeout.
Suggestion: add a case that starts Fetch with a live ctx and cancels it concurrently, e.g. go Fetch(ctx) then cancel() after a small random spin. In our runs that caught the unlocked-Signal mutation within a few thousand iterations, and a cancel-while-parked variant caught the deleted watcher on the first iteration. Also, the return value is discarded, so "canceled ctx => empty batch" isn't asserted.
There was a problem hiding this comment.
Reworked as TestFetchCancelWhileWaiting: it parks Fetch on an empty store under a live ctx, then cancels while parked, and now asserts the empty return. The pre-canceled variant never reached the wait code (an already-done ctx short-circuits the loop before parking), which is why the mutations survived it.
One caveat worth recording: after switching to context.AfterFunc (below), the lost-wakeup itself is no longer black-box reproducible. AfterFunc schedules its func only after cancel, and the predicate-check→notifyListAdd window is too narrow for that goroutine to interleave — so lock-before-Broadcast is correct-by-construction but not exercisable by a test. Confirmed 20k concurrent-cancel iterations do not catch an unlocked-Signal mutation, whereas removing the wakeup entirely deadlocks the watchdog reliably. So the new test guards "a parked Fetch is woken at all," and I documented that limitation in its doc comment rather than imply more.
| // lost, leaving Fetch blocked forever while holding bs.lock and thus | ||
| // deadlocking Pool.Close. |
There was a problem hiding this comment.
This names the wrong lock. sync.Cond.Wait releases bs.lock while parked (as line 164 says), and Pool.Close never takes bs.lock. The lock that deadlocks is the pool's:
NextRequestsholdsrp.lock.RLock()acrossFetch(pool.go:203/216).Closeblocks inrp.lock.Lock()(pool.go:284).- That pending writer blocks every new
Submit's RLock, so noInsertcan everSignal.
The #843 trace matches: the RWMutex.Lock address is the *Pool. Suggested wording: "...leaving Fetch parked forever inside Pool.NextRequests (which holds the pool's RLock), thus deadlocking Pool.Close." The test doc comment (batchstore_test.go:38-39) has the same wording.
There was a problem hiding this comment.
Fixed, in both the code comment and the test doc comment. Now: Fetch parks inside Pool.NextRequests (which holds the pool's RLock) → deadlocks Pool.Close on the pool's Lock; Cond.Wait releases bs.lock while parked.
| // Wait for a batch to become ready or for the context to be done. The loop | ||
| // re-checks the condition to tolerate a lost or spurious wakeup, including | ||
| // the case where ctx is already done when Fetch is entered. |
There was a problem hiding this comment.
This comment is stale since e64991a and credits the wrong mechanism.
- The guard at line 155 already handles "ctx is already done when Fetch is entered", so the loop never sees that case.
- A re-check loop can't recover a lost wakeup, because a lost wakeup means
Waitnever returns. What prevents it is takingbs.lockbeforeBroadcast(lines 175-177).
The loop really handles spurious or stale wakeups, e.g. a late Broadcast from an earlier Fetch's watcher. The risk is that someone trusts this comment and drops the lock around Broadcast, and no current test would catch that (see the comment on the test).
There was a problem hiding this comment.
Fixed. The comment now says the loop tolerates a spurious or stale wakeup (e.g. a late Broadcast from an earlier Fetch's AfterFunc), and that lock-before-Broadcast — not the loop — is what prevents a lost wakeup. The already-done-ctx case is handled by the loop predicate, so it never parks.
| bs.signal.Wait() | ||
|
|
||
| // Prefer a ready and full batch over a non-empty one | ||
| for len(bs.readyBatches) > 0 { |
There was a problem hiding this comment.
A behavior change to confirm. With the old if len(bs.readyBatches) > 0 { return bs.dequeueBatch() } fast path gone, Fetch drains ready batches that were emptied by Prune and falls straight through to cutting currentBatch, while ctx is still live.
Example: batch max size 2; Insert a,b,c; Prune a,b; Fetch(ctx with 2s timeout).
- Base: returns
[]immediately, then[c]~2s later, giving the batch time to fill. - This PR: returns
[c]in ~3µs.
In production this is the first NextRequests after a reconfig MemPool.Prune: one undersized batch (plus a BAF and a consensus decision) instead of waiting BatchTimeout. It isn't a correctness issue, so this may be fine, but it's worth confirming it's intended. A related minor point: draining many empty head batches is now O(n²) under the write lock in one call, because dequeueBatch copies the slice each time. That only matters with a small MaxMessageCount.
There was a problem hiding this comment.
Confirmed intended: skipping empty pruned ready-batches and falling through to cut currentBatch is deliberate. The O(n²) empty-batch drain only matters at very small MaxMessageCount (default 10000), so I've left dequeueBatch as is.
| // window is tiny, so we repeat enough times to hit it reliably: against the | ||
| // old code this deadlocks within ~1-2k iterations, so a few thousand gives a | ||
| // comfortable margin without making the passing run slow. | ||
| const iterations = 2000 |
There was a problem hiding this comment.
The comment says the old code "deadlocks within ~1-2k iterations, so a few thousand gives a comfortable margin", but the constant is exactly 2000.
- Against the base
Fetchunder-raceat GOMAXPROCS=4 (like CI), one run first deadlocked at iteration 2201, so it would pass. - Without
-race, the test caught the old bug in only ~1 of 35 runs, since hitting the window relies on-racescheduler randomization.
Consider raising the count and noting that it's only a reliable regression check under -race. A concurrent-cancel test (see the comment on line 60) would largely make this moot.
There was a problem hiding this comment.
Moot after the rewrite. The wake path is deterministic once Fetch is parked, so the iteration count is no longer a race-window tuning knob; kept a modest 500 with a comment explaining it only guards against rare scheduling where Fetch hadn't parked yet.
| // Generous watchdog: a correct Fetch on an already-canceled ctx returns in | ||
| // microseconds. On the buggy code Fetch never returns, so the first stuck | ||
| // iteration waits the full timeout, then t.Fatalf ends the test (the blocked | ||
| // Fetch and its inner goroutine are then torn down with the process). The |
There was a problem hiding this comment.
Nit: no "inner goroutine" is left behind. The old watcher returns right after its (lost) Signal, and at the PR head no watcher is started for a pre-cancelled ctx. Only the Fetch goroutine leaks, and it's parked harmlessly until the binary exits.
There was a problem hiding this comment.
The hand-rolled watcher goroutine is gone with AfterFunc, so the note no longer applies; the test comment was rewritten accordingly.
| finished := make(chan struct{}) | ||
| defer close(finished) | ||
|
|
||
| go func() { |
There was a problem hiding this comment.
Optional simplification: this goroutine + finished channel is the stdlib context.AfterFunc pattern; see ExampleAfterFunc_cond in context/example_test.go, which is exactly "cond wait until predicate or ctx done":
stop := context.AfterFunc(ctx, func() {
bs.lock.Lock()
defer bs.lock.Unlock()
bs.signal.Broadcast()
})
defer stop()
for len(bs.readyBatches) == 0 && ctx.Err() == nil {
bs.signal.Wait()
}It passed go test -race ./request/... and a 20k-iteration concurrent-cancel stress. It starts no goroutine unless ctx actually fires, so the line-155 guard (which duplicates the loop condition, and whose hot-path claim wasn't benchmarked) could go, and most of the comment block with it. Our benchmarks disagreed on cost: one measured it cheaper than the watcher goroutine, another slightly more expensive when each call gets a fresh ctx. Either way it's a small cost compared with prepareBatch's allocations.
There was a problem hiding this comment.
Adopted. Replaced the manual goroutine + finished channel with context.AfterFunc, and dropped the line-155 entry guard along with most of the comment block. go test -race ./request/... passes.
| bs.lock.Lock() | ||
| bs.signal.Broadcast() | ||
| bs.lock.Unlock() |
There was a problem hiding this comment.
Minor: if the watcher first reaches its select after Fetch has returned (finished closed) and ctx is done, select picks either case at random. Half the time it takes the write lock and broadcasts for nothing, which briefly blocks Insert, and can spuriously wake the next Fetch. The loop makes it harmless. Checking finished first (select { case <-finished: return; default: }) removes it, and so does context.AfterFunc + stop().
There was a problem hiding this comment.
Resolved by the AfterFunc switch: stop() on return removes the "watcher fires after Fetch returned" case, so there's no wasted Broadcast and no spurious wake of the next Fetch.
e64991a to
12fd062
Compare
|
On the three "not in the diff" items:
|
BatchStore.Fetch spawned a goroutine that called signal.Signal() on context cancellation without holding the lock. When the context was already cancelled (or cancelled in the window before Fetch reached signal.Wait()), the signal was delivered with no waiter registered and was therefore lost, leaving Wait() blocked forever while holding bs.lock. That stuck Fetch is reached from Pool.NextRequests under the pool's read lock, so Pool.Close (invoked from BatcherRole.Stop right after cancelBatch()) blocked forever on the pool's write lock, producing the deadlock reported in the flake. Fix the condition-variable usage: - the ctx.Done goroutine now acquires bs.lock before Signal(), so the signal cannot be delivered before Fetch is waiting on it; - Wait() is guarded by a predicate loop that also checks the context, so an already-cancelled context short-circuits the wait and a cancellation wake terminates the loop instead of re-blocking, while stray/late wakeups are re-checked rather than acted on. Add TestFetchWithCanceledContext, which deadlocks reliably against the old code and passes with the fix. Fixes hyperledger#843 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Signed-off-by: Hagar Meir <hagar.meir@ibm.com>
- Fetch spawns the ctx-watch goroutine only when it is actually about to Wait() (nothing ready and ctx still live), avoiding a goroutine create/teardown per call on the hot path where readyBatches is usually non-empty on entry. - The ctx-cancel wakeup now Broadcasts instead of Signals, so it stays correct if more than one Fetch ever waits on a store concurrently (NextRequests only holds the pool RLock; the single sequential caller is a convention, not a structural guarantee). - TestFetchWithCanceledContext: name the iteration count and watchdog as documented consts explaining the race window and the failure-path behavior. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Signed-off-by: Hagar Meir <hagar.meir@ibm.com>
Address review on the hyperledger#843 deadlock fix: - Replace the hand-rolled ctx-watch goroutine + `finished` channel in BatchStore.Fetch with `context.AfterFunc`. It schedules its wakeup only when ctx fires (so no goroutine on the ready-batch hot path), and stop() removes the select race where the watcher could Broadcast for nothing after Fetch returned. This also lets the redundant entry guard go. - The wakeup still Broadcasts under bs.lock: Wait releases bs.lock while parked, so taking it guarantees the wakeup lands only after Fetch is parked and cannot be lost. - Fix the comments that named the wrong lock: Fetch parks inside Pool.NextRequests (holding the pool's RLock), which deadlocks Pool.Close on the pool's Lock — Cond.Wait does not hold bs.lock while parked. Also correct the stale claim that the wait loop recovers a lost wakeup (it handles spurious/stale wakeups; lock-before-Broadcast prevents loss). - Rewrite the test as TestFetchCancelWhileWaiting: park Fetch on an empty store under a live ctx, then cancel while parked. The old pre-canceled variant never reached the wait code (an already-done ctx short-circuits the loop). The new test parks in Wait and requires the cancellation to wake it and return an empty batch; it deadlocks if the wakeup is removed. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Signed-off-by: Hagar Meir <hagar.meir@ibm.com>
NextRequests holds rp.lock.RLock across the blocking batchStore.Fetch, and Close/Halt take rp.lock.Lock, so stopping the pool while a NextRequests is in flight only works if the ctx is canceled first (which lets Fetch return and release the RLock). That invariant — relied on by BatcherRole.Stop/SoftStop, which call cancelBatch() before MemPool.Close()/Halt() — was undocumented; add it to the NextRequests and Close doc comments. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Signed-off-by: Hagar Meir <hagar.meir@ibm.com>
acfae39 to
b423468
Compare
Problem
Fixes #843 — a test flake caused by a deadlock when a batcher is stopped.
BatchStore.Fetchspawned a goroutine that calledsignal.Signal()oncontext cancellation without holding the lock. When the context was
already cancelled (or was cancelled in the window before
Fetchreachedsignal.Wait()), the signal was delivered with no waiter registered andwas therefore lost, leaving
Wait()blocked forever while holdingbs.lock.That stuck
Fetchis reached fromPool.NextRequestsunder the pool'sread lock, so
Pool.Close(invoked fromBatcherRole.Stopright aftercancelBatch()) blocked forever on the pool's write lock — the deadlockseen in the reported goroutine stacks:
RWMutex.LockinPool.Close(viaBatcherRole.Stop);Cond.WaitinBatchStore.Fetch(viaPool.NextRequests).Fix
Correct the condition-variable usage in
BatchStore.Fetch:ctx.Donegoroutine now acquiresbs.lockbeforeSignal(). BecauseWait()releases the lock atomically, the signal can only be deliveredafter
Fetchis actually waiting on it — it can no longer be lost.Wait()is guarded by a predicate loop that also checks the context:for len(bs.readyBatches) == 0 && ctx.Err() == nil. An already-cancelledcontext short-circuits the wait, a cancellation wake terminates the loop
(instead of re-blocking on a signal that will never come again), and any
stray/late wakeup is re-checked rather than acted on.
Both parts are required: lock-before-signal without the ctx-aware loop still
hangs on the ordinary timeout path; the loop without lock-before-signal still
loses the wakeup at entry.
Test
Added
TestFetchWithCanceledContext, which drivesFetchwith analready-cancelled context under a watchdog. It deadlocks reliably against
the old code (hangs within ~1–2k iterations) and passes with the fix.
Verification
go test -race ./request/...— passgo test -race ./node/batcher/...— passgo vet/gofmt— clean🤖 Generated with Claude Code