Repository navigation
feat(sort): add --max-memory=auto with system memory detection - #236
Conversation
01a42dc to
6cf0b8f
Compare
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #236 +/- ##
==========================================
- Coverage 88.98% 88.97% -0.01%
==========================================
Files 113 114 +1
Lines 55065 55481 +416
==========================================
+ Hits 48997 49366 +369
- Misses 6068 6115 +47 ☔ View full report in Codecov by Sentry. 🚀 New features to boost your workflow:
|
|
@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:
📝 WalkthroughWalkthrough
🚥 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 |
There was a problem hiding this comment.
Actionable comments posted: 5
🧹 Nitpick comments (4)
src/lib/sort/mod.rs (1)
45-45: Hidesegmented_bufunless it is meant to be public API.
SegmentedBufreads like an internal sort-buffer detail. Exporting the module makes it part of the crate's semver surface.Possible change
-pub mod segmented_buf; +pub(crate) mod segmented_buf;🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/mod.rs` at line 45, The module is currently exported as `pub mod segmented_buf`, which exposes internal type `SegmentedBuf` on the crate's public API; either make the module private by removing the `pub` from `pub mod segmented_buf` (so it becomes `mod segmented_buf`) or, if `SegmentedBuf` is intended to be public, keep the module private and explicitly `pub use segmented_buf::SegmentedBuf;` from a controlled public surface with a stable name; update any external references accordingly so only the intended symbols are part of the crate API.src/lib/validation.rs (1)
285-370: Keepparse_memory_size()on the typed validation error path.This is now the outlier in
validation.rs: it returnsanyhow::Resultwith ad-hoc strings while the rest of the module usesFgumiError/Result. That weakens the typed validation contract for callers.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/validation.rs` around lines 285 - 370, parse_memory_size currently returns anyhow::Result and uses anyhow::bail!, making it inconsistent with the rest of validation.rs; change its signature to use the module's typed Result (i.e. Result<u64>) and replace all anyhow::bail!/anyhow::anyhow! uses with construction of the module's FgumiError (or the appropriate validation error variant) so callers receive FgumiError-based errors; update imports to bring Result and FgumiError into scope, and convert the ByteSize parse Err to an FgumiError as well while keeping all existing messages and logic (function: parse_memory_size, type: ByteSize).src/lib/sort/segmented_buf.rs (1)
126-130: Consider adding debug_assert for segment capacity.If caller forgets
reserve_contiguous, the segment could exceedsegment_size. A debug_assert would catch misuse in tests.🔧 Optional: add debug assertion
#[inline] pub fn extend_in_place(&mut self, data: &[u8]) { + debug_assert!( + self.segments.last().map_or(false, |s| s.len() + data.len() <= self.segment_size), + "extend_in_place exceeds segment capacity; use reserve_contiguous first" + ); self.segments.last_mut().expect("segments is never empty").extend_from_slice(data); self.total_len += data.len(); }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/segmented_buf.rs` around lines 126 - 130, Add a debug-only assertion in extend_in_place to ensure the last segment won't exceed the configured segment size if caller forgot to call reserve_contiguous: inside the extend_in_place method (which calls self.segments.last_mut()), assert that last_segment.len() + data.len() <= self.segment_size (use debug_assert! so it only runs in tests/dev), and keep the existing extend_from_slice and total_len update unchanged; mention reserve_contiguous in the assertion message to guide callers.src/commands/sort.rs (1)
300-311: Minor duplication withparse_memory.Consider extracting common logic, but acceptable as-is given minimal scope.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/commands/sort.rs` around lines 300 - 311, parse_memory_reserve duplicates the non-"auto" parsing logic found in parse_memory; extract the shared logic into a small helper (e.g., parse_memory_size_to_usize or parse_bytes_from_str) that trims the input, calls parse_memory_size, converts to usize and returns Result<usize,String>, then change parse_memory_reserve to keep only the "auto" branch and call that helper for the Fixed case and update parse_memory to reuse the same helper to eliminate duplication while preserving current error messages and MemoryReserve::Auto handling.
🤖 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 438-449: When computing per_thread inside the memory_per_thread
branch (involving memory_per_thread, MIN_MEMORY_PER_THREAD, per_thread,
total_budget, available, threads), detect the case where per_thread was bumped
to MIN_MEMORY_PER_THREAD and total_budget = per_thread.saturating_mul(threads)
exceeds available; emit a warn! log (similar format to the existing info!
ByteSize message) that total_budget exceeds available and may cause earlier
spill-to-disk, including values for total_budget, available, per_thread, threads
and margin to aid debugging. Insert this check immediately after computing
total_budget and before returning it.
- Around line 799-808: The test test_resolve_memory_limit_auto relies on actual
system RAM and can fail on low-memory CI; update the test for
resolve_memory_limit (called with MemoryLimit::Auto and MemoryReserve::Auto) to
either mock or query system.total_memory() first and compute a lenient expected
minimum (e.g., use min(system.total_memory() as usize, 4 * 256 * 1024 * 1024) or
assert resolved <= system.total_memory() and resolved >= some fraction like
system.total_memory() / 4) or refactor the test to inject a fake sysinfo::System
so you can supply a controlled total_memory() value—adjust assertions
accordingly and reference resolve_memory_limit, MemoryLimit::Auto,
MemoryReserve::Auto, and system.total_memory() in the change.
In `@src/lib/sort/raw.rs`:
- Around line 755-770: effective_initial_capacity currently returns
initial_capacity unchecked which allows callers to pre-allocate more than the
configured memory/spill budget; change it to clamp the chosen initial capacity
to the configured limit by returning the minimum of the optional
initial_capacity and memory_limit (use std::cmp::min or equivalent) so
effective_initial_capacity() never exceeds self.memory_limit; update the
function that references effective_initial_capacity and the setter
initial_capacity remains unchanged.
In `@src/lib/sort/segmented_buf.rs`:
- Around line 155-167: The slice method currently uses debug_assert to check
that the requested slice (in segmented_buf::slice) does not cross a segment
boundary, which only runs in non-release builds; replace the debug_assert with a
runtime assert! (or return a Result with an explicit error) so the invariant is
enforced in release builds—update the check in pub fn slice(&self, offset:
usize, len: usize) to assert that seg_offset + len <= seg.len() and keep the
same diagnostic text (referencing locate, segments, seg_idx, seg_offset,
seg.len()) so callers get an informative failure instead of undefined behavior
or an obscure panic in release.
In `@src/lib/unified_pipeline/rebalancer.rs`:
- Around line 307-310: Update the documentation for parse_memory_limit to
reflect the current behavior after delegating to parse_memory_size: state that
bare numeric values are interpreted as MB (not bytes), the minimum allowed value
is 256 MiB, and that parsing errors come from
crate::validation::parse_memory_size; reference the parse_memory_limit function
name and the delegated parse_memory_size so reviewers can find and correct the
docblock text accordingly.
---
Nitpick comments:
In `@src/commands/sort.rs`:
- Around line 300-311: parse_memory_reserve duplicates the non-"auto" parsing
logic found in parse_memory; extract the shared logic into a small helper (e.g.,
parse_memory_size_to_usize or parse_bytes_from_str) that trims the input, calls
parse_memory_size, converts to usize and returns Result<usize,String>, then
change parse_memory_reserve to keep only the "auto" branch and call that helper
for the Fixed case and update parse_memory to reuse the same helper to eliminate
duplication while preserving current error messages and MemoryReserve::Auto
handling.
In `@src/lib/sort/mod.rs`:
- Line 45: The module is currently exported as `pub mod segmented_buf`, which
exposes internal type `SegmentedBuf` on the crate's public API; either make the
module private by removing the `pub` from `pub mod segmented_buf` (so it becomes
`mod segmented_buf`) or, if `SegmentedBuf` is intended to be public, keep the
module private and explicitly `pub use segmented_buf::SegmentedBuf;` from a
controlled public surface with a stable name; update any external references
accordingly so only the intended symbols are part of the crate API.
In `@src/lib/sort/segmented_buf.rs`:
- Around line 126-130: Add a debug-only assertion in extend_in_place to ensure
the last segment won't exceed the configured segment size if caller forgot to
call reserve_contiguous: inside the extend_in_place method (which calls
self.segments.last_mut()), assert that last_segment.len() + data.len() <=
self.segment_size (use debug_assert! so it only runs in tests/dev), and keep the
existing extend_from_slice and total_len update unchanged; mention
reserve_contiguous in the assertion message to guide callers.
In `@src/lib/validation.rs`:
- Around line 285-370: parse_memory_size currently returns anyhow::Result and
uses anyhow::bail!, making it inconsistent with the rest of validation.rs;
change its signature to use the module's typed Result (i.e. Result<u64>) and
replace all anyhow::bail!/anyhow::anyhow! uses with construction of the module's
FgumiError (or the appropriate validation error variant) so callers receive
FgumiError-based errors; update imports to bring Result and FgumiError into
scope, and convert the ByteSize parse Err to an FgumiError as well while keeping
all existing messages and logic (function: parse_memory_size, type: ByteSize).
🪄 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: 558531de-ae9a-49c4-9b90-c8959af7235e
📒 Files selected for processing (9)
Cargo.tomlsrc/commands/common.rssrc/commands/sort.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/mod.rssrc/lib/sort/raw.rssrc/lib/sort/segmented_buf.rssrc/lib/unified_pipeline/rebalancer.rssrc/lib/validation.rs
6cf0b8f to
11bf44f
Compare
There was a problem hiding this comment.
Actionable comments posted: 3
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/inline_buffer.rs (1)
153-169:⚠️ Potential issue | 🟠 MajorLarge single records now panic instead of failing cleanly.
Both write paths reserve
header + record.len()against a fixed 256 MiB segment, so any oversized BAM record aborts the sort with an assertion. If that ceiling is intentional, reject it earlier with a typed error; otherwise add an oversized-record fallback.Also applies to: 189-200, 701-713, 727-736
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/inline_buffer.rs` around lines 153 - 169, The code currently hardcodes RECORD_SEGMENT_SIZE = 256 MiB and uses SegmentedBuf::with_capacity(...) in RecordBuffer::with_capacity which causes panics when a single BAM record exceeds that segment; update the write paths (the places that reserve header + record.len() — search for methods that call SegmentedBuf::reserve or perform `reserve(header + record.len())`) to either (A) validate record size early and return a typed error (e.g., OversizedRecordError) before reserving, or (B) implement an oversized-record fallback that allocates a larger temporary segment or grows the SegmentedBuf capacity for that single write so it no longer asserts; ensure you reference RECORD_SEGMENT_SIZE, RecordBuffer::with_capacity, and the SegmentedBuf::reserve/with_capacity calls when making the changes and propagate a clear error type if choosing early rejection.
♻️ Duplicate comments (1)
src/commands/sort.rs (1)
814-887:⚠️ Potential issue | 🟡 MinorThese auto-memory tests are still host-dependent.
On low-memory runners the 256 MiB/thread floor can make
resolved > total, and the fixed-reserve cases can both collapse to the same floored value. Please derive the assertions from an injected/fake total or from the exact runtime formula instead of the host machine.
🤖 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 418-476: The resolve_memory_limit function computes total_budget
differently for MemoryLimit::Fixed and MemoryLimit::Auto but only the Auto path
checks host RAM; after you compute total_budget (both branches) perform a single
post-resolution system memory check using
sysinfo::System::new()/refresh_memory() and system.total_memory() (as usize) to
validate total_budget against actual RAM (total) and margin, logging a warning
if total_budget > total and clamping total_budget to total (or to
total.saturating_sub(resolved reserve) / MIN_MEMORY_PER_THREAD as appropriate)
before returning; update references to total_budget, total, resolve_reserve,
MemoryLimit::Fixed and MemoryLimit::Auto in resolve_memory_limit accordingly.
In `@src/lib/sort/segmented_buf.rs`:
- Around line 38-47: with_capacity currently only allocates a single first
segment (first_cap = capacity.min(segment_size)) which discards multi-segment
capacity hints and causes new() to create a zero-capacity first segment; fix by
computing how many full segments and a remainder are needed (num_full = capacity
/ segment_size, rem = capacity % segment_size) and initialize self.segments with
num_full Vec::with_capacity(segment_size) plus one Vec::with_capacity(rem) when
rem > 0, and ensure at least one segment exists (push Vec::with_capacity(0) if
capacity == 0); also change new() to call
Self::with_capacity(DEFAULT_SEGMENT_SIZE, DEFAULT_SEGMENT_SIZE) (or at least
ensure it creates a non-zero first segment) so DEFAULT_SEGMENT_SIZE is honored;
reference functions/fields: with_capacity, new, segment_size, segments,
total_len.
In `@src/lib/validation.rs`:
- Around line 285-300: Update the documentation for parse_memory_size to state
that bare numeric values are interpreted as mebibytes (MiB) rather than
megabytes (MB); specifically change descriptive lines and examples (including
the plain-number explanation and the example assertions like
parse_memory_size("768")) to use MiB units and ensure the textual list (e.g.,
the bullet that says `"768" -> 768 megabytes`) is revised to `"768" -> 768 MiB
(mebibytes)" so it matches the implementation that multiplies bare integers by
1024*1024.
---
Outside diff comments:
In `@src/lib/sort/inline_buffer.rs`:
- Around line 153-169: The code currently hardcodes RECORD_SEGMENT_SIZE = 256
MiB and uses SegmentedBuf::with_capacity(...) in RecordBuffer::with_capacity
which causes panics when a single BAM record exceeds that segment; update the
write paths (the places that reserve header + record.len() — search for methods
that call SegmentedBuf::reserve or perform `reserve(header + record.len())`) to
either (A) validate record size early and return a typed error (e.g.,
OversizedRecordError) before reserving, or (B) implement an oversized-record
fallback that allocates a larger temporary segment or grows the SegmentedBuf
capacity for that single write so it no longer asserts; ensure you reference
RECORD_SEGMENT_SIZE, RecordBuffer::with_capacity, and the
SegmentedBuf::reserve/with_capacity calls when making the changes and propagate
a clear error type if choosing early rejection.
🪄 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: 7ae4c0c0-d01e-41bb-9dd5-aa29b3054ec0
📒 Files selected for processing (10)
Cargo.tomlsrc/commands/common.rssrc/commands/sort.rssrc/lib/errors.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/mod.rssrc/lib/sort/raw.rssrc/lib/sort/segmented_buf.rssrc/lib/unified_pipeline/rebalancer.rssrc/lib/validation.rs
✅ Files skipped from review due to trivial changes (1)
- Cargo.toml
🚧 Files skipped from review as they are similar to previous changes (1)
- src/lib/sort/raw.rs
11bf44f to
b73c2dd
Compare
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (1)
src/lib/validation.rs (1)
374-379:⚠️ Potential issue | 🟡 MinorError message says "MB" but behavior is MiB.
Line 376 states plain numbers are "interpreted as MB" but code multiplies by
1024 * 1024. Fix for consistency:- - Plain numbers (interpreted as MB): '768', '4096'\n\ + - Plain numbers (interpreted as MiB): '768', '4096'\n\🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/validation.rs` around lines 374 - 379, The error message incorrectly says plain numbers are "interpreted as MB" while the parsing logic multiplies by 1024*1024 (MiB); update the reason string in validation.rs (the format! that builds reason with "{trimmed}" and the list of valid formats) to say "interpreted as MiB" (or otherwise match the 1024*1024 behavior), ensuring the human-readable examples remain accurate and consistent with the code that uses the 1024*1024 multiplier.
🧹 Nitpick comments (1)
src/lib/validation.rs (1)
329-337: Suggestion math is approximate.Line 335 suggests
{mb_value / 1000}GBbut MiB→GB conversion isn't a clean 1000:1 ratio. For 1,000,000 MiB the suggestion would be "1000GB" but actual value is ~1.05 TB. Minor UX nit—not blocking.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/validation.rs` around lines 329 - 337, The error message in the FgumiError::InvalidMemorySize branch uses mb_value / 1000 to suggest "GB", which misrepresents MiB→GB conversion; update the suggestion to compute a precise human-friendly unit from mb_value (e.g., convert MiB→GiB by dividing by 1024 or format a float with two decimals) and include the correct unit label (GiB if using 1024) so the message shown by the invalid_memory_size path (where mb_value is used) accurately reflects the converted size; adjust the format string accordingly in the code that constructs the InvalidMemorySize reason.
🤖 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/inline_buffer.rs`:
- Around line 181-194: The push_coordinate function currently panics on
oversized records; change its API and implementation to return a Result (e.g.,
Result<(), SortError>) instead of panicking: replace the assert! that compares
total_bytes to RECORD_SEGMENT_SIZE with an early Err variant (with a clear
OversizedRecord or MalformedBam error including record.len(), HEADER_SIZE, and
RECORD_SEGMENT_SIZE), and replace the u32::try_from(...).expect(...) with a
fallible conversion that maps the conversion failure into the same Err type.
Apply the same change to the other overloaded/related function range (the code
at lines referenced 730-745) so both accept/return errors rather than panicking,
and update the callers in src/lib/sort/raw.rs (the call sites you noted) to
handle the Result and implement the desired oversized-record path or propagate
the error.
In `@src/lib/sort/raw.rs`:
- Around line 1085-1091: The capacity calculation treats init_cap (from
effective_initial_capacity()) as BAM-bytes-only and underestimates per-record
size because RecordBuffer::with_capacity is given estimated_records computed as
init_cap / 200 without accounting for inline-header bytes and separate
ref-vector reservations added in inline buffer logic; update the code that
computes estimated_records (the init_cap / 200 lines before calling
RecordBuffer::with_capacity) to first compute the total per-record footprint
(BAM payload + inline-header overhead + per-record ref-vector overhead used by
the coordinate and template paths) and divide init_cap by that total footprint
so the call to RecordBuffer::with_capacity(...) budgets against the full
per-record size; apply the same adjustment to the other occurrences noted around
the similar blocks (the ones at the other two locations mentioned).
---
Duplicate comments:
In `@src/lib/validation.rs`:
- Around line 374-379: The error message incorrectly says plain numbers are
"interpreted as MB" while the parsing logic multiplies by 1024*1024 (MiB);
update the reason string in validation.rs (the format! that builds reason with
"{trimmed}" and the list of valid formats) to say "interpreted as MiB" (or
otherwise match the 1024*1024 behavior), ensuring the human-readable examples
remain accurate and consistent with the code that uses the 1024*1024 multiplier.
---
Nitpick comments:
In `@src/lib/validation.rs`:
- Around line 329-337: The error message in the FgumiError::InvalidMemorySize
branch uses mb_value / 1000 to suggest "GB", which misrepresents MiB→GB
conversion; update the suggestion to compute a precise human-friendly unit from
mb_value (e.g., convert MiB→GiB by dividing by 1024 or format a float with two
decimals) and include the correct unit label (GiB if using 1024) so the message
shown by the invalid_memory_size path (where mb_value is used) accurately
reflects the converted size; adjust the format string accordingly in the code
that constructs the InvalidMemorySize reason.
🪄 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: 650530ae-1817-492f-a22e-b087129fa47e
📒 Files selected for processing (10)
Cargo.tomlsrc/commands/common.rssrc/commands/sort.rssrc/lib/errors.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/mod.rssrc/lib/sort/raw.rssrc/lib/sort/segmented_buf.rssrc/lib/unified_pipeline/rebalancer.rssrc/lib/validation.rs
✅ Files skipped from review due to trivial changes (2)
- src/lib/sort/mod.rs
- Cargo.toml
🚧 Files skipped from review as they are similar to previous changes (3)
- src/lib/errors.rs
- src/lib/sort/segmented_buf.rs
- src/commands/sort.rs
b73c2dd to
2c82f8a
Compare
There was a problem hiding this comment.
Actionable comments posted: 2
🧹 Nitpick comments (1)
src/lib/sort/inline_buffer.rs (1)
1673-1673: Add one regression test for the newErrpath.These updates only exercise success cases. One case at
RECORD_SEGMENT_SIZE - HEADER_SIZE + 1and one atTEMPLATE_SEGMENT_SIZE - TEMPLATE_HEADER_SIZE + 1would pin the panic-to-error behavior.Also applies to: 1815-1815, 1860-1860, 1918-1920, 1951-1953
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/inline_buffer.rs` at line 1673, The tests currently only exercise the successful path of InlineBuffer::push (see calls to buffer.push(&record, key).expect(...)); add regression tests that force the new Err path by constructing records/templates whose serialized size equals RECORD_SEGMENT_SIZE - HEADER_SIZE + 1 and TEMPLATE_SEGMENT_SIZE - TEMPLATE_HEADER_SIZE + 1, then call buffer.push(&record, key) and assert it returns Err (rather than using expect). Locate the existing test module that calls buffer.push and add two new cases: one for the record-segment overflow and one for the template-segment overflow, using the same buffer setup and keys so they pin the panic-to-error behavior for push and related push_template code paths.
🤖 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 418-445: The resolve_memory_limit function currently silently
saturates on bad inputs (threads==0 and arithmetic overflows); change it to
return a Result<usize, Error> (or anyhow::Error) and validate inputs instead of
saturating: in resolve_memory_limit check threads==0 and return Err, replace
saturating_mul/unwraps with checked arithmetic (use checked_mul and
usize::try_from(total_memory).map_err(...)) and return Err when overflow or when
computed budget > available (instead of saturating to usize::MAX); keep logic
for MemoryLimit::Fixed and MemoryLimit::Auto but fail on invalid arithmetic or
when per-thread MIN_MEMORY_PER_THREAD cannot be met, and update all callers of
resolve_memory_limit to handle the Result; refer to resolve_memory_limit,
MemoryLimit::Fixed, MemoryLimit::Auto, resolve_reserve, memory_per_thread, and
MIN_MEMORY_PER_THREAD when making the changes.
In `@src/lib/sort/inline_buffer.rs`:
- Around line 184-198: push_coordinate currently calls
extract_coordinate_key_inline(record, self.nref) before validating the record is
long enough for fixed-offset reads (e.g., bam_fields::flags at offsets 14-15),
which causes panics on truncated records; fix by validating the minimum required
length (or delegating to a fallible extract that returns Result) before any
fixed-offset indexing: either add an explicit check in push_coordinate that
record.len() >= MIN_BAM_HEADER_LEN (and return an Err with context) or change
extract_coordinate_key_inline to return anyhow::Result and propagate its error
(use ?), ensuring no raw indexing into record occurs without a prior length
check (reference push_coordinate and extract_coordinate_key_inline and the
bam_fields::flags offsets).
---
Nitpick comments:
In `@src/lib/sort/inline_buffer.rs`:
- Line 1673: The tests currently only exercise the successful path of
InlineBuffer::push (see calls to buffer.push(&record, key).expect(...)); add
regression tests that force the new Err path by constructing records/templates
whose serialized size equals RECORD_SEGMENT_SIZE - HEADER_SIZE + 1 and
TEMPLATE_SEGMENT_SIZE - TEMPLATE_HEADER_SIZE + 1, then call buffer.push(&record,
key) and assert it returns Err (rather than using expect). Locate the existing
test module that calls buffer.push and add two new cases: one for the
record-segment overflow and one for the template-segment overflow, using the
same buffer setup and keys so they pin the panic-to-error behavior for push and
related push_template code paths.
🪄 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: 0ae32e35-5fde-4be5-b0b7-634f531f6770
📒 Files selected for processing (10)
Cargo.tomlsrc/commands/common.rssrc/commands/sort.rssrc/lib/errors.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/mod.rssrc/lib/sort/raw.rssrc/lib/sort/segmented_buf.rssrc/lib/unified_pipeline/rebalancer.rssrc/lib/validation.rs
✅ Files skipped from review due to trivial changes (1)
- src/lib/sort/raw.rs
🚧 Files skipped from review as they are similar to previous changes (5)
- src/lib/sort/mod.rs
- src/lib/validation.rs
- src/lib/errors.rs
- src/lib/unified_pipeline/rebalancer.rs
- src/lib/sort/segmented_buf.rs
2c82f8a to
66e268c
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
♻️ Duplicate comments (5)
src/lib/validation.rs (1)
373-379:⚠️ Potential issue | 🟡 MinorKeep the fallback plain-number help in MiB.
The parser treats bare integers as mebibytes, but this error text still says “MB”. That makes the CLI guidance inconsistent again.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/validation.rs` around lines 373 - 379, The error message for FgumiError::InvalidMemorySize is misleading: it says plain numbers are interpreted as "MB" but the parser treats bare integers as mebibytes; update the formatted reason string in the InvalidMemorySize construction (the block referencing 'trimmed' and FgumiError::InvalidMemorySize) to mention "MiB" (mebibytes) instead of "MB" and ensure the examples and list items consistently use MiB where plain integers are described.src/commands/sort.rs (2)
443-479:⚠️ Potential issue | 🟠 MajorClamp auto budgets to the memory that's actually left.
If the reserve leaves less than
MIN_MEMORY_PER_THREADtotal or per thread, this code still returns the floor even when it exceedsavailableandtotal.execute_sort()then forwards that value intoRawExternalSorter::memory_limitand the autoinitial_capacity, so--max-memory=autocan still target an impossible budget instead of spilling earlier. Clamp toavailable/total, or fail once the reserve leaves less than the minimum.Also applies to: 484-494
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/commands/sort.rs` around lines 443 - 479, The auto-budget calculation can yield budgets larger than what's actually available when the reserve leaves less than MIN_MEMORY_PER_THREAD; update the logic in the block around resolve_reserve/ memory_per_thread to clamp computed budgets to available (and per-thread to available/threads) or return an error when the reserve makes meeting MIN_MEMORY_PER_THREAD impossible. Specifically, in the memory_per_thread branch ensure per_thread = min(per_thread, available / threads) (and then compute budget = min(per_thread.checked_mul(threads)?, available)), and in the fixed-total branch set budget = min(available.max(MIN_MEMORY_PER_THREAD), available) or return an early error; adjust callers like execute_sort / RawExternalSorter::memory_limit to expect the clamped value (or propagate the error) so --max-memory=auto never exceeds total/available.
839-855:⚠️ Potential issue | 🟡 MinorThese auto-memory tests still depend on host RAM.
On low-memory CI,
test_resolve_memory_limit_autocan legitimately hit the 256 MiB/thread floor and breakresolved <= total, and both fixed-reserve cases can collapse to the same floor solarge_reserve < small_reserveno longer holds. Inject a fake total-memory value or derive the assertions from the queriedtotal.Also applies to: 899-915
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/commands/sort.rs` around lines 839 - 855, The tests (e.g., test_resolve_memory_limit_auto) rely on the real system RAM and can fail on low-memory CI; change them to not assume host memory by either injecting a fake total-memory value into the resolve_memory_limit call (or by mocking sysinfo::System) or by computing expected assertions from the captured total variable; specifically, adjust the assertions that compare resolved to MIN_MEMORY_PER_THREAD, total, and comparisons between large/small reserves to use the derived min_expected and total variables (and update the other related tests that call resolve_memory_limit with MemoryLimit::Auto / MemoryReserve::* to follow the same pattern).src/lib/sort/segmented_buf.rs (2)
38-47:⚠️ Potential issue | 🟠 Major
with_capacity()still makes >1-segment hints mostly a no-op.Anything above one segment only reserves the outer
segmentsVec; the data buffer itself still starts with a single segment, andnew()starts that first segment at capacity 0. That means the multi-GiBinitial_capacitypath fromsrc/lib/sort/raw.rsis barely honored above 256 MiB, and the default constructor still lets the first segment grow/reallocate like a plainVec. At minimum,new()should start with one fullDEFAULT_SEGMENT_SIZEsegment.Also applies to: 50-54
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/segmented_buf.rs` around lines 38 - 47, with_capacity currently only reserves the outer segments Vec and always starts with a single small first segment; update SegmentedBuf::with_capacity to allocate actual segment buffers to match the requested capacity by computing num_segments = (capacity + segment_size - 1) / segment_size, pushing that many Vec::with_capacity(...) into segments (use remainder for the last segment capacity if desired) and set total_len appropriately; similarly ensure SegmentedBuf::new initializes segments with one Vec::with_capacity(DEFAULT_SEGMENT_SIZE) instead of capacity 0 so the first segment is full-sized; refer to symbols with_capacity, new, DEFAULT_SEGMENT_SIZE, segment_size, segments, and total_len to locate and modify the code.
134-140: 🛠️ Refactor suggestion | 🟠 MajorMake
extend_in_place()enforce its precondition in release.
debug_assert!disappears in release, so an undersizedreserve_contiguous()lets the innerVecgrow pastsegment_size, which breaks the fixed-segment invariant thatlocate()assumes. Useassert!or return aResult.Patch sketch
- debug_assert!( + assert!( self.segments.last().is_some_and(|s| s.len() + data.len() <= self.segment_size), "extend_in_place exceeds segment capacity; use reserve_contiguous first" );🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/sort/segmented_buf.rs` around lines 134 - 140, extend_in_place currently uses debug_assert!, which is omitted in release builds and can allow the last inner Vec (self.segments.last_mut()) to grow beyond self.segment_size, breaking invariants relied on by locate and others; change extend_in_place to enforce the precondition in release by either replacing debug_assert! with assert! (keeping the same error message) or by changing the signature to return Result<(), Error> and explicitly check that self.segments.last().is_some_and(|s| s.len() + data.len() <= self.segment_size) before extending, returning an error if the check fails; ensure you still update self.total_len only after a successful extend and reference reserve_contiguous and locate in the error message or docs so callers know to call reserve_contiguous appropriately.
🤖 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/segmented_buf.rs`:
- Around line 1-12: Add the repo-wide unsafe guard to this new module by
inserting the crate attribute #![deny(unsafe_code)] at the top of
src/lib/sort/segmented_buf.rs (before the doc comment) so the SegmentedBuf
implementation and constants like DEFAULT_SEGMENT_SIZE are compiled under the
repository rule disallowing unsafe code. Ensure the attribute is the first
non-comment item in the file so the module follows the project's
deny(unsafe_code) policy.
---
Duplicate comments:
In `@src/commands/sort.rs`:
- Around line 443-479: The auto-budget calculation can yield budgets larger than
what's actually available when the reserve leaves less than
MIN_MEMORY_PER_THREAD; update the logic in the block around resolve_reserve/
memory_per_thread to clamp computed budgets to available (and per-thread to
available/threads) or return an error when the reserve makes meeting
MIN_MEMORY_PER_THREAD impossible. Specifically, in the memory_per_thread branch
ensure per_thread = min(per_thread, available / threads) (and then compute
budget = min(per_thread.checked_mul(threads)?, available)), and in the
fixed-total branch set budget = min(available.max(MIN_MEMORY_PER_THREAD),
available) or return an early error; adjust callers like execute_sort /
RawExternalSorter::memory_limit to expect the clamped value (or propagate the
error) so --max-memory=auto never exceeds total/available.
- Around line 839-855: The tests (e.g., test_resolve_memory_limit_auto) rely on
the real system RAM and can fail on low-memory CI; change them to not assume
host memory by either injecting a fake total-memory value into the
resolve_memory_limit call (or by mocking sysinfo::System) or by computing
expected assertions from the captured total variable; specifically, adjust the
assertions that compare resolved to MIN_MEMORY_PER_THREAD, total, and
comparisons between large/small reserves to use the derived min_expected and
total variables (and update the other related tests that call
resolve_memory_limit with MemoryLimit::Auto / MemoryReserve::* to follow the
same pattern).
In `@src/lib/sort/segmented_buf.rs`:
- Around line 38-47: with_capacity currently only reserves the outer segments
Vec and always starts with a single small first segment; update
SegmentedBuf::with_capacity to allocate actual segment buffers to match the
requested capacity by computing num_segments = (capacity + segment_size - 1) /
segment_size, pushing that many Vec::with_capacity(...) into segments (use
remainder for the last segment capacity if desired) and set total_len
appropriately; similarly ensure SegmentedBuf::new initializes segments with one
Vec::with_capacity(DEFAULT_SEGMENT_SIZE) instead of capacity 0 so the first
segment is full-sized; refer to symbols with_capacity, new,
DEFAULT_SEGMENT_SIZE, segment_size, segments, and total_len to locate and modify
the code.
- Around line 134-140: extend_in_place currently uses debug_assert!, which is
omitted in release builds and can allow the last inner Vec
(self.segments.last_mut()) to grow beyond self.segment_size, breaking invariants
relied on by locate and others; change extend_in_place to enforce the
precondition in release by either replacing debug_assert! with assert! (keeping
the same error message) or by changing the signature to return Result<(), Error>
and explicitly check that self.segments.last().is_some_and(|s| s.len() +
data.len() <= self.segment_size) before extending, returning an error if the
check fails; ensure you still update self.total_len only after a successful
extend and reference reserve_contiguous and locate in the error message or docs
so callers know to call reserve_contiguous appropriately.
In `@src/lib/validation.rs`:
- Around line 373-379: The error message for FgumiError::InvalidMemorySize is
misleading: it says plain numbers are interpreted as "MB" but the parser treats
bare integers as mebibytes; update the formatted reason string in the
InvalidMemorySize construction (the block referencing 'trimmed' and
FgumiError::InvalidMemorySize) to mention "MiB" (mebibytes) instead of "MB" and
ensure the examples and list items consistently use MiB where plain integers are
described.
🪄 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: 7bf3b282-a8fc-4333-ac66-4cdc5e7d0a4f
📒 Files selected for processing (10)
Cargo.tomlsrc/commands/common.rssrc/commands/sort.rssrc/lib/errors.rssrc/lib/sort/inline_buffer.rssrc/lib/sort/mod.rssrc/lib/sort/raw.rssrc/lib/sort/segmented_buf.rssrc/lib/unified_pipeline/rebalancer.rssrc/lib/validation.rs
✅ Files skipped from review due to trivial changes (2)
- src/lib/sort/mod.rs
- src/lib/sort/inline_buffer.rs
🚧 Files skipped from review as they are similar to previous changes (3)
- src/lib/errors.rs
- src/lib/unified_pipeline/rebalancer.rs
- Cargo.toml
Add automatic memory detection for `fgumi sort` using the `sysinfo` crate. When `--max-memory=auto` (now the default), the sort command detects total system memory at startup, subtracts a safety margin of max(4 GiB, 10%), and divides by thread count. This replaces the previous fixed 768M default that was suboptimal for most machines. Key changes: - Add `MemoryLimit` enum with `Auto` and `Fixed(usize)` variants - Add `resolve_memory_limit` for one-time memory calculation at startup - Add `initial_capacity` to `RawExternalSorter` to decouple buffer pre-allocation from spill threshold (768 MiB/thread for auto mode) - Add `SegmentedBuf` to replace `Vec<u8>` in sort buffers, eliminating the O(n) memcpy and transient 2x peak memory from Vec doubling at multi-GiB scale (256 MiB fixed-size segments, zero-copy growth) - Consolidate three separate memory parsers into a single canonical `parse_memory_size` in `fgumi_lib::validation` - Make `sysinfo` a non-optional dependency; remove feature gate from `validate_against_system_memory` Benchmarks on 6 GiB BAM (60M records, template-coordinate, 4 threads): - fgumi auto: 27s, 29.4 GiB RSS (no spills) - fgumi explicit 14.4G/thd: 26s, 26.9 GiB RSS (no spills) - fgumi explicit 4G/thd: 47s, 23.0 GiB RSS (5 spills) - samtools 14.4G/thd: 67s, 26.7 GiB RSS (no spills) - samtools 4G/thd: 87s, 22.1 GiB RSS (1 spill)
66e268c to
8c80e7d
Compare
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
* perf(sort): implement N+2 worker-pool model for parallel sort I/O 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 * fix: bound write reorder buffers to prevent unbounded memory growth 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.
Summary
--max-memory=auto(now the default) that detects system memory and subtracts a configurable reserve to leave room for co-running processes--memory-reserveoption:auto(default) reservesmin(10 GiB, 50% of RAM), or an explicit value like12GiBfor pipelines alongsidebwa memSegmentedBuf(fixed-size 256 MiB segments) to replaceVec<u8>in sort buffers, eliminating the O(n) memcpy and transient 2x peak memory from Vec doubling at multi-GiB scaleinitial_capacitytoRawExternalSorterto decouple buffer pre-allocation (768 MiB/thread) from the spill-to-disk thresholdparse_memory_sizeinfgumi_lib::validationsysinfoa non-optional dependency; always validate memory against system limitsAuto reserve defaults
Benchmarks (6 GiB BAM, 60M records, template-coordinate, 4 threads)
Test plan
cargo ci-test— all tests pass (includes newSegmentedBuf,resolve_reserve, andparse_memory_reservetests)cargo ci-fmt && cargo ci-lint— cleanfgumi sort -i large.bam -o out.bam --order template-coordinate -@ 4works with auto (default)fgumi sort -m 4GiBstill works as explicit overridefgumi sort --memory-reserve 12GiBreserves extra for pipeline usefgumi sort -m 256MiBcorrectly spills to disk