Repository navigation
perf(sort): implement N+2 worker-pool model for parallel sort I/O - #242
Conversation
|
@coderabbitai review |
1 similar comment
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughThis pull request introduces a phase-aware worker pool infrastructure for parallel BGZF block compression and decompression. A new 🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
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 |
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 9
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
src/lib/sort/read_ahead.rs (1)
276-283:⚠️ Potential issue | 🟠 Major
take_error()is not actually unified yet.
src/lib/sort/raw.rs:1708-1720now relies ontake_error()as the post-iteration failure check, but theReadAheadbranch still hardcodesNonewhile the background thread converts read failures into EOF. Corrupt or truncated input on the non-pool path is still silent.Also applies to: 503-511, 518-529
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/read_ahead.rs` around lines 276 - 283, The ReadAhead branch currently treats errors as EOF by always sending an empty Vec (hardcoded None) after an Err in the reader loop; instead, invoke take_error() on the ReadAhead reader and propagate that actual error state to the consumer so corrupt/truncated input isn't silenced: when Err(e) occurs in the loop (the branch that logs "Error reading raw BAM record"), after sending any accumulated batch (tx.send(batch)) call the reader's take_error() and send/encode its result (rather than a hardcoded empty/None) down tx so downstream logic (raw.rs's post-iteration check which calls take_error()) sees the real error; update the other similar branches (around the other two spots mentioned) to follow the same pattern using take_error() and tx/batch symbols.
🧹 Nitpick comments (1)
tests/integration/test_sort_correctness.rs (1)
179-401: Collapse this matrix withrstest.These cases only vary by
(order, threads, memory), so keeping them as separate tests makes the matrix harder to extend and update. A parameterizedrstestwould keep thesamtoolssetup in one place and make failures much easier to compare.Based on learnings, "Write unit tests alongside source code in
src/and integration tests intests/integration/. Userstestfor parameterized tests,proptestfor property-based testing, andcriterionfor benchmarks inbenches/".🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/integration/test_sort_correctness.rs` around lines 179 - 401, Collapse the repetitive test matrix into a single parameterized rstest: replace the many #[test] functions (e.g., test_sort_coordinate_in_memory, test_sort_coordinate_with_spills, test_sort_queryname_lex_in_memory, test_sort_queryname_natural_with_spills, test_sort_template_coordinate_many_spills_t4, test_sort_coordinate_consistent_across_threads, etc.) with one rstest that iterates the tuples of (order, threads, memory, maybe input_size) and runs the same body that calls sort_bam, verify_sorted and count_records; keep the existing samtools_available() guard and TempDir setup in the single test, and assert sortedness and record counts per parameter set so failures are easier to compare and the matrix is extendable.
🤖 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/bam_io.rs`:
- Around line 779-794: The code in create_raw_bam_reader_pool_integrated
currently opens a File, reads the BAM header with BgzfReader/noodles which
consumes bytes, then calls file.seek(SeekFrom::Start(0)) — this fails for
non-seekable inputs (FIFOs, /dev/stdin). Replace the seek-based approach with
the TeeReader + ChainedReader pattern used in create_bam_reader_for_pipeline:
wrap the underlying readable (File or stdin) with a TeeReader that writes
header-bytes to an in-memory buffer while noodles reads the header, then build a
ChainedReader that first serves the buffered header bytes and then the original
reader for pool workers; ensure you stop calling file.seek and instead pass the
ChainedReader to BgzfReader/noodles and to the pool so non-seekable inputs work
correctly.
In `@src/lib/sort/raw.rs`:
- Around line 928-970: The drain loop currently skips entirely when
pending_total >= max_pending_blocks, causing deadlock if decompressed_chunks
holds the target source's next_serial; change the logic so you still
pop/decompress and install any block that satisfies the target (i.e.
block.source_id == target_source && block.serial == target.next_serial) even
when the global pending cap is reached. Concretely, modify the condition around
the initial if/while that drains self.decompressed_chunks (symbols:
pending_total, max_pending_blocks, decompressed_chunks, per_source,
target_source, next_serial, pending) so that either the drain is unconditional
until target.remaining() > 0 or add a secondary path that allows popping a block
when it matches the target's next_serial, then update pending counts (recompute
new_total) and break when cap is reached otherwise. Ensure the target check that
follows can see the installed block before parking.
- Around line 1058-1064: The current read_key_from_source silently treats any
UnexpectedEof from KK::read_from (RawSortKey::read_from) as Ok(None); instead,
when KK::read_from returns an UnexpectedEof, use the SourceReadAdapter's
read_exact_or_eof helper to distinguish a clean EOF from a partial/truncated
read and only return Ok(None) for a clean EOF—propagate an error for partial
reads (i.e., return Err with context "error reading key from source {source_id}:
..."). Ensure the fix touches read_key_from_source and uses
SourceReadAdapter.read_exact_or_eof so truncated spill files produce an error
rather than being silently skipped.
In `@src/lib/sort/read_ahead.rs`:
- Around line 216-223: The threaded constructor from_reader (and its helper
from_reader_with_batch_size) exposes a threaded path that can deadlock when
passed a RawBamReader wrapping a PooledInputStream; change these APIs to not be
publicly callable from arbitrary readers: remove the pub visibility from
from_reader and from_reader_with_batch_size (and any other public factory that
spawns the background thread referenced at lines 228-236 and 404-431), leaving
only internal/private constructors that create the background thread, and/or add
an explicit runtime/type check that rejects RawBamReader<PooledInputStream>
(using the RawBamReader and PooledInputStream type names) with a clear error
message if you prefer to keep them public; also update or add documentation near
BATCH_SIZE and CHANNEL_BUFFER_SIZE to note the constraint.
In `@src/lib/sort/worker_pool.rs`:
- Around line 1023-1028: The compressor selection currently reads from
shared.phase and can race with set_phase; instead bind the compressor to the
dispatched step (SortStep::CompressSpill or SortStep::CompressOutput) or store
the intended compression target on the CompressJob when the job is created.
Modify the dispatch branch that calls Self::try_compress(shared, worker) so it
passes or sets the compressor selection derived from the matched SortStep (or
attach a target enum/level field to CompressJob at creation), then update
try_compress / CompressJob handling to use that stored target instead of reading
shared.phase at compress time; ensure references to shared.phase are removed
where try_compress/CompressJob decide compression parameters.
- Around line 1192-1193: Replace the panicking expect() calls around
decompress_block(&block, &mut worker.decompressor) with proper error handling:
on decompression failure capture the error and send it over the shared
error/completion channel used by the worker pool (the same channel the main
thread watches), ensure the worker publishes a completion signal (or
closes/sends a final message) so the main thread is unblocked, and terminate the
worker gracefully instead of panicking; apply the same change for the other
occurrence (the second expect) so malformed BGZF data never kills a thread
silently.
- Around line 1508-1510: compress_result_channel currently returns a bounded
channel and worker code uses blocking send() (see compress_result_channel and
the worker send sites), which can block shutdown when receivers stop draining;
change the send logic to avoid blocking: create the bounded channel as before
but replace blocking send() calls in the worker (the sites flagged around the
compress_result_channel usage) with a non-blocking alternative (e.g., try_send
or a send with timeout/select that detects a closed/full channel) and handle the
Err by dropping/logging the result and continuing so workers never block waiting
for a receiver during shutdown; ensure you also check for Receiver closed errors
and return/exit the worker loop gracefully.
- Around line 1-41: Add the crate-level attribute to deny unsafe code at the top
of this module: insert #![deny(unsafe_code)] immediately after the module doc
comment/header (above the existing use/imports such as read_raw_blocks,
InlineBgzfCompressor, decompress_block) so the file enforces the repository rule
disallowing unsafe code; ensure it appears before any code or use statements in
src/lib/sort/worker_pool.rs.
---
Outside diff comments:
In `@src/lib/sort/read_ahead.rs`:
- Around line 276-283: The ReadAhead branch currently treats errors as EOF by
always sending an empty Vec (hardcoded None) after an Err in the reader loop;
instead, invoke take_error() on the ReadAhead reader and propagate that actual
error state to the consumer so corrupt/truncated input isn't silenced: when
Err(e) occurs in the loop (the branch that logs "Error reading raw BAM record"),
after sending any accumulated batch (tx.send(batch)) call the reader's
take_error() and send/encode its result (rather than a hardcoded empty/None)
down tx so downstream logic (raw.rs's post-iteration check which calls
take_error()) sees the real error; update the other similar branches (around the
other two spots mentioned) to follow the same pattern using take_error() and
tx/batch symbols.
---
Nitpick comments:
In `@tests/integration/test_sort_correctness.rs`:
- Around line 179-401: Collapse the repetitive test matrix into a single
parameterized rstest: replace the many #[test] functions (e.g.,
test_sort_coordinate_in_memory, test_sort_coordinate_with_spills,
test_sort_queryname_lex_in_memory, test_sort_queryname_natural_with_spills,
test_sort_template_coordinate_many_spills_t4,
test_sort_coordinate_consistent_across_threads, etc.) with one rstest that
iterates the tuples of (order, threads, memory, maybe input_size) and runs the
same body that calls sort_bam, verify_sorted and count_records; keep the
existing samtools_available() guard and TempDir setup in the single test, and
assert sortedness and record counts per parameter set so failures are easier to
compare and the matrix is extendable.
🪄 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: 383e5054-e498-4f6c-896e-b49493d30bad
⛔ Files ignored due to path filters (1)
Cargo.lockis excluded by!**/*.lock
📒 Files selected for processing (18)
.gitignoreCargo.tomlcrates/fgumi-bgzf/src/lib.rscrates/fgumi-bgzf/src/writer.rssrc/commands/common.rssrc/commands/sort.rssrc/lib/bam_io.rssrc/lib/mod.rssrc/lib/sort/bgzf_io.rssrc/lib/sort/mod.rssrc/lib/sort/pooled_bam_writer.rssrc/lib/sort/pooled_chunk_writer.rssrc/lib/sort/raw.rssrc/lib/sort/read_ahead.rssrc/lib/sort/worker_pool.rssrc/lib/system.rstests/integration/main.rstests/integration/test_sort_correctness.rs
45653c6 to
288a49e
Compare
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #242 +/- ##
==========================================
- Coverage 88.98% 88.87% -0.11%
==========================================
Files 114 119 +5
Lines 55488 58003 +2515
==========================================
+ Hits 49375 51553 +2178
- Misses 6113 6450 +337 ☔ View full report in Codecov by Sentry. 🚀 New features to boost your workflow:
|
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 5
♻️ Duplicate comments (2)
src/lib/bam_io.rs (1)
777-813:⚠️ Potential issue | 🟠 MajorStill seek-based for non-stdin pipes.
is_stdin_path()only catches-and/dev/stdin. A FIFO or process-substitution path still enters theFile + seek(0)branch, so header parsing consumes bytes and rewind fails. Use the sameTeeReader + ChainedReaderpath for any non-seekable input, not just stdin.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/bam_io.rs` around lines 777 - 813, The code currently treats only is_stdin_path(path_ref) as non-replayable; instead detect seekability on the opened handle and use the TeeReader+ChainedReader path for any non-seekable input. After opening File::open(path_ref) use file.seek(SeekFrom::Current(0)) (or attempt a small seek) to determine if seeking is supported; if seek returns Err (non-seekable), wrap the opened file in TeeReader -> BgzfReader -> noodles::bam::io::Reader to read the header and then into_parts() -> ChainedReader (same pattern as the stdin branch); otherwise keep the existing seek-to-0 and buffered file path. Update the header/reader construction code around the file/open, header read, and rewind (file.seek(...)) logic (references: File::open, file.seek, BgzfReader, noodles::bam::io::Reader::from, TeeReader::new, ChainedReader::new, into_parts).src/lib/sort/bgzf_io.rs (1)
112-125:⚠️ Potential issue | 🟠 MajorStill unbounded on writer-side reordering.
Draining
result_rxstraight intoreorder_bufremoves the only hard backpressure cap. If an early serial stalls, later blocks keep piling up in theBTreeMap, andBufferPool::checkout()falls back to fresh allocations once reuse lags.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/bgzf_io.rs` around lines 112 - 125, The writer-side reordering is unbounded because we always drain result_rx into reorder_buf, so if an early serial stalls the BTreeMap grows without backpressure; introduce a MAX_PENDING cap and refuse to insert when reorder_buf.len() >= MAX_PENDING by pausing intake until the writer advances (i.e., block/wait on processing the next_expected serials or otherwise process incoming results until next_expected appears) instead of blindly inserting; implement this check around the result handling logic that uses result_rx, reorder_buf, result.serial, next_expected and writer.write_all, and use BufferPool::checkout semantics to ensure we stop consuming more results when the map is full to restore backpressure.
🧹 Nitpick comments (3)
src/lib/system.rs (1)
29-30: Doc comment may overstate behavior.
cgroup_limits().total_memoryreturns the cgroup limit directly, notmin(cgroup, physical). If cgroup is (unusually) set higher than physical RAM, this returns the cgroup value. The fallback only triggers when no cgroup limit exists.Consider: "Prefers cgroup memory limit when available, falling back to physical RAM."
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/system.rs` around lines 29 - 30, The doc comment is inaccurate about behavior: update the comment for the function that documents cgroup_limits().total_memory (the cgroup memory check logic) to state that the code prefers the cgroup memory limit when available and falls back to physical RAM if no cgroup limit exists, instead of claiming it always returns min(cgroup, physical); ensure the new wording explicitly notes that cgroup_limits().total_memory returns the cgroup limit directly when present.src/commands/common.rs (1)
1110-1130: Tests duplicate coverage fromsrc/lib/system.rs.These tests mirror the ones in the source module (lines 52-65). Since the re-export is trivial, testing at the source is sufficient. Consider removing to reduce redundancy.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/commands/common.rs` around lines 1110 - 1130, The two unit tests test_detect_total_memory_returns_nonzero and test_detect_total_memory_bounded_by_sysinfo in this file duplicate tests already present for detect_total_memory in the system module; remove these duplicate test functions from src/commands/common.rs (the test cases that call detect_total_memory and assert non-zero and bounded-by-sysinfo) so coverage remains only in the original src/lib/system.rs tests and avoid redundant assertions.src/lib/sort/pooled_chunk_writer.rs (1)
96-106: Reuse a scratch buffer for keyed spills.
Vec::new()here allocates once per non-embedded record on the hot spill path. The current keyed spill formats are fixed-size, so a reusable scratch buffer avoids millions of tiny allocations during large spills.Possible direction
pub struct PooledChunkWriter<K: RawSortKey> { staging: StagingBuffer, io_handle: Option<JoinHandle<Result<()>>>, + key_buf: Vec<u8>, _phantom: PhantomData<K>, } ... Ok(Self { staging: StagingBuffer::new(pool, result_tx), io_handle: Some(io_handle), + key_buf: Vec::with_capacity(K::SERIALIZED_SIZE.unwrap_or(64)), _phantom: PhantomData, }) ... - let mut key_buf = Vec::new(); - key.write_to(&mut key_buf)?; - let needed = key_buf.len() + 4 + record.len(); + self.key_buf.clear(); + key.write_to(&mut self.key_buf)?; + let needed = self.key_buf.len() + 4 + record.len(); if self.staging.buf().len() + needed > BGZF_MAX_BLOCK_SIZE { self.staging.flush(); } - self.staging.buf().extend_from_slice(&key_buf); + self.staging.buf().extend_from_slice(&self.key_buf);🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/pooled_chunk_writer.rs` around lines 96 - 106, Replace the per-record Vec::new() scratch allocation with a reusable buffer on the writer (e.g., a field like key_scratch: Vec<u8>) and use key_scratch.clear() before reusing it; call key.write_to(&mut self.key_scratch) instead of allocating key_buf, compute needed from self.key_scratch.len(), call self.staging.flush() when needed, then extend self.staging.buf() from self.key_scratch and write_chunked(record). Also ensure the scratch buffer is reserved/grown as needed to avoid repeated reallocations.
🤖 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/sort/raw.rs`:
- Around line 3692-3714: The test test_merge_bams_queryname_natural_sort is
ineffective because create_sorted_bam currently emits zero-padded names so
lexicographic and natural order match; update the test to generate
non-zero-padded numeric names (or pass a no-padding flag to create_sorted_bam)
so names like "read2" vs "read10" appear and natural ordering differs from
lexicographic, then validate order by either comparing adjacent names using the
same natural comparator (SortOrder::Queryname(QuerynameComparator::Natural) via
RawExternalSorter::new) or by parsing integers from names and asserting
non-decreasing numeric sequence (helpers: create_sorted_bam, collect_read_names,
RawExternalSorter::new, SortOrder::Queryname(QuerynameComparator::Natural)).
- Around line 968-987: The EOF branch currently runs before checking
decompression_error, which lets a drained queue hide a decompression failure and
return Ok(false); move the decompression error check to occur before the EOF
logic: in the function containing self.all_chunks_eof, self.raw_chunk_blocks,
self.decompressed_chunks and self.decompression_error (the loop that sets
source.eof on per_source), check
self.decompression_error.load(Ordering::Acquire) first and return the Err(...)
immediately if set, then perform the EOF check and mark per_source.eof as
before.
In `@src/lib/sort/read_ahead.rs`:
- Around line 538-541: The Direct enum variant currently stores an
Option<std::io::Error> but the reader loop in RawBamReader polling code
overwrites that slot on every failed next(); change the logic in the
iterator/polling code that calls next() on the RawBamReader<PooledInputStream>
so it only sets the Direct's error slot if it is currently None and then stops
further polling/iteration (return Pending or Ready(None) as appropriate) after
the first failure; ensure take_error() still returns that preserved first error
and do not overwrite the Option on subsequent errors.
- Around line 430-469: The drain logic currently empties the bounded
decompressed_input into an unbounded reorder (reorder, drain_queue), breaking
backpressure; change drain_queue (and its callers like next_block) so it only
consumes contiguous serials and stops when a gap would be created: accept the
current next_serial (and/or compute expected = next_serial + reorder.len()) and
only pop/insert while popped.serial == expected (increment expected each
insert); additionally enforce a small cap on reorder size (max_reorder) and stop
draining when that cap is reached; update next_block to pass next_serial and
respect the cap so producers remain bounded.
- Around line 442-469: next_block() can park forever when a producer sets
decompression_error or input_read_error but does not set
decompressed_input_done; add checks for those error flags inside the wait loop
so the consumer unblocks and returns to read() to handle errors. Specifically,
inside the loop in next_block() (after calling self.drain_queue() and before
std::thread::park()), check self.decompression_error and self.input_read_error
(and any other terminal producer-error flags) and if any are set return None so
the caller can observe the error state; keep existing reorder/remove,
next_serial increment, and EOF double-drain behavior unchanged. Ensure you
reference the struct fields decompressed_input_done, decompression_error,
input_read_error and methods drain_queue(), is_eof(), and the reorder map when
making the change.
---
Duplicate comments:
In `@src/lib/bam_io.rs`:
- Around line 777-813: The code currently treats only is_stdin_path(path_ref) as
non-replayable; instead detect seekability on the opened handle and use the
TeeReader+ChainedReader path for any non-seekable input. After opening
File::open(path_ref) use file.seek(SeekFrom::Current(0)) (or attempt a small
seek) to determine if seeking is supported; if seek returns Err (non-seekable),
wrap the opened file in TeeReader -> BgzfReader -> noodles::bam::io::Reader to
read the header and then into_parts() -> ChainedReader (same pattern as the
stdin branch); otherwise keep the existing seek-to-0 and buffered file path.
Update the header/reader construction code around the file/open, header read,
and rewind (file.seek(...)) logic (references: File::open, file.seek,
BgzfReader, noodles::bam::io::Reader::from, TeeReader::new, ChainedReader::new,
into_parts).
In `@src/lib/sort/bgzf_io.rs`:
- Around line 112-125: The writer-side reordering is unbounded because we always
drain result_rx into reorder_buf, so if an early serial stalls the BTreeMap
grows without backpressure; introduce a MAX_PENDING cap and refuse to insert
when reorder_buf.len() >= MAX_PENDING by pausing intake until the writer
advances (i.e., block/wait on processing the next_expected serials or otherwise
process incoming results until next_expected appears) instead of blindly
inserting; implement this check around the result handling logic that uses
result_rx, reorder_buf, result.serial, next_expected and writer.write_all, and
use BufferPool::checkout semantics to ensure we stop consuming more results when
the map is full to restore backpressure.
---
Nitpick comments:
In `@src/commands/common.rs`:
- Around line 1110-1130: The two unit tests
test_detect_total_memory_returns_nonzero and
test_detect_total_memory_bounded_by_sysinfo in this file duplicate tests already
present for detect_total_memory in the system module; remove these duplicate
test functions from src/commands/common.rs (the test cases that call
detect_total_memory and assert non-zero and bounded-by-sysinfo) so coverage
remains only in the original src/lib/system.rs tests and avoid redundant
assertions.
In `@src/lib/sort/pooled_chunk_writer.rs`:
- Around line 96-106: Replace the per-record Vec::new() scratch allocation with
a reusable buffer on the writer (e.g., a field like key_scratch: Vec<u8>) and
use key_scratch.clear() before reusing it; call key.write_to(&mut
self.key_scratch) instead of allocating key_buf, compute needed from
self.key_scratch.len(), call self.staging.flush() when needed, then extend
self.staging.buf() from self.key_scratch and write_chunked(record). Also ensure
the scratch buffer is reserved/grown as needed to avoid repeated reallocations.
In `@src/lib/system.rs`:
- Around line 29-30: The doc comment is inaccurate about behavior: update the
comment for the function that documents cgroup_limits().total_memory (the cgroup
memory check logic) to state that the code prefers the cgroup memory limit when
available and falls back to physical RAM if no cgroup limit exists, instead of
claiming it always returns min(cgroup, physical); ensure the new wording
explicitly notes that cgroup_limits().total_memory returns the cgroup limit
directly when present.
🪄 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: d3fc56d9-df91-4c0f-8e38-c6fab379da42
⛔ Files ignored due to path filters (1)
Cargo.lockis excluded by!**/*.lock
📒 Files selected for processing (18)
.gitignoreCargo.tomlcrates/fgumi-bgzf/src/lib.rscrates/fgumi-bgzf/src/writer.rssrc/commands/common.rssrc/commands/sort.rssrc/lib/bam_io.rssrc/lib/mod.rssrc/lib/sort/bgzf_io.rssrc/lib/sort/mod.rssrc/lib/sort/pooled_bam_writer.rssrc/lib/sort/pooled_chunk_writer.rssrc/lib/sort/raw.rssrc/lib/sort/read_ahead.rssrc/lib/sort/worker_pool.rssrc/lib/system.rstests/integration/main.rstests/integration/test_sort_correctness.rs
✅ Files skipped from review due to trivial changes (7)
- .gitignore
- tests/integration/main.rs
- src/lib/mod.rs
- crates/fgumi-bgzf/src/writer.rs
- src/lib/sort/mod.rs
- src/lib/sort/pooled_bam_writer.rs
- src/lib/sort/worker_pool.rs
🚧 Files skipped from review as they are similar to previous changes (4)
- Cargo.toml
- crates/fgumi-bgzf/src/lib.rs
- src/commands/sort.rs
- tests/integration/test_sort_correctness.rs
288a49e to
c0830ef
Compare
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
c0830ef to
26821dc
Compare
There was a problem hiding this comment.
Actionable comments posted: 11
♻️ Duplicate comments (4)
src/lib/sort/read_ahead.rs (2)
451-476:⚠️ Potential issue | 🟠 MajorStop parking once the producer has already failed.
After
input_read_errorordecompression_erroris set, this loop still waits only fordecompressed_input_done. The worker wakes the reader once, but the next iteration can park forever. ReturnNoneon either error flag soread()can surface the failure.Patch
// Check if the block we need is ready (O(1) via ReorderBuffer) if let Some(data) = self.reorder.try_pop_next() { return Some(data); } + if self.input_read_error.load(std::sync::atomic::Ordering::Acquire) + || self.decompression_error.load(std::sync::atomic::Ordering::Acquire) + { + return None; + } + // Nothing available — check EOF if self.is_eof() {🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/read_ahead.rs` around lines 451 - 476, In next_block(), after draining the queue and before parking (and also when checking EOF), add checks for the producer failure flags input_read_error and decompression_error and return None immediately if either is set so the reader does not park forever; specifically, update next_block() (which already calls drain_queue(), reorder.try_pop_next(), is_eof(), and uses std::thread::park()) to short-circuit and return None when input_read_error || decompression_error is true, ensuring read() can observe the failure instead of deadlocking.
573-585:⚠️ Potential issue | 🟡 MinorKeep the first direct-read error.
The
Directdocs say this slot holds the first read error, but this branch still overwrites it on every failure and keeps polling after that. Preserve the original error and short-circuit laternext()calls.Patch
match self { Self::ReadAhead(r) => r.next(), + Self::Direct(_, Some(_)) => None, Self::Direct(reader, error_slot) => { let mut record = RawRecord::default(); match reader.read_record(&mut record) { Ok(0) => None, // EOF Ok(_) => Some(record), Err(e) => { log::error!("Error reading raw BAM record: {e}"); - *error_slot = Some(e); + if error_slot.is_none() { + *error_slot = Some(e); + } None } } } }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/read_ahead.rs` around lines 573 - 585, In next(), when matching the Direct(reader, error_slot) branch, short-circuit immediately if *error_slot is already Some to stop polling after the first failure; when reader.read_record returns Err(e), only store the error into error_slot if it is currently None (preserve the first error) and then return None; reference the next method, the Direct enum variant, the error_slot, and reader.read_record/RawRecord to locate and apply this change.src/lib/bam_io.rs (1)
777-813:⚠️ Potential issue | 🟠 MajorThe “stdin/FIFOs” path still only handles stdin.
This condition is just
is_stdin_path(path_ref). Named pipes and process-substitution paths still go through theseek(0)branch, so header parsing works and the rewind fails afterward. Gate the fast path on actual seekability, or reuse the tee/chained replay path for any non-regular input.src/lib/sort/raw.rs (1)
3719-3740:⚠️ Potential issue | 🟡 MinorThis “natural” merge test is still lexicographic.
create_sorted_bam()still generates zero-padded names, and the assertion uses plain string<=, so natural and lexicographic order are identical here. A regression back to lexicographic merge would still pass.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/raw.rs` around lines 3719 - 3740, The test currently uses zero-padded names so natural vs lexicographic order is indistinguishable; update the test to generate non-zero-padded, numerically-mixed read names so Natural ordering is exercised. Specifically, change the calls to create_sorted_bam in test_merge_bams_queryname_natural_sort (or add a new helper) so the generated read names are like "r1","r2","r10" (no leading zeros) instead of zero-padded names, then run the same RawExternalSorter::new(SortOrder::Queryname(QuerynameComparator::Natural)).merge_bams(...) and keep the collect_read_names and window assertions unchanged; if create_sorted_bam cannot be parameterized, add a new helper (e.g., create_sorted_bam_unpadded) that emits unpadded numeric names and use it in this test. Ensure the produced names will fail lexicographic but pass natural ordering.
🧹 Nitpick comments (1)
src/lib/unified_pipeline/base.rs (1)
4201-4211: Prefer not to make this hook a silent no-op.The default
truemeans missing impls disable write-side admission control without any signal.OutputPipelineQueuesin this file already falls through to that default, so itswrite_reorder_stateis ignored unless a wrapper overrides the trait. Making this required, or delegating from the base impl too, would keep the invariant enforced.🤖 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 4201 - 4211, The default implementation of write_reorder_can_proceed silently returns true and thus disables write-side admission control; change the base impl to delegate to the pipeline's write_reorder_state instead of being a no-op: have write_reorder_can_proceed call into the existing write_reorder_state (and its admission-control API) so OutputPipelineQueues' write_reorder_state is honored, or alternatively make write_reorder_can_proceed a required trait method (remove the default) so implementors must explicitly enforce admission control; update the implementation points that relied on the old default accordingly.
🤖 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/commands/sort.rs`:
- Around line 853-854: The assert!(resolved <= total) can fail in auto mode when
the per-thread floor (256 MiB) makes budget > total; update resolve_memory_limit
to avoid panicking by clamping resolved to total (or replacing the hard assert
with a debug_assert and then setting resolved = total.min(resolved)) and ensure
any warning about budget > total remains; refer to resolve_memory_limit and the
variables resolved, total, and budget when locating the change.
In `@src/lib/bam_io.rs`:
- Around line 846-860: The skip_bam_header() implementation currently uses
io::copy(&mut reader.take(n), &mut io::sink()) and ignores the returned byte
count, allowing truncated headers to pass; update skip_bam_header() to check the
usize/count returned by io::copy (or use reader.read_exact into a buffer) for
each discard (the header text length l_text and each reference name length
l_name) and return an error if fewer bytes were consumed than expected; also
validate the reads for n_ref and l_ref as currently done via reader.read_exact,
and apply the same consumed-bytes check to every place using reader.take so
truncated input yields an Err rather than silently succeeding.
In `@src/lib/sort/bgzf_io.rs`:
- Around line 143-178: The writer paths can return early and leak permits held
by producers (acquired in StagingBuffer::flush), so modify the block that
performs writer.write_all (both inside the main recv loop, the cascade that
drains reorder_buf, and the final drain loop) to ensure the permit pool is
signalled on any error exit: on any Err from writer.write_all or writer.flush,
first call the permit-pool shutdown/abort/close API (e.g., permit_pool.abort()
or permit_pool.close() — whatever wakes/wakes_all waiters) and/or release any
outstanding permits you own, then propagate the error; also ensure the final
BGZF_EOF write path does the same so producers do not deadlock. Locate uses of
result_rx, buffer_pool.checkin, permit_pool.release, reorder_buf.remove/insert,
writer.write_all, writer.flush and add the abort/cleanup call(s) in all error
branches before returning Err.
In `@src/lib/sort/pooled_bam_writer.rs`:
- Line 1: Add the crate-level attribute to forbid unsafe code in this new
module: insert #![deny(unsafe_code)] at the top of
src/lib/sort/pooled_bam_writer.rs (above the module doc comment) so the file
opts into the repository's unsafe-code guard; then run a build to ensure no
unsafe blocks remain in functions/types such as any SortWorkerPool or
PooledBamWriter implementations and fix or remove any explicit unsafe usage if
reported.
In `@src/lib/sort/pooled_chunk_writer.rs`:
- Around line 106-117: The check using anyhow::ensure!(needed <=
BGZF_MAX_BLOCK_SIZE, ...) incorrectly prevents large keyed records from being
written because StagingBuffer::write_chunked() already supports payloads that
span BGZF blocks; remove or narrow that ensure! so it does not short‑circuit the
non-embedded path. Specifically, in the method that computes needed from
self.key_buf and record, delete (or conditionally skip) the ensure! that
compares needed to BGZF_MAX_BLOCK_SIZE and leave the existing logic that flushes
the staging buffer when self.staging.buf().len() + needed > BGZF_MAX_BLOCK_SIZE,
then continue to extend self.staging.buf() with self.key_buf and the length
bytes and call self.staging.write_chunked(record) so large records are handled
by write_chunked.
- Around line 40-45: PooledChunkWriter currently lacks a Drop implementation so
if write_record() returns Err and the writer is dropped the io_handle JoinHandle
is detached and its result lost; implement a Drop for PooledChunkWriter<K> that
calls start_finish()/finish() (or otherwise flushes/stops the background task)
and joins the io_handle, or alternatively change the API to force explicit
finalization (e.g. a builder/Sealed trait requiring finish()). Locate symbols
PooledChunkWriter, write_record, start_finish, finish and the io_handle
JoinHandle in pooled_chunk_writer.rs and ensure Drop either invokes
start_finish()+finish() and then awaits/join() the io_handle Result or
logs/propagates the task error deterministically so the background thread is not
left detached on error paths.
In `@src/lib/sort/raw.rs`:
- Around line 1794-1815: The pending_spill local (holding PendingSpill with its
SpillWriteHandle) must not be kept live across early-return paths; ensure the
handle is dropped before any fallible checks like record_source.take_error() or
before returning on error. Modify the code around the pending_spill assignment
so you either (a) move the record_source.take_error() / early-return logic above
creating/assigning pending_spill, or (b) explicitly drop or set pending_spill =
None immediately before calling record_source.take_error() and any other `?`
paths, and guarantee drain_pending_spill(&mut pending_spill, ...) is only called
when no earlier returns remain. Refer to the pending_spill variable,
PendingSpill type, record_source.take_error(), and the
drain_pending_spill::<RawCoordinateKey>() call when making the change.
In `@src/lib/system.rs`:
- Around line 53-56: The test test_detect_total_memory_nonzero uses a hard lower
bound (>= 64 MiB) which can fail in constrained CI; change the assertion to only
require detect_total_memory() returns non-zero (e.g., assert!(total > 0) or
assert_ne!(total, 0)) and optionally rename the test to reflect nonzero
behavior; update the assertion in the test_detect_total_memory_nonzero function
that calls detect_total_memory() accordingly.
- Around line 35-36: The code currently uses
system.cgroup_limits().map_or_else(|| system.total_memory(), |c| c.total_memory)
which can return a cgroup-reported total larger than the host physical memory;
change it to clamp the chosen cgroup value to system.total_memory() before
converting to usize (i.e. take the minimum of the cgroup total and
system.total_memory()), then perform usize::try_from on that clamped value and
keep the existing unwrap_or(usize::MAX) fallback; update the logic around
system.cgroup_limits() and usize::try_from to use this min-clamped value.
In `@src/lib/unified_pipeline/bam.rs`:
- Around line 1142-1144: The write-reorder gating is wired via
write_reorder_state and write_reorder_can_proceed but try_step_write still uses
insert/try_pop_next which leaves the reorder heap size and next_seq stale;
change try_step_write to use the size-aware APIs: replace insert with
insert_with_size, replace try_pop_next with try_pop_next_with_size, and update
heap accounting by calling add_heap_bytes when inserting, sub_heap_bytes when
removing, and update_next_seq when advancing the next sequence; ensure
write_reorder_state methods (insert_with_size, try_pop_next_with_size,
add_heap_bytes, sub_heap_bytes, update_next_seq) are invoked in the same places
where insert/try_pop_next were used so can_proceed() reflects actual
heap/next_seq state and write-side backpressure works correctly.
In `@src/lib/unified_pipeline/base.rs`:
- Around line 1763-1769: The write-side admission gate is being initialized with
ReorderBufferState::new(0), which hardcodes the 512 MiB default and ignores
PipelineConfig.queue_memory_limit; change OutputPipelineQueues::new call sites
to accept and forward the configured limit and construct the reorder gate with
ReorderBufferState::new(config.queue_memory_limit) (update the constructor
signature of OutputPipelineQueues::new if necessary), and ensure all callers
pass PipelineConfig.queue_memory_limit so the admission control honors the
configured queue_memory_limit.
---
Duplicate comments:
In `@src/lib/sort/raw.rs`:
- Around line 3719-3740: The test currently uses zero-padded names so natural vs
lexicographic order is indistinguishable; update the test to generate
non-zero-padded, numerically-mixed read names so Natural ordering is exercised.
Specifically, change the calls to create_sorted_bam in
test_merge_bams_queryname_natural_sort (or add a new helper) so the generated
read names are like "r1","r2","r10" (no leading zeros) instead of zero-padded
names, then run the same
RawExternalSorter::new(SortOrder::Queryname(QuerynameComparator::Natural)).merge_bams(...)
and keep the collect_read_names and window assertions unchanged; if
create_sorted_bam cannot be parameterized, add a new helper (e.g.,
create_sorted_bam_unpadded) that emits unpadded numeric names and use it in this
test. Ensure the produced names will fail lexicographic but pass natural
ordering.
In `@src/lib/sort/read_ahead.rs`:
- Around line 451-476: In next_block(), after draining the queue and before
parking (and also when checking EOF), add checks for the producer failure flags
input_read_error and decompression_error and return None immediately if either
is set so the reader does not park forever; specifically, update next_block()
(which already calls drain_queue(), reorder.try_pop_next(), is_eof(), and uses
std::thread::park()) to short-circuit and return None when input_read_error ||
decompression_error is true, ensuring read() can observe the failure instead of
deadlocking.
- Around line 573-585: In next(), when matching the Direct(reader, error_slot)
branch, short-circuit immediately if *error_slot is already Some to stop polling
after the first failure; when reader.read_record returns Err(e), only store the
error into error_slot if it is currently None (preserve the first error) and
then return None; reference the next method, the Direct enum variant, the
error_slot, and reader.read_record/RawRecord to locate and apply this change.
---
Nitpick comments:
In `@src/lib/unified_pipeline/base.rs`:
- Around line 4201-4211: The default implementation of write_reorder_can_proceed
silently returns true and thus disables write-side admission control; change the
base impl to delegate to the pipeline's write_reorder_state instead of being a
no-op: have write_reorder_can_proceed call into the existing write_reorder_state
(and its admission-control API) so OutputPipelineQueues' write_reorder_state is
honored, or alternatively make write_reorder_can_proceed a required trait method
(remove the default) so implementors must explicitly enforce admission control;
update the implementation points that relied on the old default accordingly.
🪄 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: 53654e05-3e0d-456b-93fb-8ba93213774d
⛔ Files ignored due to path filters (1)
Cargo.lockis excluded by!**/*.lock
📒 Files selected for processing (23)
.gitignoreCargo.tomlcrates/fgumi-bgzf/src/lib.rscrates/fgumi-bgzf/src/writer.rssrc/commands/common.rssrc/commands/sort.rssrc/lib/bam_io.rssrc/lib/mod.rssrc/lib/sort/bgzf_io.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/mod.rssrc/lib/sort/pooled_bam_writer.rssrc/lib/sort/pooled_chunk_writer.rssrc/lib/sort/raw.rssrc/lib/sort/read_ahead.rssrc/lib/sort/segmented_buf.rssrc/lib/sort/worker_pool.rssrc/lib/system.rssrc/lib/unified_pipeline/bam.rssrc/lib/unified_pipeline/base.rssrc/lib/unified_pipeline/fastq.rstests/integration/main.rstests/integration/test_sort_correctness.rs
✅ Files skipped from review due to trivial changes (3)
- .gitignore
- tests/integration/main.rs
- src/lib/sort/mod.rs
🚧 Files skipped from review as they are similar to previous changes (4)
- crates/fgumi-bgzf/src/writer.rs
- crates/fgumi-bgzf/src/lib.rs
- src/commands/common.rs
- tests/integration/test_sort_correctness.rs
There was a problem hiding this comment.
Actionable comments posted: 1
♻️ Duplicate comments (1)
src/lib/bam_io.rs (1)
911-925:⚠️ Potential issue | 🟠 Major
skip_bam_header()still silently accepts truncated headers.
io::copyreturns bytes copied, not an error on short read. Truncated input succeeds silently.Patch
- io::copy(&mut reader.take(l_text as u64), &mut io::sink())?; + let copied = io::copy(&mut reader.take(l_text as u64), &mut io::sink())?; + anyhow::ensure!(copied == l_text as u64, "Truncated BAM header text"); @@ - io::copy(&mut reader.take(l_name as u64), &mut io::sink())?; + let copied = io::copy(&mut reader.take(l_name as u64), &mut io::sink())?; + anyhow::ensure!(copied == l_name as u64, "Truncated BAM reference entry");🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/bam_io.rs` around lines 911 - 925, skip_bam_header currently uses io::copy(&mut reader.take(len), &mut io::sink()) which silently succeeds on truncated input because io::copy returns bytes copied rather than erroring; change the calls that skip variable-length fields (the l_text and l_name skips) to verify the number of bytes consumed matches the expected length: capture the u64 returned by io::copy for reader.take(l_text as u64) and reader.take(l_name as u64) and if it is less than the requested length return an Err(io::Error::new(io::ErrorKind::UnexpectedEof, "truncated BAM header")) (or replace with reader.read_exact into a temporary buffer in skip_bam_header to guarantee a short read errors), referencing the skip_bam_header function, the l_text and l_name variables, and the io::copy(&mut reader.take(...), &mut io::sink()) calls.
🧹 Nitpick comments (2)
src/lib/sort/raw.rs (2)
2737-2750: Indexing path intentionally uses legacy readers — document follow-up.The TODO comments note this path needs
PooledIndexingBamWriter. Consider opening an issue to track this follow-on work so it doesn't get lost.Want me to open a GitHub issue to track pool-integrated indexing?
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/raw.rs` around lines 2737 - 2750, The indexing path currently uses legacy readers (GenericKeyedChunkReader) and the TODO mentions replacing create_indexing_bam_writer with PooledIndexingBamWriter; please open a GitHub issue to track implementing pool-integrated decompression for the indexing path, then update the TODO in src/lib/sort/raw.rs near the build_chunk_sources call (referencing create_indexing_bam_writer, PooledIndexingBamWriter, and GenericKeyedChunkReader) to include the new issue number and a short actionable description so the follow-up work isn't lost.
1178-1184: Recursiveread()call is safe but adds cognitive load.The recursion at line 1183 is bounded (poll returns
Ok(true)only when target has data), but an iterative approach would be clearer. Consider aloopwithcontinueinstead.Iterative alternative
- } else { - // poll said data is ready but the source is still empty — retry rather - // than returning Ok(0) which read_key_from_source treats as clean EOF. - self.read(buf) - } + } else { + // poll said data is ready but the source is still empty — retry + continue; // would require wrapping in loop { ... } + }Or wrap the entire function body in
loop { ... }and replace the recursion withcontinue.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/raw.rs` around lines 1178 - 1184, The recursive call self.read(buf) in the read implementation is safe but should be made iterative: refactor the function (the read method that updates self.bytes_read and calls poll to check readiness) to wrap its body in a loop and replace the recursive call with continue; keep the same behavior—when poll reports ready but take == 0, continue the loop, when take > 0 increment self.bytes_read and return Ok(take), and otherwise return Err/Ok as before—so the logic around self.bytes_read, the readiness check and the interaction with read_key_from_source remains unchanged but uses loop/continue instead of recursion.
🤖 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/sort/raw.rs`:
- Around line 156-161: The log_summary method computes percentages by dividing
by overall which can be 0.0; update fn log_summary(&self, threads: usize) to
guard against division by zero: compute overall as currently done, then if
overall > 0.0 calculate read_pct, sort_pct, spill_pct normally (100.0 * ... /
overall) else set those percentage variables to 0.0 (or a safe default) before
using them in the log output; reference overall_start, read_secs, sort_secs,
spill_write_secs, and spill_count to locate the calculations to change.
---
Duplicate comments:
In `@src/lib/bam_io.rs`:
- Around line 911-925: skip_bam_header currently uses io::copy(&mut
reader.take(len), &mut io::sink()) which silently succeeds on truncated input
because io::copy returns bytes copied rather than erroring; change the calls
that skip variable-length fields (the l_text and l_name skips) to verify the
number of bytes consumed matches the expected length: capture the u64 returned
by io::copy for reader.take(l_text as u64) and reader.take(l_name as u64) and if
it is less than the requested length return an
Err(io::Error::new(io::ErrorKind::UnexpectedEof, "truncated BAM header")) (or
replace with reader.read_exact into a temporary buffer in skip_bam_header to
guarantee a short read errors), referencing the skip_bam_header function, the
l_text and l_name variables, and the io::copy(&mut reader.take(...), &mut
io::sink()) calls.
---
Nitpick comments:
In `@src/lib/sort/raw.rs`:
- Around line 2737-2750: The indexing path currently uses legacy readers
(GenericKeyedChunkReader) and the TODO mentions replacing
create_indexing_bam_writer with PooledIndexingBamWriter; please open a GitHub
issue to track implementing pool-integrated decompression for the indexing path,
then update the TODO in src/lib/sort/raw.rs near the build_chunk_sources call
(referencing create_indexing_bam_writer, PooledIndexingBamWriter, and
GenericKeyedChunkReader) to include the new issue number and a short actionable
description so the follow-up work isn't lost.
- Around line 1178-1184: The recursive call self.read(buf) in the read
implementation is safe but should be made iterative: refactor the function (the
read method that updates self.bytes_read and calls poll to check readiness) to
wrap its body in a loop and replace the recursive call with continue; keep the
same behavior—when poll reports ready but take == 0, continue the loop, when
take > 0 increment self.bytes_read and return Ok(take), and otherwise return
Err/Ok as before—so the logic around self.bytes_read, the readiness check and
the interaction with read_key_from_source remains unchanged but uses
loop/continue instead of recursion.
🪄 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: 146eddcc-5e3b-4dba-a765-261eb2345753
📒 Files selected for processing (4)
src/lib/bam_io.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/raw.rssrc/lib/sort/worker_pool.rs
26821dc to
a7f59e0
Compare
There was a problem hiding this comment.
Actionable comments posted: 5
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/fastq.rs (1)
2533-2550:⚠️ Potential issue | 🟡 MinorHeap bytes not decremented on write error.
If
write_allfails at line 2537, we return early without callingsub_heap_bytes. The batch is already popped fromreorder, so accounting becomes inconsistent.Since errors typically mean shutdown, this is low-risk but leaves stale bytes in
write_reorder_state.🛡️ Proposed fix
while let Some((batch, heap_size)) = reorder.try_pop_next_with_size() { let mut batch_bytes: u64 = 0; for block in &batch.blocks { batch_bytes += block.data.len() as u64; if let Err(e) = output.write_all(&block.data) { + state.output.write_reorder_state.sub_heap_bytes(heap_size as u64); state.set_error(e); return false; } } state.output.write_reorder_state.sub_heap_bytes(heap_size as u64);🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/unified_pipeline/fastq.rs` around lines 2533 - 2550, When a write_all call inside the loop handling reorder.try_pop_next_with_size() fails, the code returns early after state.set_error(e) without calling state.output.write_reorder_state.sub_heap_bytes(heap_size as u64), leaving the heap accounting stale; modify the error path so that before returning false you always call state.output.write_reorder_state.sub_heap_bytes(heap_size as u64) (and optionally update_next_seq(reorder.next_seq()) if semantics require) — locate the block in the loop around reorder.try_pop_next_with_size(), output.write_all(&block.data), state.set_error(e) and ensure sub_heap_bytes is invoked even on write errors (use a small scope guard or explicit call before return).
♻️ Duplicate comments (1)
src/lib/unified_pipeline/base.rs (1)
5139-5142:⚠️ Potential issue | 🟠 MajorAccount write-side bytes before the writer drains Q6.
write_reorder_can_proceed()is checked on the compress side before enqueue, but the counter only grows here aftercompressed.pop(). That leaves the gate blind to bytes already buffered in Q6 or sitting in held compressed batches, so a lowqueue_memory_limitcan still be overshot by the whole write-side backlog. Start accounting when output is admitted to the write side, or include Q6 bytes in this gate.Also applies to: 5167-5169
🤖 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 5139 - 5142, The write-side byte counter is updated only after compressed.pop(), so write_reorder_can_proceed() can admit more output than queue_memory_limit allows; move or duplicate the heap accounting to occur at the moment the output is admitted to the write side (i.e., before/when calling reorder.insert_with_size or immediately prior to enqueue) so state.write_reorder_state().add_heap_bytes(heap_size as u64) happens before or atomically with the insert; alternatively ensure write_reorder_can_proceed() considers Q6/held-compressed bytes by including those bytes in its check. Apply the same change to the other occurrence around the reorder.insert_with_size/state.write_reorder_state().add_heap_bytes pair (the 5167-5169 region).
🧹 Nitpick comments (3)
src/lib/sort/pooled_chunk_writer.rs (1)
288-291: Consider unconditional shutdown for test cleanup.
Arc::try_unwrapsilently skipsshutdown()if other refs exist. IfSortWorkerPool::drophandles thread cleanup, this is fine; otherwise tests could leak threads on failure paths.♻️ Alternative pattern
- if let Ok(pool) = Arc::try_unwrap(pool) { - pool.shutdown(); - } + Arc::try_unwrap(pool) + .expect("pool should have single owner at test end") + .shutdown();🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/pooled_chunk_writer.rs` around lines 288 - 291, The current use of Arc::try_unwrap(pool) means pool.shutdown() is skipped when other Arc refs exist, risking leaked threads in tests; change the code to call shutdown unconditionally (i.e., always invoke pool.shutdown() even if Arc::try_unwrap fails) by making SortWorkerPool::shutdown accept &self (or add a thread-safe public shutdown(&self) method) so you can call Arc::as_ref().shutdown() regardless of ownership; alternatively ensure SortWorkerPool::drop performs the same cleanup so no shutdown call is skipped.src/lib/unified_pipeline/base.rs (1)
5893-5943: Add a behavior test for non-zeroqueue_memory_limit.These test updates only cover the new constructor signature. A tiny limit with out-of-order serials would lock in the intended write-side backpressure behavior.
🤖 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 5893 - 5943, Add a unit test that verifies write-side backpressure when constructing OutputPipelineQueues with a non-zero queue_memory_limit: create an OutputPipelineQueues via OutputPipelineQueues::new with a very small queue_memory_limit, push TestProcessed/TestSerialized items with out-of-order serials (use the existing push methods used across tests), and assert that queue_depths() and are_queues_empty()/has_error()/is_draining() behave to show the pipeline does not grow beyond the memory limit (i.e., depths respect the small limit and producers are constrained); name it e.g. test_output_queues_backpressure_limit and reference OutputPipelineQueues::new, queue_depths, and the push methods to locate where to add the test.src/lib/unified_pipeline/bam.rs (1)
2836-2841: Avoid duplicate heap-size estimation in the hot write drain loop.
estimate_heap_size()is computed twice per batch. Reuse one value.Suggested diff
- while let Some((serial, batch)) = state.output.compressed.pop() { - let q7_heap = batch.estimate_heap_size() as u64; - state.q6_track_pop(q7_heap); - state.deadlock_state.record_q7_pop(); - let heap_size = batch.estimate_heap_size(); - reorder.insert_with_size(serial, batch, heap_size); - state.output.write_reorder_state.add_heap_bytes(heap_size as u64); - } + while let Some((serial, batch)) = state.output.compressed.pop() { + let heap_size = batch.estimate_heap_size(); + state.q6_track_pop(heap_size as u64); + state.deadlock_state.record_q7_pop(); + reorder.insert_with_size(serial, batch, heap_size); + state.output.write_reorder_state.add_heap_bytes(heap_size as u64); + }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/unified_pipeline/bam.rs` around lines 2836 - 2841, In the hot drain loop starting at state.output.compressed.pop(), avoid calling batch.estimate_heap_size() twice: call it once into a local variable (e.g., let heap = batch.estimate_heap_size()), reuse that value for q7_heap (convert to u64 for state.q6_track_pop if needed) and pass the original heap to reorder.insert_with_size(serial, batch, heap); keep the calls to state.q6_track_pop(...) and state.deadlock_state.record_q7_pop() unchanged but use the single cached heap value for both places that currently recompute it.
🤖 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/bam_io.rs`:
- Around line 815-817: The call sequence in
create_raw_bam_reader_pool_integrated unconditionally calls
pool.set_input_file(reader) and pool.set_phase(phase::PHASE1), which can corrupt
state if the pool is already in PHASE1 and workers are active; add a guard or
assertion to check the current phase before changing it (for example query
pool.phase() or similar) and either return/error if already in PHASE1 or
document the precondition in the function docs; reference the
create_raw_bam_reader_pool_integrated function and the pool.set_input_file /
pool.set_phase (phase::PHASE1) calls when implementing the check or adding the
precondition note.
In `@src/lib/commands/common.rs`:
- Around line 570-578: The total-memory calculation here is inconsistent with
detect_total_memory(): instead of using system.cgroup_limits().map(|c|
c.total_memory).unwrap_or_else(|| system.total_memory()), call or reuse
detect_total_memory() (the function in system.rs) or apply the same clamp logic
by taking c.total_memory.min(physical) so cgroup limits are capped at physical
RAM; update this calculation so it matches resolve_memory_limit's semantics and
references detect_total_memory()/the local physical value to compute
total_memory_bytes.
In `@src/lib/sort/pooled_bam_writer.rs`:
- Around line 69-80: The spawned I/O thread (io_handle from thread::spawn
calling io_writer_loop) is being dropped if staging.write_chunked or
staging.flush returns Err in new() and start_finish(), which detaches the
thread; fix by ensuring the JoinHandle is joined before returning an error:
either immediately wrap the handle in its final owner only after all fallible
operations succeed, or on an early error drop staging/resources and call
io_handle.join() to wait for the thread and propagate its result, following the
same pattern used for SpillWriteHandle in pooled_chunk_writer (i.e., ensure
io_handle is joined or transferred only after successful header write and flush
so the background writer cannot outlive the error return).
- Around line 278-282: The test currently uses reader.records().count(), which
counts Err variants too because records() yields Result<Record, _>; change both
occurrences (around the reader.read_header()/assert_eq block and the similar
block at the other location) to explicitly consume each Result and fail on
errors—iterate over reader.records(), call .expect(...) or propagate the error
for each item before incrementing a counter (or collect into a Vec via map(|r|
r.expect(...)).count()) so corrupted records fail the test instead of being
silently counted.
In `@src/lib/sort/pooled_chunk_writer.rs`:
- Around line 18-21: The doc comment in pooled_chunk_writer.rs incorrectly
states permits are released "after each disk write"; update the wording to
reflect that PermitPool::release is invoked after the buffered write hands data
to the OS (i.e., after writer.write_all() on a BufWriter), not after an actual
flush to disk. Edit the comment around StagingBuffer::flush and io_writer_loop
to say permits are released after each block is handed to the OS write buffer
(or similar phrasing) and optionally reference that bgzf_io.rs calls
permit_pool.release() immediately after writer.write_all() to make the behavior
explicit.
---
Outside diff comments:
In `@src/lib/unified_pipeline/fastq.rs`:
- Around line 2533-2550: When a write_all call inside the loop handling
reorder.try_pop_next_with_size() fails, the code returns early after
state.set_error(e) without calling
state.output.write_reorder_state.sub_heap_bytes(heap_size as u64), leaving the
heap accounting stale; modify the error path so that before returning false you
always call state.output.write_reorder_state.sub_heap_bytes(heap_size as u64)
(and optionally update_next_seq(reorder.next_seq()) if semantics require) —
locate the block in the loop around reorder.try_pop_next_with_size(),
output.write_all(&block.data), state.set_error(e) and ensure sub_heap_bytes is
invoked even on write errors (use a small scope guard or explicit call before
return).
---
Duplicate comments:
In `@src/lib/unified_pipeline/base.rs`:
- Around line 5139-5142: The write-side byte counter is updated only after
compressed.pop(), so write_reorder_can_proceed() can admit more output than
queue_memory_limit allows; move or duplicate the heap accounting to occur at the
moment the output is admitted to the write side (i.e., before/when calling
reorder.insert_with_size or immediately prior to enqueue) so
state.write_reorder_state().add_heap_bytes(heap_size as u64) happens before or
atomically with the insert; alternatively ensure write_reorder_can_proceed()
considers Q6/held-compressed bytes by including those bytes in its check. Apply
the same change to the other occurrence around the
reorder.insert_with_size/state.write_reorder_state().add_heap_bytes pair (the
5167-5169 region).
---
Nitpick comments:
In `@src/lib/sort/pooled_chunk_writer.rs`:
- Around line 288-291: The current use of Arc::try_unwrap(pool) means
pool.shutdown() is skipped when other Arc refs exist, risking leaked threads in
tests; change the code to call shutdown unconditionally (i.e., always invoke
pool.shutdown() even if Arc::try_unwrap fails) by making
SortWorkerPool::shutdown accept &self (or add a thread-safe public
shutdown(&self) method) so you can call Arc::as_ref().shutdown() regardless of
ownership; alternatively ensure SortWorkerPool::drop performs the same cleanup
so no shutdown call is skipped.
In `@src/lib/unified_pipeline/bam.rs`:
- Around line 2836-2841: In the hot drain loop starting at
state.output.compressed.pop(), avoid calling batch.estimate_heap_size() twice:
call it once into a local variable (e.g., let heap =
batch.estimate_heap_size()), reuse that value for q7_heap (convert to u64 for
state.q6_track_pop if needed) and pass the original heap to
reorder.insert_with_size(serial, batch, heap); keep the calls to
state.q6_track_pop(...) and state.deadlock_state.record_q7_pop() unchanged but
use the single cached heap value for both places that currently recompute it.
In `@src/lib/unified_pipeline/base.rs`:
- Around line 5893-5943: Add a unit test that verifies write-side backpressure
when constructing OutputPipelineQueues with a non-zero queue_memory_limit:
create an OutputPipelineQueues via OutputPipelineQueues::new with a very small
queue_memory_limit, push TestProcessed/TestSerialized items with out-of-order
serials (use the existing push methods used across tests), and assert that
queue_depths() and are_queues_empty()/has_error()/is_draining() behave to show
the pipeline does not grow beyond the memory limit (i.e., depths respect the
small limit and producers are constrained); name it e.g.
test_output_queues_backpressure_limit and reference OutputPipelineQueues::new,
queue_depths, and the push methods to locate where to add the test.
🪄 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: 551c59c2-738e-49de-acd4-e5e51e707ccf
⛔ Files ignored due to path filters (1)
Cargo.lockis excluded by!**/*.lock
📒 Files selected for processing (23)
.gitignoreCargo.tomlcrates/fgumi-bgzf/src/lib.rscrates/fgumi-bgzf/src/writer.rssrc/lib/bam_io.rssrc/lib/commands/common.rssrc/lib/commands/sort.rssrc/lib/mod.rssrc/lib/sort/bgzf_io.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/mod.rssrc/lib/sort/pooled_bam_writer.rssrc/lib/sort/pooled_chunk_writer.rssrc/lib/sort/raw.rssrc/lib/sort/read_ahead.rssrc/lib/sort/segmented_buf.rssrc/lib/sort/worker_pool.rssrc/lib/system.rssrc/lib/unified_pipeline/bam.rssrc/lib/unified_pipeline/base.rssrc/lib/unified_pipeline/fastq.rstests/integration/main.rstests/integration/test_sort_correctness.rs
✅ Files skipped from review due to trivial changes (8)
- .gitignore
- Cargo.toml
- crates/fgumi-bgzf/src/lib.rs
- src/lib/sort/mod.rs
- crates/fgumi-bgzf/src/writer.rs
- src/lib/system.rs
- src/lib/sort/worker_pool.rs
- src/lib/sort/raw.rs
🚧 Files skipped from review as they are similar to previous changes (5)
- tests/integration/main.rs
- src/lib/sort/segmented_buf.rs
- src/lib/sort/bgzf_io.rs
- tests/integration/test_sort_correctness.rs
- src/lib/sort/inline_buffer.rs
Replace ad-hoc thread spawning in fgumi sort with a structured N+2 model: N pool workers + 1 I/O writer thread + 1 main thread. Worker pool (src/lib/sort/worker_pool.rs): - ArrayQueue (lock-free bounded) replaces crossbeam channels for all inter-step queues; eliminates blocking push/pop in the worker loop - generic_worker_loop pattern: held-item advancement → backpressure priorities → exclusive step ownership → adaptive backoff (10µs min, yield_now + jitter, up to 1ms) - Phase 1: Worker 0 owns ReadInputBlocks exclusively; all workers share DecompressInput and Compress steps - Phase 2: each worker owns its own disjoint chunk files for ReadChunkBlocks (no exclusivity needed); workers share DecompressChunks and Compress steps - park()/unpark() on main thread: workers call unpark() after pushing to main-thread-facing queues, eliminating polling delays - Workers stay alive until SHUTDOWN (not phase completion) to prevent premature exit deadlocks during phase transitions Raw sorter (src/lib/sort/raw.rs): - pool_decompress field and use_sequential_sort() helper gate whether the worker pool handles read/decompress I/O; when enabled, rayon par_sort is replaced with single-threaded sort to preserve N+2 - Fix: merge_chunks_with_index passes pool_decompress=false to build_chunk_sources since the indexing path has no pool consumer set up - Remove dead max_memory_total field (superseded by MemoryLimit::Auto from #236); --max-memory=auto continues to work as intended Staging buffer (src/lib/sort/bgzf_io.rs): - Remove dead is_spill field from CompressJob and StagingBuffer; the worker pool does not need to distinguish spill vs output jobs Correctness: all 2318 tests pass; outputs verified against baseline
Sort pipeline: add PermitPool (num_workers × 4 tokens) acquired by StagingBuffer::flush() before each compress job and released by io_writer_loop after each disk write. Bounds the BTreeMap reorder buffer to at most budget × BGZF_MAX_BLOCK_SIZE and corrects the false backpressure claim in PooledChunkWriter's doc comment. Unified pipeline: add write_reorder_state: ReorderBufferState alongside write_reorder, mirroring the Q2/Q3 admission-control pattern. The compress step checks write_reorder_can_proceed before pushing to Q6; the write step tracks heap_bytes and next_seq with Release/Acquire ordering consistent with the existing Q2/Q3 implementation.
a7f59e0 to
4d364a3
Compare
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
Replaces ad-hoc thread spawning in
fgumi sortwith a structured N+2 model: N pool workers + 1 I/O writer thread + 1 main thread. At--threads 4, this reduces the thread count from ~30 to exactly 6 and cuts wall time by 35–40% vs the previous implementation.Worker pool (
src/lib/sort/worker_pool.rs)ArrayQueueeverywhere: lock-free bounded queues replace crossbeam channels for all inter-step queues, eliminating blocking push/pop in the worker loopgeneric_worker_looppattern: held-item advancement → backpressure priorities → exclusive step ownership → adaptive backoff (10 µs min,yield_now+ jitter, up to 1 ms)ReadInputBlocksexclusively; all workers shareDecompressInputandCompressReadChunkBlocks(no exclusivity needed); workers shareDecompressChunksandCompresspark()/unpark(): workers callunpark()after pushing to main-thread-facing queues, eliminating polling delays between the main thread and the poolRaw sorter (
src/lib/sort/raw.rs)pool_decompressfield anduse_sequential_sort()helper gate worker-pool I/O; when enabled, rayonpar_sortis replaced with single-threaded sort to preserve the N+2 budgetmerge_chunks_with_index(coordinate sort with--write-index) was passingpool_decompress=truetobuild_chunk_sourcesbut then callingnext_record(…, None), which panics forPoolDisksources. Fixed to always usepool_decompress=falsein that path until the indexing path gets a pooled writer (see TODO comment)max_memory_totalfield (superseded byMemoryLimit::Autofrom feat(sort): add --max-memory=auto with system memory detection #236);--max-memory=autocontinues to work as intendedStaging buffer (
src/lib/sort/bgzf_io.rs)is_spillfield fromCompressJobandStagingBuffer; the worker pool does not distinguish spill vs output jobsPerformance (coordinate sort, 60 M records, MacBook Pro M-series)
fgumi is slower than samtools at t1 (pool overhead dominates; single-thread fast path is a follow-on). From t2 onward fgumi leads and the gap widens with thread count.
Test plan
cargo ci-test)cargo ci-fmt && cargo ci-lint)fgumi sort --verifypasses for all four sort orders at t1/t2/t4/t8--write-indexno longer panics (regression for themerge_chunks_with_indexfix)ps -M)Follow-on work
--threads 1, skip the pool entirely (matches samtools t1 behavior; tracked as a separate issue)io_writer_loopBTreeMapreorder buffer is unbounded and can hold many in-flight compressed blocks; mimalloc arena retention also inflates RSS. Will file a separate issue.merge_chunks_with_indexcurrently falls back to non-pooled I/O; replacing withPooledIndexingBamWriterwould extend the N+2 benefit to--write-indexpaths.