From 1ba16c61556f9e1f0cd0223feedb62a0d5226bad Mon Sep 17 00:00:00 2001 From: tfenne Date: Tue, 16 Jun 2026 12:02:40 -0600 Subject: [PATCH 1/2] Fix lost-wakeup deadlock in parallel BGZF reader shutdown BgzfMtReader's destructor stored mStop and then called notify_all() on the produce/decompress/consume condition variables without holding their mutexes. mStop is part of the wait predicate for all three CVs, so mutating it outside the locks defeats the usual guarantee that a notification cannot be lost between a waiter's predicate check and its park. A reader or decompress worker that had evaluated its predicate as false but had not yet parked could miss the shutdown notification and block forever; the destructor then hung in join() and the whole process deadlocked while tearing down the BGZF input reader. Publish mStop while holding all three mutexes before notifying, so any waiter sitting in the predicate-check-to-park window is serialized against the stop. The teardown runs whenever a BgzfMtReader is destroyed: the evaluator constructs and discards several per run, and the main read pass tears one down at the end, so the hang could surface either at startup or at the very end of a run. It reproduced intermittently at low frequency -- on the order of one in a thousand runs -- with BGZF-compressed inputs, and was independent of output settings (it is purely on the input-decompression side). An isolated stress harness that loops construct/read/destroy confirms both halves: the stock teardown hangs within tens of thousands of cycles with the captured stack (decompress worker parked in condition_variable::wait, main thread blocked in ~BgzfMtReader -> join), while the fixed teardown runs cleanly across millions of cycles. This affects every release since the parallel BGZF reader was introduced. --- src/bgzf.h | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/src/bgzf.h b/src/bgzf.h index 9b02852a..15fa8547 100644 --- a/src/bgzf.h +++ b/src/bgzf.h @@ -66,7 +66,18 @@ class BgzfMtReader { } ~BgzfMtReader() { - mStop = true; + // Publish mStop while holding every CV's mutex before notifying. mStop is a + // wait predicate for all three CVs; a thread that has evaluated its predicate + // as false but not yet parked still holds that mutex, so taking all three here + // serializes the stop+notify against that gap. Storing mStop outside the locks + // (as before) lets notify_all() slip into that window and be missed, leaving a + // worker parked forever and the join() below hung — an intermittent deadlock. + { + std::lock_guard l1(mProduceMtx); + std::lock_guard l2(mDecompMtx); + std::lock_guard l3(mConsumeMtx); + mStop = true; + } mDecompCv.notify_all(); mProduceCv.notify_all(); mConsumeCv.notify_all(); From d0d6aa4f87e6f292528bc0a76dd2c110583dcbf2 Mon Sep 17 00:00:00 2001 From: tfenne Date: Tue, 16 Jun 2026 12:46:22 -0600 Subject: [PATCH 2/2] Fix remaining lost-wakeup races on BgzfMtReader's internal CVs The previous commit fixed the shutdown (mStop) lost wakeup, but the same pattern remained at the producer/consumer hand-off points: each stage published a slot's new state with an atomic store and then notified the next stage's condition variable without holding that CV's mutex. A waiter that had evaluated its predicate as false but not yet parked could miss the notification and block forever. Continued hammering surfaced a second deadlock with BGZF input, again at low frequency (on the order of one in a few thousand runs): the reader thread blocked forever in BgzfMtReader::read() on mConsumeCv while every downstream worker sat idle in the backpressure wait and the main thread hung joining the reader. The lost wakeup was the READY/DONE publication for the slot the consumer was waiting on. Publish each slot transition under the mutex of the CV whose predicate reads it, before notifying: - read(): FREE under mProduceMtx (readerLoop waits on mProduceCv) - readerLoop(): COMPRESSED under mDecompMtx (workers wait on mDecompCv) - decompWorker(): READY under mConsumeMtx (consumer waits on mConsumeCv) - markDone(): DONE under mConsumeMtx (terminal EOF notify to consumer) No thread holds more than one of these mutexes at a time, so there is no new lock-ordering hazard, and the locks sit outside the per-block decompress work so the added cost is negligible. --- src/bgzf.h | 32 +++++++++++++++++++++++++++----- 1 file changed, 27 insertions(+), 5 deletions(-) diff --git a/src/bgzf.h b/src/bgzf.h index 15fa8547..849f39b2 100644 --- a/src/bgzf.h +++ b/src/bgzf.h @@ -106,7 +106,12 @@ class BgzfMtReader { mConsumeOffset += tocopy; if (mConsumeOffset >= s.decompLen) { - s.state.store(FREE, std::memory_order_release); + // Publish FREE under mProduceMtx so readerLoop, which waits on + // mProduceCv for this slot to free up, cannot miss the wakeup. + { + std::lock_guard lk(mProduceMtx); + s.state.store(FREE, std::memory_order_release); + } mConsumeOffset = 0; mConsumeIdx++; mProduceCv.notify_one(); @@ -145,8 +150,13 @@ class BgzfMtReader { } s.compLen = bsize; - s.state.store(COMPRESSED, std::memory_order_release); - mProduceIdx++; + // Publish the COMPRESSED slot (state + mProduceIdx, which claimSlot scans) + // under mDecompMtx so a decompWorker parked on mDecompCv cannot miss it. + { + std::lock_guard lk(mDecompMtx); + s.state.store(COMPRESSED, std::memory_order_release); + mProduceIdx++; + } mDecompCv.notify_all(); } } @@ -173,7 +183,12 @@ class BgzfMtReader { int ret = isal_inflate_stateless(&ist); target->decompLen = (ret == ISAL_DECOMP_OK) ? (int)ist.total_out : 0; - target->state.store(READY, std::memory_order_release); + // Publish READY under mConsumeMtx so the consumer parked in read() on + // mConsumeCv cannot miss the wakeup for this slot. + { + std::lock_guard lk(mConsumeMtx); + target->state.store(READY, std::memory_order_release); + } mConsumeCv.notify_one(); } } @@ -194,7 +209,14 @@ class BgzfMtReader { } void markDone(Slot& s) { - s.state.store(DONE, std::memory_order_release); + // Publish DONE under mConsumeMtx before notifying: the consumer in read() + // waits on mConsumeCv for READY/DONE, and at EOF this is the terminal + // notification (no further producer follows), so a lost wakeup here would + // strand the reader thread in read() forever. + { + std::lock_guard lk(mConsumeMtx); + s.state.store(DONE, std::memory_order_release); + } mConsumeCv.notify_all(); mDecompCv.notify_all(); }