Repository navigation
fix: prevent pipeline deadlock when held items block reorder buffer progress - #251
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #251 +/- ##
==========================================
- Coverage 89.28% 89.26% -0.02%
==========================================
Files 119 119
Lines 57770 57914 +144
==========================================
+ Hits 51582 51699 +117
- Misses 6188 6215 +27 ☔ View full report in Codecov by Sentry. 🚀 New features to boost your workflow:
|
📝 WalkthroughWalkthroughDecompression, decoding, and compression steps were changed to track a 🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Warning Review ran into problems🔥 ProblemsGit: Failed to clone repository. Please run the Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
src/lib/unified_pipeline/base.rs (1)
4726-4736:⚠️ Potential issue | 🟠 MajorPotential data loss when Q6 push fails with
held_blocked=true.If
held_blockedis true (line 4652 already stored the original held item back), andq6_pushfails here, line 4733 overwritesheld_compressedwith the new batch, silently losing the original held item.Race scenario: Q6 fills between Priority 2 check (skipped when
held_blocked) and this push.Consider detecting and handling this edge case:
Suggested defensive check
Err((serial, batch)) => { + if held_blocked { + // Already holding something - can't store two batches + // This shouldn't happen if Q6 is adequately sized, but if it does, + // fail explicitly rather than lose the original held batch + state.set_error(io::Error::other( + "compress: Q6 full while held_blocked - cannot hold two batches" + )); + return StepResult::OutputFull; + } // Output full - hold the result for next attempt let heap_size = batch.estimate_heap_size(); *worker.held_compressed_mut() = Some((serial, batch, heap_size)); StepResult::OutputFull }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/unified_pipeline/base.rs` around lines 4726 - 4736, The current q6_push error path unconditionally sets worker.held_compressed_mut() which can overwrite an original held item when held_blocked is true; modify the Err((serial, batch)) branch in the match after state.q6_push to first check worker.held_compressed() and worker.held_blocked (or held_blocked flag) and if there is already Some(...) and held_blocked is true, avoid overwriting—either return StepResult::OutputFull without replacing the held item or merge/queue the new (serial, batch) into a safe buffer; ensure you still preserve heap_size via batch.estimate_heap_size() only when storing a new held item, and keep existing behavior for record_q7_push_progress/StepResult::Success triggered only on Ok.
🧹 Nitpick comments (1)
src/lib/unified_pipeline/base.rs (1)
6180-6213: Test coverage looks good for the happy path.Verifies the core fix: held out-of-order batch doesn't block compressing the next-seq batch. Correctly asserts serial 0 reaches Q6 while serial 5 remains held.
Missing edge case: test where Q6 becomes full during the compress operation (triggering the potential data loss path mentioned above). Consider adding:
#[test] fn test_compress_held_blocked_q6_full_no_data_loss() { let state = MockCompressState::new(1, 1000); // capacity=1 state.reorder.add_heap_bytes(600); // Fill Q6 with a dummy batch let dummy = CompressedBlockBatch::new(); assert!(state.q6.push((99, dummy)).is_ok()); // Put next-seq in Q5 let batch_0 = SerializedBatch { data: vec![0u8; 10], record_count: 1, secondary_data: None }; assert!(state.q5.push((0, batch_0)).is_ok()); let mut worker = MockCompressWorker::new(); worker.held_compressed = Some((5, CompressedBlockBatch::new(), 100)); let result = shared_try_step_compress(&state, &mut worker); // Should NOT lose serial 5 // (Behavior depends on fix applied) }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/unified_pipeline/base.rs` around lines 6180 - 6213, Add a new unit test that exercises the edge case where Q6 is full while compressing the next-seq batch to ensure held out-of-order batches are not lost: create MockCompressState with capacity 1 (MockCompressState::new(1, 1000)), mark heap bytes high via reorder.add_heap_bytes(600), push a dummy CompressedBlockBatch into state.q6, push the next-seq SerializedBatch into state.q5, set worker.held_compressed = Some((5, CompressedBlockBatch::new(), 100)) on a MockCompressWorker, call shared_try_step_compress(&state, &mut worker) and then assert that the held_compressed still contains serial 5 and that the next-seq batch was either moved to q6 (if space was freed) or appropriately handled per current behavior (use StepResult and state.q6.pop()/state.q6.peek() to validate no data loss).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@src/lib/unified_pipeline/bam.rs`:
- Around line 1824-1825: When held_blocked is true the code restores the old
held batch into worker.held_decompressed (worker.held_decompressed =
Some((serial, held, heap_size))) but then later on a failed push path can
overwrite that same slot (causing a dropped batch); change the flow so that when
held_blocked is true you do not reuse worker.held_decompressed for the newly
decompressed batch—store the fresh batch into a separate temporary slot or
insert it directly into the reorder/queue path instead of writing into
worker.held_decompressed. Locate the restore and reassign sites (the assignments
to worker.held_decompressed around the held_blocked branch and the push error
handling code that currently writes back at the later path) and implement a
second storage variable or direct reorder insertion for the new batch to avoid
overwriting the restored held batch.
- Around line 2163-2164: The decode path currently overwrites
worker.held_decoded (set via worker.held_decoded = Some((serial, held,
heap_size)) using serial/held/heap_size) which can lose an already-held decoded
batch when a later push fails and the same slot is replaced; instead preserve
both the existing held batch and the new decoded result: change the storage to
hold multiple entries (e.g., Vec or a small pair) or append the new tuple rather
than replace, and update the push-failure path that currently assigns to
worker.held_decoded to merge/append into that collection (mirror the decompress
fix behavior used elsewhere), ensuring both batches are retained and later
re-queued rather than silently overwritten.
---
Outside diff comments:
In `@src/lib/unified_pipeline/base.rs`:
- Around line 4726-4736: The current q6_push error path unconditionally sets
worker.held_compressed_mut() which can overwrite an original held item when
held_blocked is true; modify the Err((serial, batch)) branch in the match after
state.q6_push to first check worker.held_compressed() and worker.held_blocked
(or held_blocked flag) and if there is already Some(...) and held_blocked is
true, avoid overwriting—either return StepResult::OutputFull without replacing
the held item or merge/queue the new (serial, batch) into a safe buffer; ensure
you still preserve heap_size via batch.estimate_heap_size() only when storing a
new held item, and keep existing behavior for
record_q7_push_progress/StepResult::Success triggered only on Ok.
---
Nitpick comments:
In `@src/lib/unified_pipeline/base.rs`:
- Around line 6180-6213: Add a new unit test that exercises the edge case where
Q6 is full while compressing the next-seq batch to ensure held out-of-order
batches are not lost: create MockCompressState with capacity 1
(MockCompressState::new(1, 1000)), mark heap bytes high via
reorder.add_heap_bytes(600), push a dummy CompressedBlockBatch into state.q6,
push the next-seq SerializedBatch into state.q5, set worker.held_compressed =
Some((5, CompressedBlockBatch::new(), 100)) on a MockCompressWorker, call
shared_try_step_compress(&state, &mut worker) and then assert that the
held_compressed still contains serial 5 and that the next-seq batch was either
moved to q6 (if space was freed) or appropriately handled per current behavior
(use StepResult and state.q6.pop()/state.q6.peek() to validate no data loss).
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: b32ccbba-f358-4053-9756-63601cd199a9
📒 Files selected for processing (2)
src/lib/unified_pipeline/bam.rssrc/lib/unified_pipeline/base.rs
…rogress When a worker holds an out-of-order compressed/decompressed/decoded batch that is blocked by reorder buffer memory backpressure, the step function would return immediately at Priority 1 without ever reaching Priority 3 (pop new input). This meant the next-in-sequence batch — which the reorder buffer always admits — would sit unprocessed in the input queue, starving the downstream consumer and causing a total deadlock. The fix removes the early return at Priority 1 when the held item is blocked by backpressure. Instead, the function continues to Priority 3 to pop and process new work. At Priority 5, when the worker already holds a blocked item, the newly processed batch bypasses the backpressure check and is pushed directly to the output queue — the reorder buffer consumer will drain it, and this prevents the deadlock. Three instances fixed: - shared_try_step_compress (base.rs) — affects BAM + FASTQ pipelines - try_step_decompress (bam.rs) — BAM pipeline only - try_step_decode (bam.rs) — BAM pipeline only
5a04b4a to
46af9f0
Compare
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@src/lib/unified_pipeline/base.rs`:
- Around line 4717-4726: The current bypass admits any Q5 head when held_blocked
is true and can push out-of-order batches into Q6; modify the branch that
currently does match state.q6_push((serial, batch)) so it first verifies the
batch is the writer's expected next sequence (e.g., compare serial to the
writer's next_seq via an existing accessor like state.next_expected_seq() or
state.writer_next_seq()); if serial == next_seq then proceed with
state.q6_push((serial, batch)), otherwise preserve the non-next_seq batch
instead of force-admitting it (for example, store it back into
worker.held_compressed_mut() or requeue it to Q5/Q5-backlog) and return
StepResult::OutputFull, ensuring you do not bypass
state.write_reorder_can_proceed(serial) for arbitrary Q5 items and add an
assertion/log if you prefer to fail-fast when assumptions are violated.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: e81faf25-e3ca-4944-89bf-4c20bb43f958
📒 Files selected for processing (2)
src/lib/unified_pipeline/bam.rssrc/lib/unified_pipeline/base.rs
🚧 Files skipped from review as they are similar to previous changes (1)
- src/lib/unified_pipeline/bam.rs
…ine deadlock Remove can_proceed / write_reorder_can_proceed checks from Priority 1 (held batch push) and Priority 5 (new batch push) in Decompress, Decode, and Compress steps. Held batches now push unconditionally, subject only to physical queue capacity — the same pattern Process and Serialize already use. This eliminates three classes of deadlock: 1. All workers hold non-next_seq batches blocked by can_proceed at P1, nobody produces next_seq (the original PR #251 deadlock target) 2. held_blocked bypass + TOCTOU race causing double-hold aborts 3. push_with_toctou_retry trapping exclusive-step owner threads The held_blocked mechanism and push_with_toctou_retry are fully removed.
Decompress, Decode, and Compress steps checked can_proceed (memory backpressure) before pushing held batches at Priority 1 and new batches at Priority 5. When all workers held non-next_seq batches, can_proceed blocked everyone at Priority 1 and nobody could produce next_seq — deadlock. The held_blocked mechanism (PR #251) worked around this by letting workers bypass memory checks and produce a second batch, but that introduced TOCTOU races and double-hold failures. The push_with_toctou_retry loop (PR #253) fixed the TOCTOU race but trapped exclusive-step owner threads, causing new deadlocks. Remove can_proceed from Priority 1 and Priority 5 entirely. Held batches push unconditionally subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry.
Decompress, Decode, and Compress steps checked can_proceed (memory backpressure) before pushing held batches at Priority 1 and new batches at Priority 5. When all workers held non-next_seq batches, can_proceed blocked everyone at Priority 1 and nobody could produce next_seq — deadlock. The held_blocked mechanism (PR #251) worked around this by letting workers bypass memory checks and produce a second batch, but that introduced TOCTOU races and double-hold failures. The push_with_toctou_retry loop (PR #253) fixed the TOCTOU race but trapped exclusive-step owner threads, causing new deadlocks. Remove can_proceed from Priority 1 and Priority 5 entirely. Held batches push unconditionally subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry.
Decompress, Decode, and Compress steps checked can_proceed (memory backpressure) before pushing held batches at Priority 1 and new batches at Priority 5. When all workers held non-next_seq batches, can_proceed blocked everyone at Priority 1 and nobody could produce next_seq — deadlock. The held_blocked mechanism (PR #251) worked around this by letting workers bypass memory checks and produce a second batch, but that introduced TOCTOU races and double-hold failures. The push_with_toctou_retry loop (PR #253) fixed the TOCTOU race but trapped exclusive-step owner threads, causing new deadlocks. Remove can_proceed from Priority 1 and Priority 5 entirely. Held batches push unconditionally subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry.
Decompress, Decode, and Compress steps checked can_proceed (memory backpressure) before pushing held batches at Priority 1 and new batches at Priority 5. When all workers held non-next_seq batches, can_proceed blocked everyone at Priority 1 and nobody could produce next_seq — deadlock. The held_blocked mechanism (PR #251) worked around this by letting workers bypass memory checks and produce a second batch, but that introduced TOCTOU races and double-hold failures. The push_with_toctou_retry loop (PR #253) fixed the TOCTOU race but trapped exclusive-step owner threads, causing new deadlocks. Remove can_proceed from Priority 1 and Priority 5 entirely. Held batches push unconditionally subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry.
Decompress, Decode, and Compress steps checked can_proceed (memory backpressure) before pushing held batches at Priority 1 and new batches at Priority 5. When all workers held non-next_seq batches, can_proceed blocked everyone at Priority 1 and nobody could produce next_seq — deadlock. The held_blocked mechanism (PR #251) worked around this by letting workers bypass memory checks and produce a second batch, but that introduced TOCTOU races and double-hold failures. The push_with_toctou_retry loop (PR #253) fixed the TOCTOU race but trapped exclusive-step owner threads, causing new deadlocks. Remove can_proceed from Priority 1 and Priority 5 entirely. Held batches push unconditionally subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry.
Decompress, Decode, and Compress steps checked can_proceed (memory backpressure) before pushing held batches at Priority 1 and new batches at Priority 5. When all workers held non-next_seq batches, can_proceed blocked everyone at Priority 1 and nobody could produce next_seq — deadlock. The held_blocked mechanism (PR #251) worked around this by letting workers bypass memory checks and produce a second batch, but that introduced TOCTOU races and double-hold failures. The push_with_toctou_retry loop (PR #253) fixed the TOCTOU race but trapped exclusive-step owner threads, causing new deadlocks. Remove can_proceed from Priority 1 and Priority 5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at Priority 2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. This is deadlock-safe: Q1 is FIFO, so next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group 2t/16t and dedup 8t all complete without OOM on 29M-pair input that previously OOM-killed at all thread counts.
PR #242 (5df7e17) added write_reorder_can_proceed checks to the Compress step at P1 and P5, extending the can_proceed pattern that Decompress and Decode had since the initial commit. This made the pipeline deadlock-prone: when all workers hold non-next_seq batches, can_proceed blocks everyone at P1 and nobody produces next_seq. The held_blocked mechanism (PR #251) worked around this but introduced TOCTOU races. The push_with_toctou_retry loop (PR #253 v1) fixed the race but trapped exclusive-step owner threads. Remove can_proceed from P1 and P5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at P2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. Deadlock-safe because Q1 is FIFO: next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group 2t/16t and dedup 8t all complete without OOM on 29M-pair input that previously OOM-killed at all thread counts.
PR #242 (5df7e17) added write_reorder_can_proceed checks to the Compress step at P1 and P5, extending the can_proceed pattern that Decompress and Decode had since the initial commit. This made the pipeline deadlock-prone: when all workers hold non-next_seq batches, can_proceed blocks everyone at P1 and nobody produces next_seq. The held_blocked mechanism (PR #251) worked around this but introduced TOCTOU races. The push_with_toctou_retry loop (PR #253 v1) fixed the race but trapped exclusive-step owner threads. Remove can_proceed from P1 and P5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at P2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. Deadlock-safe because Q1 is FIFO: next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group 2t/16t and dedup 8t all complete without OOM on 29M-pair input that previously OOM-killed at all thread counts.
PR #242 (5df7e17) added write_reorder_can_proceed checks to the Compress step at P1 and P5, extending the can_proceed pattern that Decompress and Decode had since the initial commit. This made the pipeline deadlock-prone: when all workers hold non-next_seq batches, can_proceed blocks everyone at P1 and nobody produces next_seq. The held_blocked mechanism (PR #251) worked around this but introduced TOCTOU races. The push_with_toctou_retry loop (PR #253 v1) fixed the race but trapped exclusive-step owner threads. Remove can_proceed from P1 and P5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at P2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. Deadlock-safe because Q1 is FIFO: next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Wire FASTQ pipeline scheduler backpressure to write reorder state. The FASTQ pipeline hardcoded memory_high: false and memory_drained: true, so the scheduler never redirected threads to drain Compress/Write. Replace bounded retry loops in try_step_group finalization with non-blocking flush attempts. The previous MAX_RETRIES=10,000 with yield_now() gave only ~milliseconds of wall time, causing spurious "Failed to flush final groups" errors when Process takes >10ms per batch (e.g. duplex consensus on large groups). On flush failure, the thread now returns to the scheduler, does useful Process work to drain Q4, and the scheduler retries Group via the is_finished() check. The deadlock detector handles true deadlocks. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group and dedup at 2/4/8/16 threads show no wall-clock regressions vs baseline (f0b83c6). Dedup at 8 threads is 5% faster with 20% less RSS.
PR #242 (5df7e17) added write_reorder_can_proceed checks to the Compress step at P1 and P5, extending the can_proceed pattern that Decompress and Decode had since the initial commit. This made the pipeline deadlock-prone: when all workers hold non-next_seq batches, can_proceed blocks everyone at P1 and nobody produces next_seq. The held_blocked mechanism (PR #251) worked around this but introduced TOCTOU races. The push_with_toctou_retry loop (PR #253 v1) fixed the race but trapped exclusive-step owner threads. Remove can_proceed from P1 and P5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at P2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. Deadlock-safe because Q1 is FIFO: next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Wire FASTQ pipeline scheduler backpressure to write reorder state. The FASTQ pipeline hardcoded memory_high: false and memory_drained: true, so the scheduler never redirected threads to drain Compress/Write. Replace bounded retry loops in try_step_group finalization with non-blocking flush attempts. The previous MAX_RETRIES=10,000 with yield_now() gave only ~milliseconds of wall time, causing spurious "Failed to flush final groups" errors when Process takes >10ms per batch (e.g. duplex consensus on large groups). On flush failure, the thread now returns to the scheduler, does useful Process work to drain Q4, and the scheduler retries Group via the is_finished() check. The deadlock detector handles true deadlocks. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group and dedup at 2/4/8/16 threads show no wall-clock regressions vs baseline (f0b83c6). Dedup at 8 threads is 5% faster with 20% less RSS.
PR #242 (5df7e17) added write_reorder_can_proceed checks to the Compress step at P1 and P5, extending the can_proceed pattern that Decompress and Decode had since the initial commit. This made the pipeline deadlock-prone: when all workers hold non-next_seq batches, can_proceed blocks everyone at P1 and nobody produces next_seq. The held_blocked mechanism (PR #251) worked around this but introduced TOCTOU races. The push_with_toctou_retry loop (PR #253 v1) fixed the race but trapped exclusive-step owner threads. Remove can_proceed from P1 and P5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at P2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. Deadlock-safe because Q1 is FIFO: next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Wire FASTQ pipeline scheduler backpressure to write reorder state. The FASTQ pipeline hardcoded memory_high: false and memory_drained: true, so the scheduler never redirected threads to drain Compress/Write. Replace bounded retry loops in try_step_group finalization with non-blocking flush attempts. The previous MAX_RETRIES=10,000 with yield_now() gave only ~milliseconds of wall time, causing spurious "Failed to flush final groups" errors when Process takes >10ms per batch (e.g. duplex consensus on large groups). On flush failure, the thread now returns to the scheduler, does useful Process work to drain Q4, and the scheduler retries Group via the is_finished() check. The deadlock detector handles true deadlocks. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group and dedup at 2/4/8/16 threads show no wall-clock regressions vs baseline (f0b83c6). Dedup at 8 threads is 5% faster with 20% less RSS.
PR #242 (5df7e17) added write_reorder_can_proceed checks to the Compress step at P1 and P5, extending the can_proceed pattern that Decompress and Decode had since the initial commit. This made the pipeline deadlock-prone: when all workers hold non-next_seq batches, can_proceed blocks everyone at P1 and nobody produces next_seq. The held_blocked mechanism (PR #251) worked around this but introduced TOCTOU races. The push_with_toctou_retry loop (PR #253 v1) fixed the race but trapped exclusive-step owner threads. Remove can_proceed from P1 and P5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at P2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. Deadlock-safe because Q1 is FIFO: next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Wire FASTQ pipeline scheduler backpressure to write reorder state. The FASTQ pipeline hardcoded memory_high: false and memory_drained: true, so the scheduler never redirected threads to drain Compress/Write. Replace bounded retry loops in try_step_group finalization with non-blocking flush attempts. The previous MAX_RETRIES=10,000 with yield_now() gave only ~milliseconds of wall time, causing spurious "Failed to flush final groups" errors when Process takes >10ms per batch (e.g. duplex consensus on large groups). On flush failure, the thread now returns to the scheduler, does useful Process work to drain Q4, and the scheduler retries Group via the is_finished() check. The deadlock detector handles true deadlocks. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group and dedup at 2/4/8/16 threads show no wall-clock regressions vs baseline (f0b83c6). Dedup at 8 threads is 5% faster with 20% less RSS.
PR #242 (5df7e17) added write_reorder_can_proceed checks to the Compress step at P1 and P5, extending the can_proceed pattern that Decompress and Decode had since the initial commit. This made the pipeline deadlock-prone: when all workers hold non-next_seq batches, can_proceed blocks everyone at P1 and nobody produces next_seq. The held_blocked mechanism (PR #251) worked around this but introduced TOCTOU races. The push_with_toctou_retry loop (PR #253 v1) fixed the race but trapped exclusive-step owner threads. Remove can_proceed from P1 and P5 entirely. Held batches push unconditionally at P1 subject only to physical queue capacity, matching the pattern Process and Serialize already use. Remove held_blocked and push_with_toctou_retry. Restore memory backpressure at P2 via is_memory_high() on reorder buffer state. This gates new work (not held pushes) when reorder buffers are large, preventing unbounded memory growth. Deadlock-safe because Q1 is FIFO: next_seq has already been produced and is in the ArrayQueue, reorder buffer, or held slot; exclusive steps continue draining, memory drops, and workers resume. Wire FASTQ pipeline scheduler backpressure to write reorder state. The FASTQ pipeline hardcoded memory_high: false and memory_drained: true, so the scheduler never redirected threads to drain Compress/Write. Replace bounded retry loops in try_step_group finalization with non-blocking flush attempts. The previous MAX_RETRIES=10,000 with yield_now() gave only ~milliseconds of wall time, causing spurious "Failed to flush final groups" errors when Process takes >10ms per batch (e.g. duplex consensus on large groups). On flush failure, the thread now returns to the scheduler, does useful Process work to drain Q4, and the scheduler retries Group via the is_finished() check. The deadlock detector handles true deadlocks. Verified on EC2 (c7g.4xlarge, 30 GB RAM): group and dedup at 2/4/8/16 threads show no wall-clock regressions vs baseline (f0b83c6). Dedup at 8 threads is 5% faster with 20% less RSS.
Summary
shared_try_step_compress(base.rs, BAM + FASTQ),try_step_decompress(bam.rs),try_step_decode(bam.rs)Deadlock chain (before fix)
serial != next_seq— blocked bycan_proceed()returning falseOutputFull/false), never reaching Priority 3next_seqbatch sits unprocessed in the input queue — no worker can pop itFix
held_blocked = trueinstead of returning earlyheld_blocked(must still try to process new work)held_blockedis true, bypass thecan_proceedbackpressure check for the freshly processed batch — push it directly to the output queue to break the deadlockTest plan
test_compress_held_item_does_not_block_next_seq— reproduces the exact deadlock scenario (RED before fix, GREEN after)test_compress_held_next_seq_advances_normally— verifies held next-seq items still advance correctlytest_compress_no_held_empty_q5_returns_input_empty— verifies normal behavior when no held itemscargo ci-fmtcleancargo ci-lintclean