Skip to content

request: fix deadlock between Pool.Close and BatchStore.Fetch (#843) - #1211

Merged
HagarMeir merged 4 commits into
hyperledger:mainfrom
HagarMeir:fix/843-fetch-close-deadlock
Sep 30, 2026
Merged

HagarMeir merged 4 commits into
hyperledger:mainfrom
HagarMeir:fix/843-fetch-close-deadlock

Conversation

@HagarMeir

Copy link
Copy Markdown
Contributor

Problem

Fixes #843 — a test flake caused by a deadlock when a batcher is stopped.

BatchStore.Fetch spawned a goroutine that called signal.Signal() on
context cancellation without holding the lock. When the context was
already cancelled (or was 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 — the deadlock
seen in the reported goroutine stacks:

  • one goroutine on RWMutex.Lock in Pool.Close (via BatcherRole.Stop);
  • one goroutine on Cond.Wait in BatchStore.Fetch (via Pool.NextRequests).

Fix

Correct the condition-variable usage in BatchStore.Fetch:

  1. The ctx.Done goroutine now acquires bs.lock before Signal(). Because
    Wait() releases the lock atomically, the signal can only be delivered
    after Fetch is actually waiting on it — it can no longer be lost.
  2. Wait() is guarded by a predicate loop that also checks the context:
    for len(bs.readyBatches) == 0 && ctx.Err() == nil. An already-cancelled
    context 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 drives Fetch with an
already-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/... — pass
  • go test -race ./node/batcher/... — pass
  • go vet / gofmt — clean

🤖 Generated with Claude Code

@HagarMeir
HagarMeir force-pushed the fix/843-fetch-close-deadlock branch 5 times, most recently from 849dddf to c91739f Compare September 17, 2026 17:29
@HagarMeir
HagarMeir force-pushed the fix/843-fetch-close-deadlock branch 2 times, most recently from 87c3851 to a55f277 Compare September 24, 2026 10:00

@tock-ibm tock-ibm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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() == nil is false, so Wait() is skipped and Fetch returns promptly — the deadlock is gone.
  • ctx cancels in the window between the condition check and parking in Wait(): no lost signal, because Fetch holds the write lock continuously across that window and the ctx-watch goroutine cannot Lock() to Signal() until Wait() 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

Comment thread request/batchstore.go Outdated
// lost, leaving Fetch blocked forever while holding bs.lock and thus
// deadlocking Pool.Close.
bs.lock.Lock()
bs.signal.Signal()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Comment thread request/batchstore.go Outdated
@@ -157,21 +157,27 @@ func (bs *BatchStore) Fetch(ctx context.Context) ([]interface{}, []interface{})
go func() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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().

Comment thread request/batchstore_test.go Outdated
func TestFetchWithCanceledContext(t *testing.T) {
sugaredLogger := testutil.CreateLogger(t, 0)

for i := 0; i < 2000; i++ {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Two small test-quality points:

  • The 2000 iteration 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.Fatalf fires from the main test goroutine while the blocked Fetch goroutine (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.

@HagarMeir
HagarMeir force-pushed the fix/843-fetch-close-deadlock branch 2 times, most recently from 9d4ebbd to 9f7cc5a Compare September 28, 2026 11:50
@tock-ibm
tock-ibm force-pushed the fix/843-fetch-close-deadlock branch from 2de67a2 to e64991a Compare September 29, 2026 12:52

@tock-ibm tock-ibm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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:

  • NextRequests holds rp.lock.RLock() across a blocking Fetch (pool.go:216), so Pool.Close only works because callers cancel the ctx first (cancelBatch() before MemPool.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) calls cancelContextFunc(); stateCond.Broadcast() without holding stateCond.L, while PopOrWait/Put check the ctx under L and then Wait(). We reproduced it in 3/3 runs: Collator.Stop hangs, so Assembler.Stop/SoftStop hangs while holding a.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 calls onDelete (which releases the semaphore permit and decrements the size) and never clears keys2Batches. After a reconfig prune in batching mode, pruned IDs are rejected as duplicates, and capacity leaks until Submit times out. Reproduced in a scratch test.

🤖 Reviewed with Claude Code

Comment thread request/batchstore_test.go Outdated
bs := NewBatchStore(100, 100*8, func(string) {}, sugaredLogger)

ctx, cancel := context.WithCancel(context.Background())
cancel() // context is already canceled before Fetch is called

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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 ./request package still pass.
  • Deleting the watcher entirely: this test still passes. Only TestBatchStore catches it, by hanging until the go test timeout.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

Comment thread request/batchstore.go Outdated
Comment on lines +165 to +166
// lost, leaving Fetch blocked forever while holding bs.lock and thus
// deadlocking Pool.Close.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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:

  • NextRequests holds rp.lock.RLock() across Fetch (pool.go:203/216).
  • Close blocks in rp.lock.Lock() (pool.go:284).
  • That pending writer blocks every new Submit's RLock, so no Insert can ever Signal.

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

Comment thread request/batchstore.go Outdated
Comment on lines +182 to +184
// 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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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 Wait never returns. What prevents it is taking bs.lock before Broadcast (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).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

Comment thread request/batchstore.go
bs.signal.Wait()

// Prefer a ready and full batch over a non-empty one
for len(bs.readyBatches) > 0 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

Comment thread request/batchstore_test.go Outdated
Comment on lines +45 to +48
// 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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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 Fetch under -race at 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 -race scheduler 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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

Comment thread request/batchstore_test.go Outdated
// 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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

The hand-rolled watcher goroutine is gone with AfterFunc, so the note no longer applies; the test comment was rewritten accordingly.

Comment thread request/batchstore.go Outdated
finished := make(chan struct{})
defer close(finished)

go func() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

Comment thread request/batchstore.go Outdated
Comment on lines +175 to +177
bs.lock.Lock()
bs.signal.Broadcast()
bs.lock.Unlock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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().

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

@HagarMeir

Copy link
Copy Markdown
Contributor Author

On the three "not in the diff" items:

HagarMeir and others added 4 commits September 30, 2026 13:53
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>
@tock-ibm
tock-ibm force-pushed the fix/843-fetch-close-deadlock branch from acfae39 to b423468 Compare September 30, 2026 10:53
@HagarMeir
HagarMeir merged commit 1a2c81b into hyperledger:main Sep 30, 2026
16 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Test flake - deadlock close and fetch from pool

2 participants