Repository navigation
fix(pipeline): stream rejects to disk in simplex/duplex/codec/correct - #293
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #293 +/- ##
==========================================
+ Coverage 90.19% 90.26% +0.06%
==========================================
Files 124 124
Lines 60598 60618 +20
==========================================
+ Hits 54657 54716 +59
+ Misses 5941 5902 -39 ☔ View full report in Codecov by Sentry. 🚀 New features to boost your workflow:
|
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (5)
🚧 Files skipped from review as they are similar to previous changes (5)
📝 WalkthroughWalkthroughRejects are no longer buffered per-batch; when 🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
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.
|
3072739 to
e7ed89c
Compare
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
e7ed89c to
cadc96c
Compare
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
src/lib/commands/codec.rs (1)
566-669:⚠️ Potential issue | 🟠 MajorFinalize the rejects writer even when the pipeline fails.
Line 654 returns before Lines 667-669 on any pipeline error, so a partially written rejects BAM can miss finalization/EOF. Capture the pipeline result, finalize rejects, then propagate the pipeline error.
Proposed fix
- let groups_processed = run_bam_pipeline_from_reader( + let pipeline_result = run_bam_pipeline_from_reader( pipeline_config, reader, input_header, &self.io.output, @@ ) - .map_err(|e| anyhow::anyhow!("Pipeline error: {e}"))?; + .map_err(|e| anyhow::anyhow!("Pipeline error: {e}")); - // ========== Post-pipeline: Aggregate metrics and close rejects writer ========== + let finish_rejects_result: Result<()> = (|| { + if let Some(rw_arc) = rejects_writer { + if let Some(writer) = rw_arc.lock().take() { + writer.finish().context("Failed to finish rejects file")?; + info!("Rejected reads streamed to rejects file during processing"); + } + } + Ok(()) + })(); + + let groups_processed = pipeline_result?; + finish_rejects_result?; + + // ========== Post-pipeline: Aggregate metrics ========== let mut total_groups = 0u64; let mut merged_stats = CodecConsensusStats::default(); @@ - // Rejects were streamed during serialize; now finalize the writer. - if let Some(rw_arc) = rejects_writer { - if let Some(writer) = rw_arc.lock().take() { - writer.finish().context("Failed to finish rejects file")?; - info!("Rejected reads streamed to rejects file during processing"); - } - } - // Log statistics🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/commands/codec.rs` around lines 566 - 669, The pipeline call currently uses .map_err(...) which returns early on error before finalizing the rejects writer (rejects_writer / rw_arc), so a partially-written rejects BAM may never be finished; change to capture the Result from run_bam_pipeline_from_reader into a variable (e.g., let pipeline_result = run_bam_pipeline_from_reader(...);), then after aggregating metrics always finalize the rejects writer (if let Some(rw_arc) = rejects_writer { if let Some(writer) = rw_arc.lock().take() { writer.finish().context("Failed to finish rejects file")?; } }), and finally propagate the pipeline_result error (e.g., pipeline_result.map_err(|e| anyhow::anyhow!("Pipeline error: {e}"))?), ensuring rejects_writer finalization runs regardless of pipeline success or failure.src/lib/commands/duplex.rs (1)
702-740:⚠️ Potential issue | 🟠 MajorRejects are still buffered per processed batch.
serialize_fnstreams rejects only afterprocess_fnhas builtall_rejects, so a high-rejection MI batch can still retain all rejected raw records until serialization. Consider flushing rejects per MI group through the shared writer or otherwise keepingDuplexProcessedBatchfree of reject payloads.Also applies to: 755-765
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/commands/duplex.rs` around lines 702 - 740, The batch loop accumulates per-MI rejects into all_rejects and returns them in DuplexProcessedBatch, causing high-memory buffering; instead, when track_rejects is enabled, immediately drain caller.take_rejected_reads() for each RawMiGroup and write/stream them via the shared reject writer (the same writer used by serialize_fn) before moving to the next MI, and remove the reject payload from DuplexProcessedBatch (or keep only lightweight metadata). Update process_fn to call the shared writer inside the loop (after caller.take_rejected_reads()) and stop extending all_rejects so DuplexProcessedBatch no longer holds raw reject payloads; keep merging stats (batch_stats, batch_overlapping) as before.
🧹 Nitpick comments (2)
src/lib/commands/codec.rs (1)
506-508: Update the stale rejects comment.The comment still says rejects buffering is preserved and tracked as follow-up, but this implementation now streams rejects during serialization.
Proposed fix
- // Per-thread metrics accumulator: bounded metric memory, no unbounded - // queue. Rejects buffering semantics are preserved (see follow-up). + // Per-thread metrics accumulator: bounded metric memory, no unbounded + // queue. Rejects are streamed separately during serialization.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/commands/codec.rs` around lines 506 - 508, The comment above the PerThreadAccumulator::<CollectedCodecMetrics>::new(num_threads) instantiation is outdated: it claims "Rejects buffering semantics are preserved (see follow-up)" but the current implementation streams rejects during serialization; update the comment to accurately reflect that rejects are streamed during serialization (not buffered) and that the accumulator enforces bounded per-thread metric memory. Locate the comment near the collected_metrics / PerThreadAccumulator::<CollectedCodecMetrics>::new(...) and change the wording to state that rejects are streamed during serialization and that the accumulator enforces bounded metric memory with no unbounded queue.src/lib/commands/simplex.rs (1)
499-500: Update the stale rejects comment.Rejects are no longer “preserved” for a follow-up; they are streamed during serialization now.
Proposed comment update
- // Per-thread metrics accumulator: bounded metric memory, no unbounded - // queue. Rejects buffering semantics are preserved (see follow-up). + // Per-thread metrics accumulator: bounded metric memory, no unbounded + // queue. Rejects are streamed during serialize, not stored here.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/commands/simplex.rs` around lines 499 - 500, Update the stale inline comment next to the Per-thread metrics accumulator so it no longer says "Rejects buffering semantics are preserved (see follow-up)"; change it to state that rejects are streamed during serialization (e.g., replace that clause with "Rejects are streamed during serialization.") so the comment accurately reflects current behavior in the code surrounding the Per-thread metrics accumulator.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@CHANGELOG.md`:
- Around line 10-11: Remove the stale clause "Rejects buffering in the consensus
commands is unchanged and tracked separately." from the release-note bullet that
begins "Replace unbounded `SegQueue<CollectedXxxMetrics>`..." and either delete
the clause entirely or replace it with a corrected statement noting that
rejected records are now streamed to disk (as documented in the next bullet) so
reject buffering is no longer unchanged; ensure the two bullets are consistent
so readers aren't confused about reject-buffering behavior.
In `@src/lib/per_thread_accumulator.rs`:
- Around line 81-91: into_slots currently mutates shared slots via
std::mem::take when Arc::try_unwrap fails; change its behavior to avoid mutating
state held by other Arc owners by altering the signature to return
Result<Vec<A>, Arc<Self>> (keep the receiver as Arc<Self>), use
Arc::try_unwrap(self) => Ok(inner) branch to consume and return the Vec as
before, and in the Err(arc) branch return Err(arc) without performing
std::mem::take or locking/mutating slots — callers can then call the
non-consuming slots() accessor to inspect or reduce values safely; refer to
into_slots, Arc::try_unwrap, and slots() when making the change.
---
Outside diff comments:
In `@src/lib/commands/codec.rs`:
- Around line 566-669: The pipeline call currently uses .map_err(...) which
returns early on error before finalizing the rejects writer (rejects_writer /
rw_arc), so a partially-written rejects BAM may never be finished; change to
capture the Result from run_bam_pipeline_from_reader into a variable (e.g., let
pipeline_result = run_bam_pipeline_from_reader(...);), then after aggregating
metrics always finalize the rejects writer (if let Some(rw_arc) = rejects_writer
{ if let Some(writer) = rw_arc.lock().take() { writer.finish().context("Failed
to finish rejects file")?; } }), and finally propagate the pipeline_result error
(e.g., pipeline_result.map_err(|e| anyhow::anyhow!("Pipeline error: {e}"))?),
ensuring rejects_writer finalization runs regardless of pipeline success or
failure.
In `@src/lib/commands/duplex.rs`:
- Around line 702-740: The batch loop accumulates per-MI rejects into
all_rejects and returns them in DuplexProcessedBatch, causing high-memory
buffering; instead, when track_rejects is enabled, immediately drain
caller.take_rejected_reads() for each RawMiGroup and write/stream them via the
shared reject writer (the same writer used by serialize_fn) before moving to the
next MI, and remove the reject payload from DuplexProcessedBatch (or keep only
lightweight metadata). Update process_fn to call the shared writer inside the
loop (after caller.take_rejected_reads()) and stop extending all_rejects so
DuplexProcessedBatch no longer holds raw reject payloads; keep merging stats
(batch_stats, batch_overlapping) as before.
---
Nitpick comments:
In `@src/lib/commands/codec.rs`:
- Around line 506-508: The comment above the
PerThreadAccumulator::<CollectedCodecMetrics>::new(num_threads) instantiation is
outdated: it claims "Rejects buffering semantics are preserved (see follow-up)"
but the current implementation streams rejects during serialization; update the
comment to accurately reflect that rejects are streamed during serialization
(not buffered) and that the accumulator enforces bounded per-thread metric
memory. Locate the comment near the collected_metrics /
PerThreadAccumulator::<CollectedCodecMetrics>::new(...) and change the wording
to state that rejects are streamed during serialization and that the accumulator
enforces bounded metric memory with no unbounded queue.
In `@src/lib/commands/simplex.rs`:
- Around line 499-500: Update the stale inline comment next to the Per-thread
metrics accumulator so it no longer says "Rejects buffering semantics are
preserved (see follow-up)"; change it to state that rejects are streamed during
serialization (e.g., replace that clause with "Rejects are streamed during
serialization.") so the comment accurately reflects current behavior in the code
surrounding the Per-thread metrics accumulator.
🪄 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: 5504d2e3-1ad7-46e8-ac38-7f82aefbc72c
📒 Files selected for processing (10)
CHANGELOG.mdsrc/lib/commands/clip.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/duplex.rssrc/lib/commands/filter.rssrc/lib/commands/group.rssrc/lib/commands/simplex.rssrc/lib/mod.rssrc/lib/per_thread_accumulator.rs
cadc96c to
ee7ac26
Compare
4afa6f6 to
d907ac6
Compare
|
Addressed review feedback and rebased:
Branch is rebased onto current The |
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 6
🤖 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/commands/correct.rs`:
- Around line 859-875: The current flush_rejects closure (used in process_fn and
writing via rejects_writer_for_process) writes rejects under a mutex but outside
the pipeline's reorder/serialize stage, so reject BAMs can be out-of-order while
still carrying a sorted header; fix by either routing reject records through the
same ordered pipeline/serialize stage (i.e., stop calling flush_rejects
pre-serialize and instead enqueue or emit rejects alongside normal records so
the reorder/serialize step handles them) or explicitly mark the rejects BAM
header as unsorted/unknown when creating/initializing rejects_writer_for_process
(update the header HD SO tag to "unsorted" or "unknown") so downstream tools
won’t assume sort order.
- Around line 1032-1042: The current code uses pipeline_result? before checking
rejects_finish_result which hides a rejects.finish() error; change the control
flow to first capture pipeline_result without early-return (e.g., let
pipeline_res = pipeline_result;), then if let Some(result) =
rejects_finish_result { result? } to propagate any finish() error first, and
only after successful rejects finalization return or ? the pipeline result
(e.g., match pipeline_res { Ok(records_written) => records_written, Err(e) =>
return Err(e) }). Ensure you reference and use the existing symbols
rejects_writer, rejects_finish_result, pipeline_result, and writer.finish() so
that a rejects.finish() failure is not dropped when the pipeline also fails.
In `@src/lib/commands/duplex.rs`:
- Around line 703-719: The current flush_rejects closure writes worker-produced
reject records under a mutex but not in input order, which can falsify sort
metadata in the rejects BAM; fix by either (A) routing rejects through an
ordered stage: have worker code send (sequence_id, Vec<Vec<u8>>) to a single
ordered_consumer on the main thread which calls w.write_raw_record in sequence
(introduce a channel and an ordered drain in place of direct flush_rejects
usage), or (B) make the rejects header explicitly unordered before any writes by
mutating input_header (remove or set the HD/SO tag to "unsorted") so the
contract matches nondeterministic ordering; locate changes around flush_rejects,
rejects_writer_for_process, w.write_raw_record, and input_header.
- Around line 786-797: The pipeline error can short-circuit before we observe
rejects_writer.finish() failures: ensure the rejects writer is finalized and any
finish() error is surfaced even when pipeline_result is Err. Change control flow
around rejects_finish_result and pipeline_result so you first compute
rejects_finish_result (already done), then handle both results: if
rejects_finish_result is Some(Err) return that error (or combine it), otherwise
if pipeline_result is Err return the pipeline error; only proceed to set
groups_processed when pipeline_result is Ok. Reference symbols: rejects_writer,
rejects_finish_result, pipeline_result, groups_processed, and the
writer.finish() call — ensure rejects_writer.lock().take() ->
writer.finish().context(...) is awaited/checked before returning
pipeline_result.map_err(...)? so finish() failures are not masked.
In `@src/lib/commands/simplex.rs`:
- Around line 599-615: The current flush_rejects closure writes per-worker
rejects under a mutex but does so in arrival order, allowing output order to
diverge from input order when process_fn runs concurrently; change to preserve
input ordering by tagging each rejected record with its input sequence/index and
funneling all rejects into a single ordered writer that emits records in
increasing sequence (e.g., have workers push (seq, Vec<u8>) into a shared
ordered queue or BTreeMap and have rejects_writer_for_process take/flush entries
only when the next expected sequence is present), or alternatively buffer
per-worker batches and merge/flush them by sequence before calling
write_raw_record; apply the same change to the other reject-flush site (lines
~627-655) so rejects are always emitted in input order.
- Around line 699-710: The code returns early on pipeline_result with `?` before
checking `rejects_finish_result`, so a failed `rejects_writer.finish()` can be
lost; change control flow to evaluate and propagate `rejects_finish_result` (the
result of `rejects_writer.and_then(...).map(|writer| writer.finish()...)`)
before unwrapping `pipeline_result` into `groups_processed`—i.e., first if
`rejects_finish_result` is Some, call `result?` and log, then map/`?` the
`pipeline_result` into `groups_processed` so any rejects finalization error is
surfaced prior to returning the pipeline error. Ensure you reference and update
`rejects_finish_result`, `rejects_writer`, and `pipeline_result` in that order.
🪄 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: 11819411-946d-43dc-a693-5ff0cecc6069
📒 Files selected for processing (4)
src/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/duplex.rssrc/lib/commands/simplex.rs
✅ Files skipped from review due to trivial changes (1)
- src/lib/commands/codec.rs
d907ac6 to
e338c91
Compare
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 1
🧹 Nitpick comments (1)
tests/integration/test_correct_command.rs (1)
187-187: Nit: collapse literal.Diff
- let expected_rejects = 10 + 10 + 10; + let expected_rejects = 30;Or derive from the family sizes to keep them in sync.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@tests/integration/test_correct_command.rs` at line 187, The test sets expected_rejects as a sum of three identical literals (let expected_rejects = 10 + 10 + 10;) — replace this with a single literal (e.g. 30) or, preferably, compute it from the existing family-size variables used in the test (so expected_rejects is derived from those family size identifiers) to keep the values in sync and avoid magic repetition.
🤖 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/commands/duplex.rs`:
- Around line 77-79: Update the stale comments that say rejects are streamed
during serialize to reflect the current behavior: rejects are now streamed
during process_fn. Specifically, edit the doc comment on the per-thread
accumulator (the comment above the accumulator struct/variable in duplex.rs
referencing "Rejected records are streamed directly to the rejects BAM during
serialize") and the later comment near the code referencing serialize (around
the block related to serialize/process flow) to say "streamed during process_fn"
and, if helpful, mention the process_fn function name where the streaming now
occurs so readers can locate the implementation.
---
Nitpick comments:
In `@tests/integration/test_correct_command.rs`:
- Line 187: The test sets expected_rejects as a sum of three identical literals
(let expected_rejects = 10 + 10 + 10;) — replace this with a single literal
(e.g. 30) or, preferably, compute it from the existing family-size variables
used in the test (so expected_rejects is derived from those family size
identifiers) to keep the values in sync and avoid magic repetition.
🪄 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: 00399063-a873-409a-af88-d01e0e1164f2
📒 Files selected for processing (8)
crates/fgumi-consensus/src/duplex_caller.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/duplex.rstests/integration/helpers/assertions.rstests/integration/test_bgzf_eof.rstests/integration/test_correct_command.rstests/integration/test_duplex_command.rs
🚧 Files skipped from review as they are similar to previous changes (4)
- tests/integration/helpers/assertions.rs
- tests/integration/test_bgzf_eof.rs
- tests/integration/test_duplex_command.rs
- src/lib/commands/correct.rs
779a38d to
d6a96fe
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
🧹 Nitpick comments (1)
src/lib/commands/duplex.rs (1)
706-722: Nit: mutex is held across the entire group's record loop.This is fine (and arguably desirable — keeps a single MI group's rejects contiguous), but worth noting that for very large groups this serializes workers. The BGZF thread pool still compresses concurrently, so it's unlikely to matter. No change required.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/commands/duplex.rs` around lines 706 - 722, The mutex in the closure flush_rejects is held for the entire loop over a group's records (rw_arc lock is acquired before iterating and calling write_raw_record), which can serialize workers for very large groups; to reduce contention, acquire the lock only for the minimal write operation (e.g., take a short-lived guard per record or move/copy each raw record out and lock just to call write_raw_record), ensuring you still call write_raw_record on the same writer (guard.as_mut()) and preserve group contiguity as needed; update the flush_rejects closure (and uses of rw_arc, guard, and write_raw_record) 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 `@crates/fgumi-consensus/src/duplex_caller.rs`:
- Around line 1748-1757: The closure collect_rejected_raw currently clones
a_records and b_records causing double memory; change it to move the raw records
out instead of cloning: when track is true, drain/consume a_records and
b_records and return the combined Vec<Vec<u8>> (and return an empty Vec when
false). Update the closure/call sites so it takes ownership (make it FnOnce or
otherwise consume the vectors) and adjust each rejection return that uses
collect_rejected_raw (including the single-strand paths via ss_caller and the
other listed rejection sites) to pass tracking=true and return the moved Vec
rather than cloned copies; references you’ll need to modify:
collect_rejected_raw, a_records, b_records, and the rejection return points (and
ss_caller usage).
---
Nitpick comments:
In `@src/lib/commands/duplex.rs`:
- Around line 706-722: The mutex in the closure flush_rejects is held for the
entire loop over a group's records (rw_arc lock is acquired before iterating and
calling write_raw_record), which can serialize workers for very large groups; to
reduce contention, acquire the lock only for the minimal write operation (e.g.,
take a short-lived guard per record or move/copy each raw record out and lock
just to call write_raw_record), ensuring you still call write_raw_record on the
same writer (guard.as_mut()) and preserve group contiguity as needed; update the
flush_rejects closure (and uses of rw_arc, guard, and write_raw_record)
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: 886e51f0-e887-4d77-8dbf-48cc78696959
📒 Files selected for processing (8)
crates/fgumi-consensus/src/duplex_caller.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/duplex.rstests/integration/helpers/assertions.rstests/integration/test_bgzf_eof.rstests/integration/test_correct_command.rstests/integration/test_duplex_command.rs
✅ Files skipped from review due to trivial changes (1)
- tests/integration/test_correct_command.rs
🚧 Files skipped from review as they are similar to previous changes (2)
- tests/integration/helpers/assertions.rs
- src/lib/commands/correct.rs
d6a96fe to
fa60481
Compare
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
🧹 Nitpick comments (1)
src/lib/commands/codec.rs (1)
594-606: Mutex held acrosswrite_raw_recordcalls — watch worker contention.
flush_rejectsholds theparking_lot::Mutexfor the full loop overrecs, sowrite_raw_recordback-pressure (BGZF channel full) serializes all workers. Fine for typical reject volumes given the PR's stated goals, but if a high-rejection input regresses throughput vs. the old buffered path, consider a bounded MPSC to a dedicated rejects-writer thread so workers never block each other.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/commands/codec.rs` around lines 594 - 606, flush_rejects currently holds the parking_lot::Mutex (rejects_writer_for_process) across the entire loop calling write_raw_record, which causes workers to block on back-pressure; change this so the mutex is not held while doing I/O: either (A) implement a bounded MPSC channel and spawn a dedicated rejects-writer thread (create a sender stored in rejects_writer_for_process that workers .try_send() / .send() to, and have the writer read from the channel and call write_raw_record) or (B) if you want a minimal change, lock only briefly to take ownership of the inner writer (e.g., swap out the Option<Writer> or clone/move it out of the mutex) and then release the lock before iterating recs and calling write_raw_record; reference flush_rejects, rejects_writer_for_process, and write_raw_record when making the change.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Nitpick comments:
In `@src/lib/commands/codec.rs`:
- Around line 594-606: flush_rejects currently holds the parking_lot::Mutex
(rejects_writer_for_process) across the entire loop calling write_raw_record,
which causes workers to block on back-pressure; change this so the mutex is not
held while doing I/O: either (A) implement a bounded MPSC channel and spawn a
dedicated rejects-writer thread (create a sender stored in
rejects_writer_for_process that workers .try_send() / .send() to, and have the
writer read from the channel and call write_raw_record) or (B) if you want a
minimal change, lock only briefly to take ownership of the inner writer (e.g.,
swap out the Option<Writer> or clone/move it out of the mutex) and then release
the lock before iterating recs and calling write_raw_record; reference
flush_rejects, rejects_writer_for_process, and write_raw_record when making the
change.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 8729633b-29c6-4e41-ab20-24bf6863103c
📒 Files selected for processing (8)
crates/fgumi-consensus/src/duplex_caller.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/duplex.rstests/integration/helpers/assertions.rstests/integration/test_bgzf_eof.rstests/integration/test_correct_command.rstests/integration/test_duplex_command.rs
✅ Files skipped from review due to trivial changes (1)
- tests/integration/test_correct_command.rs
🚧 Files skipped from review as they are similar to previous changes (2)
- tests/integration/helpers/assertions.rs
- src/lib/commands/correct.rs
Rejects are now written to the rejects BAM directly from serialize_fn via a shared BGZF writer guarded by a parking_lot mutex. The mutex only serializes the byte append; BGZF compression runs on the writer's own thread pool. This removes the last unbounded Vec<Vec<u8>> growth path from simplex, complementing the metric-memory fix in #285/#290. Peak reject memory drops from proportional to the reject count to a small fixed working set of the BGZF block buffer.
Same pattern as the simplex change: a shared BGZF writer guarded by a parking_lot mutex receives reject bytes directly from serialize_fn, so rejected raw BAM records no longer accumulate in per-thread Vecs.
Same pattern as the simplex/duplex changes: shared BGZF writer guarded by a parking_lot mutex receives reject bytes directly from serialize_fn.
fa60481 to
8ad2f8f
Compare
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
src/lib/commands/codec.rs (1)
500-502:⚠️ Potential issue | 🟡 MinorStale comment: rejects are no longer buffered.
The "Rejects buffering semantics are preserved (see follow-up)" note contradicts this PR — rejects now stream straight to disk, not via per-thread buffers. Update or drop.
♻️ Suggested tweak
- // Per-thread metrics accumulator: bounded metric memory, no unbounded - // queue. Rejects buffering semantics are preserved (see follow-up). + // Per-thread metrics accumulator: bounded metric memory, no unbounded + // queue. Rejects are streamed directly to disk in process_fn (no buffering). let collected_metrics = PerThreadAccumulator::<CollectedCodecMetrics>::new(num_threads);🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@src/lib/commands/codec.rs` around lines 500 - 502, The inline comment on the PerThreadAccumulator creation is stale: update or remove the phrase "Rejects buffering semantics are preserved (see follow-up)" because rejects now stream directly to disk instead of being buffered per-thread; locate the instantiation of PerThreadAccumulator::<CollectedCodecMetrics> (bound to collected_metrics) in codec.rs and either remove the misleading clause or replace it with a brief accurate note (e.g., "Per-thread metrics accumulator; rejects stream directly to disk, not buffered") so the comment matches current behavior.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@tests/integration/test_correct_command.rs`:
- Around line 179-239: The test currently only creates 63 input records so it
never spans CorrectUmis’ 1000-template batch and asserts any 30 unique rejects;
to fix, increase the uncorrectable family sizes so the total number of templates
exceeds the batch size (e.g. set far_family_size such that far_family_size * 3 +
other templates > 1000), keep expected_rejects = far_family_size * 3, run the
same Command::new("correct") invocation, then in the reader/seen loop assert
that total == expected_rejects and that the set of seen read names exactly
matches the specific uncorrectable reads produced by the three far families (use
the create_umi_family naming pattern for far_a/far_b/far_c to generate the
expected read names) rather than allowing any 30 records.
---
Outside diff comments:
In `@src/lib/commands/codec.rs`:
- Around line 500-502: The inline comment on the PerThreadAccumulator creation
is stale: update or remove the phrase "Rejects buffering semantics are preserved
(see follow-up)" because rejects now stream directly to disk instead of being
buffered per-thread; locate the instantiation of
PerThreadAccumulator::<CollectedCodecMetrics> (bound to collected_metrics) in
codec.rs and either remove the misleading clause or replace it with a brief
accurate note (e.g., "Per-thread metrics accumulator; rejects stream directly to
disk, not buffered") so the comment matches current behavior.
🪄 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: dfd68707-ba19-4706-9d77-7eba06147a9c
📒 Files selected for processing (12)
crates/fgumi-consensus/src/duplex_caller.rscrates/fgumi-sam/src/lib.rsdocs/src/guide/migration-from-fgbio.mdsrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/duplex.rssrc/lib/commands/simplex.rssrc/lib/sam/mod.rstests/integration/helpers/assertions.rstests/integration/test_bgzf_eof.rstests/integration/test_correct_command.rstests/integration/test_duplex_command.rs
✅ Files skipped from review due to trivial changes (1)
- crates/fgumi-consensus/src/duplex_caller.rs
🚧 Files skipped from review as they are similar to previous changes (6)
- src/lib/sam/mod.rs
- tests/integration/helpers/assertions.rs
- docs/src/guide/migration-from-fgbio.md
- tests/integration/test_bgzf_eof.rs
- src/lib/commands/simplex.rs
- tests/integration/test_duplex_command.rs
Rejected raw records are now written to the rejects BAM directly from serialize_fn via a shared BGZF writer guarded by a parking_lot mutex, completing the rejects-streaming work started in simplex/duplex/codec. Removes the last per-thread Vec<Vec<u8>> accumulator in the pipeline commands.
8ad2f8f to
b7dceb0
Compare
Summary
Stacked on #290. Removes the per-thread `Vec<Vec>` reject-buffering pattern from the four commands that keep it (after #290 removed the metric-buffering pattern from all seven). Rejected raw BAM records are now streamed directly to the rejects BAM during `serialize_fn` via a shared BGZF writer guarded by a `parking_lot::Mutex`. The mutex only serializes the byte append; BGZF compression runs on the writer's own thread pool.
Commits (one per command):
Why
With #287/#290 landed, metric memory in the consensus commands and `correct` is capped at O(threads × distinct keys). The remaining unbounded growth vector is reject buffering: every rejected raw BAM record was held in per-thread Vecs for the life of the run, then written after the pipeline completed. On protocols with high rejection rates (e.g. duplex calls that fail `/A`/`/B` partnering, correct runs where most UMIs don't match the whitelist), this can retain multiple GB of raw BAM bytes.
Streaming during the pipeline caps reject memory at roughly the BGZF block size per worker, independent of reject count.
Why not the pipeline's `with_secondary` API
`run_bam_pipeline_from_reader_with_secondary` exists and `filter` uses it, but it writes the secondary output using `output_header`. simplex/duplex/codec build an unmapped-consensus output header that differs from the input header — rejects need the input header since they're original mapped reads. Extending the pipeline API to accept a separate secondary header is a larger surgery; the shared-Mutex pattern here is surgical and keeps the pipeline internals untouched. Worth revisiting once all four commands are in the same shape and a generalization has a clear target.
Test plan
Note
The final changelog commit is unsigned (three 1Password biometric signing attempts failed locally). It should be re-signed before merge.