Skip to content

perf: remove filter single-threaded fallback, fix MemoryEstimate accuracy, and add ordered reject output - #109

Merged
nh13 merged 1 commit into
mainfrom
nhomer/perf-filter-and-memory-estimate-fixes
Feb 17, 2026
Merged

nh13 merged 1 commit into
mainfrom
nhomer/perf-filter-and-memory-estimate-fixes

Conversation

@nh13

@nh13 nh13 commented Feb 16, 2026 •

Copy link
Copy Markdown
Member

Summary

  • Remove filter single-threaded fallback (~200 lines): The execute_single_threaded() path decoded records into RecordBuf via noodles then re-encoded to raw bytes. The unified pipeline with threads=1 is strictly faster and already validated equivalent by test_threading_modes.
  • Fix MemoryEstimate accuracy across 9 types: Multiple implementations used .len() instead of .capacity() and missed Vec element overhead, causing backpressure to activate late and risking memory spikes. Fixed types: SimplexProcessedBatch, DuplexProcessedBatch, CodecProcessedBatch, Vec<RecordBuf>, ClipProcessedBatch, ProcessedDedupGroup, FastqDecompressedBatch, FastqBoundaryBatch, FastqParsedBatch.
  • Add ordered reject output for the filter command's --rejects option. Previously rejected records were written via a Mutex in arbitrary thread completion order. Now rejects are routed through the pipeline's secondary output infrastructure (secondary_data on SerializedBatch/CompressedBlockBatch, secondary_serialize_fn on PipelineFunctions) so both kept and rejected BAM files maintain input order.
  • Remove unused --sort-order CLI arg from filter (was parsed but never applied).
  • Extract FilterProcessCaptures to reduce duplication between single-read and template pipeline modes.
  • Additional correctness fixes: preserve primary error when secondary finalization fails, pre-allocate secondary buffer capacity, reset secondary_data in batch clear() methods, use aggregated failed_reads consistently.

Test plan

  • All 1778 tests pass (cargo ci-test)
  • cargo ci-fmt clean
  • cargo ci-lint clean
  • test_threading_modes validates None, Some(1), Some(2) equivalence
  • New unit tests for all 9 fixed MemoryEstimate implementations
  • Manual verification: filter with --rejects and --threads 4 on 8.6M records, reject BAM maintains input order, fgumi compare bams confirms identical kept output vs baseline

@nh13
nh13 temporarily deployed to github-actions February 16, 2026 03:24 — with GitHub Actions Inactive
@codecov

codecov Bot commented Feb 16, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 92.67887% with 44 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.38%. Comparing base (a5b17d8) to head (7fe3b0b).
⚠️ Report is 1 commits behind head on main.

Files with missing lines Patch % Lines
src/lib/unified_pipeline/bam.rs 80.92% 37 Missing ⚠️
src/lib/unified_pipeline/base.rs 92.98% 4 Missing ⚠️
src/commands/filter.rs 98.36% 3 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main     #109      +/-   ##
==========================================
+ Coverage   82.14%   82.38%   +0.24%     
==========================================
  Files         127      127              
  Lines       51200    51277      +77     
==========================================
+ Hits        42058    42245     +187     
+ Misses       9142     9032     -110     

☔ View full report in Codecov by Sentry.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@coderabbitai

coderabbitai Bot commented Feb 16, 2026 •

Copy link
Copy Markdown

Note

Reviews paused

It 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 reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

This PR expands heap-size estimations to include Vec capacities and container overhead across many processing batches (clip, codec, dedup, duplex, simplex, fastq, and related types) and adds unit tests validating those estimates. The Filter command is refactored into a unified 7-step streaming pipeline. The unified pipeline gains optional dual-output support (secondary_data, secondary serialization plumbing and writer wiring). A new RawBamWriter::write_raw_bytes method was added.

🚥 Pre-merge checks | ✅ 4
✅ Passed checks (4 passed)
Check name Status Explanation
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Merge Conflict Detection ✅ Passed ✅ No merge conflicts detected when merging into main
Title check ✅ Passed Title accurately summarizes the main changes: removing filter single-threaded fallback, fixing MemoryEstimate accuracy, and adding ordered reject output.
Description check ✅ Passed Description clearly details all major changes, test coverage, and specific types fixed, directly related to the changeset.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing touches
  • 📝 Generate docstrings
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Post copyable unit tests in a comment
  • Commit unit tests in branch nhomer/perf-filter-and-memory-estimate-fixes

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands and usage tips.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (2)
src/commands/filter.rs (2)

406-421: ⚠️ Potential issue | 🟠 Major

Rejected records buffered in memory until post-pipeline.

CollectedFilterMetrics.rejects accumulates all rejected raw records in the SegQueue and they're only flushed to disk in aggregate_and_finalize_metrics. For inputs with high rejection rates, this can spike memory substantially — bypassing the pipeline's backpressure mechanism.

Consider streaming rejects to disk within the serialize step, or at minimum documenting this trade-off.


206-214: ⚠️ Potential issue | 🟡 Minor

Missing outer Vec allocation in heap estimate.

The implementation accounts for overhead only for len elements (via .iter()), but kept_records and rejected_records each allocate heap for capacity slots. This undercounts when either vector has spare capacity.

Other MemoryEstimate impls in the codebase explicitly add vec.capacity() * sizeof(Element) to account for this. Follow that pattern:

Proposed fix
 impl MemoryEstimate for FilterProcessedBatchRaw {
     fn estimate_heap_size(&self) -> usize {
         let vec_overhead = std::mem::size_of::<Vec<u8>>();
-        let kept: usize = self.kept_records.iter().map(|v| v.capacity() + vec_overhead).sum();
-        let rejected: usize =
-            self.rejected_records.iter().map(|v| v.capacity() + vec_overhead).sum();
-        kept + rejected
+        let kept_outer = self.kept_records.capacity() * vec_overhead;
+        let kept_inner: usize = self.kept_records.iter().map(|v| v.capacity()).sum();
+        let rejected_outer = self.rejected_records.capacity() * vec_overhead;
+        let rejected_inner: usize = self.rejected_records.iter().map(|v| v.capacity()).sum();
+        kept_outer + kept_inner + rejected_outer + rejected_inner
     }
 }
🧹 Nitpick comments (2)
src/commands/filter.rs (2)

304-310: Significant duplication between single-read and template pipeline setup.

Both methods duplicate: pipeline config, reference loading, metrics collection setup, and the serialize closure. Consider extracting common setup into a shared helper.

Also applies to: 447-453


962-976: Dead code: get_threshold is unused.

Marked #[allow(dead_code)]. If it's no longer needed after the refactor, consider removing it entirely rather than carrying dead code.

@nh13
nh13 temporarily deployed to github-actions February 16, 2026 09:03 — with GitHub Actions Inactive
@nh13
nh13 force-pushed the nhomer/perf-filter-and-memory-estimate-fixes branch from e87144f to 80f40ed Compare February 16, 2026 09:04
@nh13
nh13 temporarily deployed to github-actions February 16, 2026 09:04 — with GitHub Actions Inactive
@nh13
nh13 force-pushed the nhomer/perf-filter-and-memory-estimate-fixes branch from 80f40ed to ed768d4 Compare February 16, 2026 18:28
@nh13
nh13 temporarily deployed to github-actions February 16, 2026 18:28 — with GitHub Actions Inactive

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 (2)
src/lib/unified_pipeline/base.rs (2)

930-933: ⚠️ Potential issue | 🟡 Minor

clear() does not reset secondary_data.

If a CompressedBlockBatch is cleared and reused, stale secondary_data from a previous batch will persist.

Proposed fix
     pub fn clear(&mut self) {
         self.blocks.clear();
         self.record_count = 0;
+        self.secondary_data = None;
     }

1644-1648: ⚠️ Potential issue | 🟡 Minor

Same issue: clear() leaves secondary_data stale.

Proposed fix
     pub fn clear(&mut self) {
         self.data.clear();
         self.record_count = 0;
+        self.secondary_data = None;
     }
🤖 Fix all issues with AI agents
In `@src/commands/filter.rs`:
- Around line 207-217: The estimate_heap_size implementation for
FilterProcessedBatchRaw double-counts the Vec<u8> struct overhead: vec_overhead
is applied once via kept_outer/rejected_outer (capacity * vec_overhead) and then
again inside the iterator sums. In the impl MemoryEstimate for
FilterProcessedBatchRaw, change the kept_inner and rejected_inner computations
(currently mapping |v| v.capacity() + vec_overhead) to only sum v.capacity(),
leaving vec_overhead applied only in kept_outer and rejected_outer; keep the
final sum logic unchanged.

In `@src/lib/unified_pipeline/bam.rs`:
- Around line 3464-3476: finalize_pipeline's error can be lost if secondary
writer.finish() also fails; capture the outcome of finalize_pipeline(&*state)
into result, then attempt to finalize the secondary output
(state.output.secondary_output / writer.finish()) but do not immediately `?` on
its Err — instead, if finalize_pipeline returned Err, return that primary error;
if finalize_pipeline was Ok and writer.finish() failed, return the secondary
error; if both failed and you want to surface both, combine them into a single
io::Error message (e.g., include both errors' messages) and return that,
ensuring you reference the existing symbols result, finalize_pipeline,
state.output.secondary_output and writer.finish().
- Around line 4099-4210: The secondary BAM may be missing the BGZF EOF block
because run_bam_pipeline_from_reader_with_secondary relies on run_bam_pipeline
and the writer.finish() inside it, which may not write the EOF; explicitly
append BGZF_EOF to the secondary_output_path after the pipeline completes
(mirror the primary logic that re-opens output_path and writes BGZF_EOF) —
locate run_bam_pipeline_from_reader_with_secondary and the secondary_writer
created via crate::bam_io::create_raw_bam_writer and, when result.is_ok(), open
secondary_output_path with OpenOptions::new().append(true).open(...) and
write_all(&BGZF_EOF), returning propagated io::Errors analogous to the primary
EOF write.
🧹 Nitpick comments (5)
src/commands/filter.rs (2)

419-443: failed_reads is redundant with failed_reads_total.

failed_reads (line 422/428) accumulates the same value as total_reads - passed_reads (line 436). Consider dropping the separate accumulator and using a single variable.


354-417: run_filter_pipeline — when track_rejects is true, the else branch (no rejects path) silently returns zero reads.

Line 385 checks track_rejects, then line 386 checks self.rejects. Since track_rejects is always derived from self.rejects.is_some() at the call site, the inner if let is technically redundant. If that invariant ever breaks, the function silently writes no output. Consider an else with an unreachable/error, or just remove the redundant guard.

src/lib/unified_pipeline/fastq.rs (1)

2172-2178: Heap tracking stub persists under memory-debug feature.

q4_heap is hardcoded to 0, so groups_heap_bytes never reflects actual memory. The comment acknowledges this, but it means the memory-debug feature gives misleading data for this queue. Consider a TODO/issue reference so it doesn't stay forgotten.

src/lib/unified_pipeline/bam.rs (2)

2744-2749: std::mem::take discards the secondary buffer's allocation every batch.

Unlike the primary buffer (line 2739) which uses std::mem::replace with a pre-allocated Vec, the secondary buffer loses its capacity after each non-empty batch, causing repeated allocations.

♻️ Retain capacity like the primary buffer
     // Build secondary data if any was serialized
     let secondary_data = if worker.core.secondary_serialization_buffer.is_empty() {
         None
     } else {
-        Some(std::mem::take(&mut worker.core.secondary_serialization_buffer))
+        Some(std::mem::replace(
+            &mut worker.core.secondary_serialization_buffer,
+            Vec::with_capacity(SERIALIZATION_BUFFER_CAPACITY),
+        ))
     };

3119-3141: Secondary serialize+write pattern is repeated 3 times in the single-threaded path.

Consider extracting a helper closure to reduce duplication. Low priority — the logic is straightforward and each site is short.

Also applies to: 3156-3172, 3184-3206

Comment thread src/commands/filter.rs
Comment thread src/lib/unified_pipeline/bam.rs
Comment thread src/lib/unified_pipeline/bam.rs
@nh13
nh13 temporarily deployed to github-actions February 17, 2026 02:25 — with GitHub Actions Inactive

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

🧹 Nitpick comments (6)
src/lib/unified_pipeline/bam.rs (3)

2712-2749: secondary_serialization_buffer loses capacity after std::mem::take.

Line 2748 uses std::mem::take, which replaces the buffer with a zero-capacity Vec. The primary buffer (line 2739-2742) is replaced with a pre-allocated Vec::with_capacity(SERIALIZATION_BUFFER_CAPACITY). The secondary buffer will re-allocate from scratch on every batch that produces secondary data.

♻️ Restore capacity for the secondary buffer
-        Some(std::mem::take(&mut worker.core.secondary_serialization_buffer))
+        Some(std::mem::replace(
+            &mut worker.core.secondary_serialization_buffer,
+            Vec::with_capacity(SERIALIZATION_BUFFER_CAPACITY),
+        ))
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/unified_pipeline/bam.rs` around lines 2712 - 2749, The
secondary_serialization_buffer is being drained with std::mem::take which leaves
a zero-capacity Vec and causes reallocation; change the logic that builds
secondary_data to use std::mem::replace and pre-allocate capacity like the
primary buffer (i.e. mirror the pattern used for
worker.core.serialization_buffer), e.g. replace the std::mem::take call with
std::mem::replace(&mut worker.core.secondary_serialization_buffer,
Vec::with_capacity(SERIALIZATION_BUFFER_CAPACITY)) so secondary_data is
Some(replaced_vec) and the buffer keeps reserved capacity for future batches
(refer to worker.core.secondary_serialization_buffer,
SERIALIZATION_BUFFER_CAPACITY, and the combined_data replacement pattern).

2833-2846: Nested lock acquisition: output lock → secondary_mutex.

The Write step holds the output lock (line 2795) when it acquires secondary_mutex (line 2837). The finalization path (line 3467) acquires secondary_mutex without the output lock. This is safe (no circular dependency), but worth documenting the lock ordering to prevent future issues.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/unified_pipeline/bam.rs` around lines 2833 - 2846, Add an explicit
comment documenting the lock ordering around state.output and secondary_output
to prevent future deadlocks: note that the Write step (the block that reads
batch.secondary_data and locks state.output.secondary_output via secondary_mutex
and calls sw.write_raw_bytes) acquires the output lock first and then
secondary_mutex (output lock → secondary_mutex), while the finalization path
locks secondary_mutex without holding the output lock; state.set_error and
sw.write_raw_bytes are called under secondary_mutex. Insert this comment
adjacent to the Write block and the finalization code paths so maintainers see
the intended ordering and the rationale that no circular dependency exists.

3119-3141: Repeated secondary serialize+write block appears three times.

The pattern (clear secondary buffer → call secondary_fn → write to secondary_writer) is duplicated across three locations in the single-threaded path. Consider extracting a helper closure or function.

Also applies to: 3156-3172, 3184-3206

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/unified_pipeline/bam.rs` around lines 3119 - 3141, The repeated
pattern clearing buffers.secondary, invoking fns.secondary_serialize_fn, and
writing bytes to secondary_writer should be extracted into a small helper
(closure or private fn) to avoid duplication; implement something like a helper
named flush_secondary that takes &processed, &mut buffers.secondary, a reference
to fns.secondary_serialize_fn, and secondary_writer.as_mut(), does
buffers.secondary.clear(), calls the secondary_fn when Some, and if
buffers.secondary is not empty writes via sw.write_raw_bytes(&buffers.secondary)
returning the same Result type, then replace the three duplicated blocks with a
single call to this helper (use the existing symbols buffers.secondary,
fns.secondary_serialize_fn, secondary_writer and sw.write_raw_bytes to locate
and wire the helper).
src/commands/filter.rs (3)

384-404: Redundant guard: track_rejects already implies self.rejects.is_some().

Line 288 sets track_rejects = self.rejects.is_some(), so the inner if let Some(rejects_path) on line 385 can never be None when track_rejects is true. Consider flattening to a single if let Some(rejects_path) = &self.rejects.

Simplify
-        if track_rejects {
-            if let Some(rejects_path) = &self.rejects {
-                // Secondary serialize: write rejected records
-                let secondary_serialize_fn =
-                    |batch: &FilterProcessedBatchRaw, buf: &mut Vec<u8>| -> io::Result<u64> {
-                        serialize_raw_records(&batch.rejected_records, buf)
-                    };
-
-                run_bam_pipeline_from_reader_with_secondary(
-                    setup.pipeline_config,
-                    reader,
-                    header,
-                    &self.io.output,
-                    None,
-                    rejects_path,
-                    grouper_fn,
-                    process_fn,
-                    serialize_fn,
-                    secondary_serialize_fn,
-                )?;
-            }
-        } else {
+        if let Some(rejects_path) = &self.rejects {
+            let secondary_serialize_fn =
+                |batch: &FilterProcessedBatchRaw, buf: &mut Vec<u8>| -> io::Result<u64> {
+                    serialize_raw_records(&batch.rejected_records, buf)
+                };
+
+            run_bam_pipeline_from_reader_with_secondary(
+                setup.pipeline_config,
+                reader,
+                header,
+                &self.io.output,
+                None,
+                rejects_path,
+                grouper_fn,
+                process_fn,
+                serialize_fn,
+                secondary_serialize_fn,
+            )?;
+        } else {
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 384 - 404, Remove the redundant boolean
guard: instead of checking `if track_rejects { if let Some(rejects_path) =
&self.rejects { ... } }`, collapse to a single `if let Some(rejects_path) =
&self.rejects { ... }` because `track_rejects` is derived from
`self.rejects.is_some()`; keep the inner block unchanged (secondary_serialize_fn
and the call to run_bam_pipeline_from_reader_with_secondary) and drop all
references to `track_rejects` in this branch to avoid the redundant condition.

310-346: Reference is unconditionally loaded even though the field is Option<Arc<ReferenceReader>>.

Line 332-333 always wraps in Some(...). If the intent is to support optional references in the future, this is fine scaffolding. But currently the Option wrapper adds unnecessary unwrapping downstream without benefit — reference on the CLI is required (#[arg] without Option).

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 310 - 346, The code always wraps the
reference in Some(Arc::new(ReferenceReader::new(...))) even though the CLI
requires a reference, so remove the unnecessary Option wrapper: in
setup_pipeline(), replace the local variable reference:
Option<Arc<ReferenceReader>> = Some(…) with reference: Arc<ReferenceReader> =
Arc::new(ReferenceReader::new(&self.reference)?), remove the Some(...) wrapper,
and update the FilterPipelineSetup type (and any other places) to expect
Arc<ReferenceReader> instead of Option<Arc<ReferenceReader>> so downstream
unwrapping is no longer required; keep the ReferenceReader::new call and the
info logging as-is.

418-442: failed_reads aggregated but unused for logging; failed_reads_total recomputed instead.

failed_reads (line 421/427) is only passed to write_filter_stats, while line 435 recomputes total_reads - passed_reads for the info log. If these ever diverge it'll be confusing. Pick one source of truth.

Use aggregated value
-        let failed_reads_total = total_reads - passed_reads;
-        info!(
-            "Processed {total_reads} reads; kept {passed_reads} and rejected {failed_reads_total}"
-        );
-        if track_rejects && failed_reads_total > 0 {
-            info!("Wrote {failed_reads_total} rejected records to rejects file");
+        info!(
+            "Processed {total_reads} reads; kept {passed_reads} and rejected {failed_reads}"
+        );
+        if track_rejects && failed_reads > 0 {
+            info!("Wrote {failed_reads} rejected records to rejects file");
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 418 - 442, The code aggregates
failed_reads but then recomputes failed_reads_total from total_reads and
passed_reads, which can diverge; use the aggregated failed_reads as the single
source of truth. Replace the computed failed_reads_total = total_reads -
passed_reads usage in the info! logs and the track_rejects check with the
aggregated failed_reads variable (and pass failed_reads into write_filter_stats
as you already do), ensuring all messaging (Processed/kept/rejected, rejects
file message) and the "Wrote ..." condition reference failed_reads consistently;
update any variable names if you keep failed_reads_total for clarity so it is
assigned from failed_reads instead of recomputing.
🤖 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/filter.rs`:
- Around line 282-296: The CLI option sort_order is parsed into the field
sort_order but never used; propagate it from setup_pipeline/run_filter_pipeline
into the pipeline and apply it when opening/constructing the BAM reader or
output header. Concretely, thread the sort_order Option<String> into
create_bam_reader_for_pipeline or into the code that mutates the SAM/BAM header
(e.g., after add_pg_record) and, when Some("coordinate") (or other valid
values), modify the header's SO (sort order) tag or configure the reader/writer
to output in that order so --sort-order takes effect; update calls to
create_bam_reader_for_pipeline, execute_threads_mode_template, and
execute_threads_mode_single_read to accept and forward sort_order so the chosen
sort order is honored.

---

Duplicate comments:
In `@src/lib/unified_pipeline/bam.rs`:
- Around line 4105-4216: The secondary BAM may be missing the BGZF EOF block
because you only append BGZF_EOF for the primary; after run_bam_pipeline
succeeds you should also finalize the secondary output similarly. In
run_bam_pipeline_from_reader_with_secondary, after checking result.is_ok(), open
secondary_output_path with append and write_all(&BGZF_EOF) (handling/map_err
like the primary) or call/ensure the secondary writer's finish path writes the
EOF; reference the secondary_writer passed into run_bam_pipeline and the
BGZF_EOF constant to implement the symmetric EOF append for the secondary
output.

---

Nitpick comments:
In `@src/commands/filter.rs`:
- Around line 384-404: Remove the redundant boolean guard: instead of checking
`if track_rejects { if let Some(rejects_path) = &self.rejects { ... } }`,
collapse to a single `if let Some(rejects_path) = &self.rejects { ... }` because
`track_rejects` is derived from `self.rejects.is_some()`; keep the inner block
unchanged (secondary_serialize_fn and the call to
run_bam_pipeline_from_reader_with_secondary) and drop all references to
`track_rejects` in this branch to avoid the redundant condition.
- Around line 310-346: The code always wraps the reference in
Some(Arc::new(ReferenceReader::new(...))) even though the CLI requires a
reference, so remove the unnecessary Option wrapper: in setup_pipeline(),
replace the local variable reference: Option<Arc<ReferenceReader>> = Some(…)
with reference: Arc<ReferenceReader> =
Arc::new(ReferenceReader::new(&self.reference)?), remove the Some(...) wrapper,
and update the FilterPipelineSetup type (and any other places) to expect
Arc<ReferenceReader> instead of Option<Arc<ReferenceReader>> so downstream
unwrapping is no longer required; keep the ReferenceReader::new call and the
info logging as-is.
- Around line 418-442: The code aggregates failed_reads but then recomputes
failed_reads_total from total_reads and passed_reads, which can diverge; use the
aggregated failed_reads as the single source of truth. Replace the computed
failed_reads_total = total_reads - passed_reads usage in the info! logs and the
track_rejects check with the aggregated failed_reads variable (and pass
failed_reads into write_filter_stats as you already do), ensuring all messaging
(Processed/kept/rejected, rejects file message) and the "Wrote ..." condition
reference failed_reads consistently; update any variable names if you keep
failed_reads_total for clarity so it is assigned from failed_reads instead of
recomputing.

In `@src/lib/unified_pipeline/bam.rs`:
- Around line 2712-2749: The secondary_serialization_buffer is being drained
with std::mem::take which leaves a zero-capacity Vec and causes reallocation;
change the logic that builds secondary_data to use std::mem::replace and
pre-allocate capacity like the primary buffer (i.e. mirror the pattern used for
worker.core.serialization_buffer), e.g. replace the std::mem::take call with
std::mem::replace(&mut worker.core.secondary_serialization_buffer,
Vec::with_capacity(SERIALIZATION_BUFFER_CAPACITY)) so secondary_data is
Some(replaced_vec) and the buffer keeps reserved capacity for future batches
(refer to worker.core.secondary_serialization_buffer,
SERIALIZATION_BUFFER_CAPACITY, and the combined_data replacement pattern).
- Around line 2833-2846: Add an explicit comment documenting the lock ordering
around state.output and secondary_output to prevent future deadlocks: note that
the Write step (the block that reads batch.secondary_data and locks
state.output.secondary_output via secondary_mutex and calls sw.write_raw_bytes)
acquires the output lock first and then secondary_mutex (output lock →
secondary_mutex), while the finalization path locks secondary_mutex without
holding the output lock; state.set_error and sw.write_raw_bytes are called under
secondary_mutex. Insert this comment adjacent to the Write block and the
finalization code paths so maintainers see the intended ordering and the
rationale that no circular dependency exists.
- Around line 3119-3141: The repeated pattern clearing buffers.secondary,
invoking fns.secondary_serialize_fn, and writing bytes to secondary_writer
should be extracted into a small helper (closure or private fn) to avoid
duplication; implement something like a helper named flush_secondary that takes
&processed, &mut buffers.secondary, a reference to fns.secondary_serialize_fn,
and secondary_writer.as_mut(), does buffers.secondary.clear(), calls the
secondary_fn when Some, and if buffers.secondary is not empty writes via
sw.write_raw_bytes(&buffers.secondary) returning the same Result type, then
replace the three duplicated blocks with a single call to this helper (use the
existing symbols buffers.secondary, fns.secondary_serialize_fn, secondary_writer
and sw.write_raw_bytes to locate and wire the helper).

Comment thread src/commands/filter.rs
@nh13
nh13 force-pushed the nhomer/perf-filter-and-memory-estimate-fixes branch from b6324d8 to 5b6488c Compare February 17, 2026 04:38
@nh13
nh13 temporarily deployed to github-actions February 17, 2026 04:38 — with GitHub Actions Inactive

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick comments (4)
src/commands/filter.rs (3)

496-533: Per-record Vec::new() allocations.

Each single record creates two fresh Vecs (kept_records, rejected_records). For single-read mode this means one allocation pair per record. The pipeline likely amortizes this, but Vec::with_capacity(1) / Vec::new() for rejected when not tracking would avoid a realloc on push.

Minor, not blocking.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 496 - 533, The closure process_fn
currently allocates kept_records and rejected_records with Vec::new() for every
record; change it to pre-size the vectors to avoid per-record reallocs by using
Vec::with_capacity(1) for kept_records and only creating rejected_records
(Vec::with_capacity(1)) when track_rejects is true (or reserve(1) on first
push), then populate and return FilterProcessedBatchRaw as before; reference
process_fn, kept_records, rejected_records, track_rejects, and
FilterProcessedBatchRaw to locate where to apply the change.

416-436: Redundant guard — track_rejects already implies self.rejects.is_some().

track_rejects is set from self.rejects.is_some() at line 301, so the inner if let Some(rejects_path) on line 417 is always true here. If track_rejects is true but somehow self.rejects is None, the pipeline silently falls through without writing kept or rejected records.

Consider collapsing or adding an else-bail:

Suggested simplification
-        if track_rejects {
-            if let Some(rejects_path) = &self.rejects {
+        if let Some(rejects_path) = &self.rejects {
+            {
                 // Secondary serialize: write rejected records
                 ...
-            }
-        } else {
+            }
+        } else {
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 416 - 436, The inner `if let
Some(rejects_path)` is redundant because `track_rejects` is derived from
`self.rejects.is_some()`; replace the inner optional match by unwrapping
`self.rejects` when `track_rejects` is true (e.g., let rejects_path =
self.rejects.as_ref().expect("track_rejects true implies rejects present")) and
then call run_bam_pipeline_from_reader_with_secondary with that rejects_path and
the existing secondary_serialize_fn (which uses FilterProcessedBatchRaw and
serialize_raw_records); alternatively, if you prefer explicit failure, change
the branch to bail with a clear error when `self.rejects` is None instead of
silently skipping writing rejects.

967-1996: Consider using create_filter_with_paths (or a builder) to reduce test boilerplate.

~20 tests construct Filter inline with near-identical fields. Most only vary 1–2 fields. create_filter_with_paths already exists (line 908) but is barely used. Many of these tests only assert that struct fields store what was passed in — consider whether they add value.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 967 - 1996, Many tests repeatedly
construct nearly identical Filter structs; replace repeated inline constructions
in tests with the existing helper create_filter_with_paths (or add a builder) to
reduce boilerplate and make intent clearer: update tests that create Filter
directly (search for Filter { io: BamIoOptions { input:, output: }, reference:,
... } blocks) to call create_filter_with_paths(...) or a new builder that
accepts overrides for the few varying fields (e.g., min_reads, threading,
rejects, stats, max_no_call_fraction), and collapse or remove trivial
“store-and-assert” tests that only confirm field assignment in favor of
asserting behavior in validation/masking functions (e.g., validate_parameters,
mask_bases, filter_read).
src/lib/unified_pipeline/bam.rs (1)

3119-3143: Extract the repeated secondary-serialize-then-write block into a helper.

The secondary serialize + write pattern is duplicated three times in run_bam_pipeline_single_threaded (main loop, final_batch loop, final_group). A small closure or helper would reduce the surface area for future divergence.

♻️ Sketch
+    // Helper: secondary serialize + write
+    let mut do_secondary = |processed: &P, buffers: &mut SingleThreadedBuffers, sw: &mut Option<crate::bam_io::RawBamWriter>| -> io::Result<()> {
+        buffers.secondary.clear();
+        if let Some(ref sec_fn) = fns.secondary_serialize_fn {
+            (sec_fn)(processed, &mut buffers.secondary)?;
+        }
+        if !buffers.secondary.is_empty() {
+            if let Some(ref mut w) = sw {
+                w.write_raw_bytes(&buffers.secondary)?;
+            }
+        }
+        Ok(())
+    };

Then replace each occurrence with do_secondary(&processed, &mut buffers, &mut secondary_writer)?;.

Also applies to: 3156-3173, 3184-3206

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/unified_pipeline/bam.rs` around lines 3119 - 3143, The repeated
pattern that runs the secondary serialize into buffers.secondary and then writes
it (when non-empty) appears multiple times in run_bam_pipeline_single_threaded;
extract that logic into a small helper (e.g., do_secondary) that takes
(&processed, &mut buffers, &mut secondary_writer) and performs: clear
buffers.secondary, call fns.secondary_serialize_fn if Some to fill
buffers.secondary, and if buffers.secondary is non-empty and secondary_writer is
Some call sw.write_raw_bytes(&buffers.secondary)?; then replace each duplicated
block (the occurrences around buffers.secondary.clear(), secondary_serialize_fn,
and conditional sw.write_raw_bytes) with a call to the new helper to keep
behavior identical and reduce duplication.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Duplicate comments:
In `@src/commands/filter.rs`:
- Around line 152-154: The CLI option sort_order (struct field sort_order) is
parsed but not used when producing the output BAM; update the pipeline code that
constructs/writes the output BAM (e.g., the function that assembles the writer
or invokes the sort step — locate functions like run(), execute(),
build_pipeline(), or write_bam()) to read self.sort_order (or the passed
sort_order) and apply it to the writer/sort stage (map Option<String> to the
corresponding enum/parameter and pass it into the sorter/writer initialization).
Ensure the submitted sort_order value is validated/mapped to the expected sort
enum and forwarded wherever the pipeline configures output sorting so the CLI
flag actually changes output ordering.

In `@src/lib/unified_pipeline/bam.rs`:
- Around line 4200-4216: The BGZF EOF is only appended to output_path but not to
secondary_output_path, risking truncated secondary BAMs; either ensure
RawBamWriter::finish() (called in run_bam_pipeline / writer.finish()) actually
writes the 28-byte BGZF_EOF or mirror the primary logic by appending BGZF_EOF to
secondary_output_path when result.is_ok(); specifically add the same
OpenOptions::new().append(true).open(secondary_output_path) and
write_all(&BGZF_EOF) with equivalent map_err messages, or alternatively update
RawBamWriter::finish()/BgzfWriterEnum implementation to explicitly write
BGZF_EOF so downstream tools (e.g., samtools) see a complete BGZF EOF marker.

---

Nitpick comments:
In `@src/commands/filter.rs`:
- Around line 496-533: The closure process_fn currently allocates kept_records
and rejected_records with Vec::new() for every record; change it to pre-size the
vectors to avoid per-record reallocs by using Vec::with_capacity(1) for
kept_records and only creating rejected_records (Vec::with_capacity(1)) when
track_rejects is true (or reserve(1) on first push), then populate and return
FilterProcessedBatchRaw as before; reference process_fn, kept_records,
rejected_records, track_rejects, and FilterProcessedBatchRaw to locate where to
apply the change.
- Around line 416-436: The inner `if let Some(rejects_path)` is redundant
because `track_rejects` is derived from `self.rejects.is_some()`; replace the
inner optional match by unwrapping `self.rejects` when `track_rejects` is true
(e.g., let rejects_path = self.rejects.as_ref().expect("track_rejects true
implies rejects present")) and then call
run_bam_pipeline_from_reader_with_secondary with that rejects_path and the
existing secondary_serialize_fn (which uses FilterProcessedBatchRaw and
serialize_raw_records); alternatively, if you prefer explicit failure, change
the branch to bail with a clear error when `self.rejects` is None instead of
silently skipping writing rejects.
- Around line 967-1996: Many tests repeatedly construct nearly identical Filter
structs; replace repeated inline constructions in tests with the existing helper
create_filter_with_paths (or add a builder) to reduce boilerplate and make
intent clearer: update tests that create Filter directly (search for Filter {
io: BamIoOptions { input:, output: }, reference:, ... } blocks) to call
create_filter_with_paths(...) or a new builder that accepts overrides for the
few varying fields (e.g., min_reads, threading, rejects, stats,
max_no_call_fraction), and collapse or remove trivial “store-and-assert” tests
that only confirm field assignment in favor of asserting behavior in
validation/masking functions (e.g., validate_parameters, mask_bases,
filter_read).

In `@src/lib/unified_pipeline/bam.rs`:
- Around line 3119-3143: The repeated pattern that runs the secondary serialize
into buffers.secondary and then writes it (when non-empty) appears multiple
times in run_bam_pipeline_single_threaded; extract that logic into a small
helper (e.g., do_secondary) that takes (&processed, &mut buffers, &mut
secondary_writer) and performs: clear buffers.secondary, call
fns.secondary_serialize_fn if Some to fill buffers.secondary, and if
buffers.secondary is non-empty and secondary_writer is Some call
sw.write_raw_bytes(&buffers.secondary)?; then replace each duplicated block (the
occurrences around buffers.secondary.clear(), secondary_serialize_fn, and
conditional sw.write_raw_bytes) with a call to the new helper to keep behavior
identical and reduce duplication.

@nh13
nh13 force-pushed the nhomer/perf-filter-and-memory-estimate-fixes branch from 5b6488c to 7420717 Compare February 17, 2026 04:47
@nh13
nh13 temporarily deployed to github-actions February 17, 2026 04:47 — with GitHub Actions Inactive

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
src/commands/filter.rs (1)

93-95: ⚠️ Potential issue | 🟡 Minor

Stale --sort-order reference in help text.

The long_about still mentions --sort-order (line 94), but that option was removed from the struct. Users reading --help will see a nonexistent flag.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 93 - 95, The help text in the command's
long_about still references a removed flag (--sort-order); update the long_about
string in src/commands/filter.rs (the constant/assignment used to define
long_about for the Filter command) to remove or reword the sentence mentioning
--sort-order so it no longer points to a nonexistent option (alternatively, if
you intended to keep the flag, reintroduce the corresponding struct field and
clap attribute in the FilterArgs/struct that defines the options). Ensure the
change touches the long_about definition associated with the filter command and
that the help text now accurately reflects available flags.
🧹 Nitpick comments (5)
src/commands/filter.rs (2)

486-523: Per-record Vec allocation in single-read mode.

Each record creates two fresh Vecs (lines 487-488). Since single-read mode processes one record at a time, consider Vec::with_capacity(1) for the expected-pass path to avoid a realloc on push. Minor, though.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 486 - 523, The per-record closure
process_fn currently allocates empty Vecs for kept_records and rejected_records;
to avoid a realloc in the common single-read pass path, initialize kept_records
with Vec::with_capacity(1) (and optionally rejected_records with
Vec::with_capacity(1) when track_rejects is true) before pushing the record,
ensuring FilterProcessedBatchRaw is constructed the same way.

341-343: reference is always Some here.

Since validate_file_exists already bails on missing reference (line 269), this is always Some(Arc::new(...)). The Option wrapper is cosmetic overhead — harmless but slightly misleading.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/commands/filter.rs` around lines 341 - 343, The local binding reference
is wrapped in Option but is always Some because validate_file_exists ensures it
exists; change the type from Option<Arc<ReferenceReader>> to
Arc<ReferenceReader> and construct it directly with
Arc::new(ReferenceReader::new(&self.reference)?); update any uses of reference
(and function signatures or structs that currently expect
Option<Arc<ReferenceReader>>) to accept Arc<ReferenceReader> instead (search for
the symbol reference and the type Option<Arc<ReferenceReader>>), removing the
needless Some(...) wrapper and any matches/unwraps that handled None.
src/lib/unified_pipeline/bam.rs (2)

3118-3146: Repeated secondary-serialize-then-write block appears three times.

The pattern (clear → secondary serialize → primary serialize → compress → write → write secondary) is duplicated across the main loop, the final-batch loop, and the final-group handler. Consider extracting a helper like process_and_write_group(...).

♻️ Sketch
// Define once, e.g.:
fn process_group(
    group: G,
    fns: &PipelineFunctions<G, P>,
    buffers: &mut SingleThreadedBuffers,
    compressor: &mut InlineBgzfCompressor,
    output: &mut dyn Write,
    secondary_writer: &mut Option<crate::bam_io::RawBamWriter>,
    progress: &ProgressTracker,
) -> io::Result<()> {
    let processed = (fns.process_fn)(group)?;
    buffers.secondary.clear();
    if let Some(ref sf) = fns.secondary_serialize_fn {
        (sf)(&processed, &mut buffers.secondary)?;
    }
    buffers.serialized.clear();
    let rc = (fns.serialize_fn)(processed, &mut buffers.serialized)?;
    compressor.write_all(&buffers.serialized)?;
    compressor.maybe_compress()?;
    compressor.write_blocks_to(output)?;
    if !buffers.secondary.is_empty() {
        if let Some(ref mut sw) = secondary_writer {
            sw.write_raw_bytes(&buffers.secondary)?;
        }
    }
    progress.log_if_needed(rc);
    Ok(())
}

Also applies to: 3157-3177, 3183-3212

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/unified_pipeline/bam.rs` around lines 3118 - 3146, Extract the
repeated "clear → secondary serialize → primary serialize → compress → write →
write secondary → log" sequence into a single helper (e.g.
process_and_write_group or process_group) that takes the group, fns (used for
process_fn, secondary_serialize_fn, serialize_fn), buffers (buffers.secondary,
buffers.serialized), compressor, output, secondary_writer and progress; inside
the helper call (fns.process_fn)(...), clear and run the optional
(fns.secondary_serialize_fn) into buffers.secondary, run (fns.serialize_fn) into
buffers.serialized to get record_count, then use
compressor.write_all/maybe_compress/write_blocks_to(output) and write_raw_bytes
to secondary_writer if buffers.secondary not empty, finally call
progress.log_if_needed(record_count) and return io::Result; replace the three
duplicated blocks in the main loop, final-batch loop and final-group handler
with calls to this helper to remove duplication.

2836-2849: Silent no-op if secondary_output mutex holds None.

If the secondary writer was already taken (e.g., during finalization or error cleanup), *sw_guard is None and secondary data is silently dropped. This is probably intentional during teardown, but worth noting that secondary records can be lost without error if the writer is prematurely consumed.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/unified_pipeline/bam.rs` around lines 2836 - 2849, The code silently
drops secondary_data when the secondary_output mutex contains None; update the
branch handling *sw_guard == None to surface this condition instead of no-op:
detect when state.output.secondary_output.lock() yields Some(None) and either
call state.set_error(...) with a descriptive error or return a failure tuple
(false, false) so callers know data was lost; reference the
batch.secondary_data, state.output.secondary_output, sw_guard/*sw*/,
sw.write_raw_bytes(...) and state.set_error(...) symbols when making the change.
src/lib/unified_pipeline/base.rs (1)

1645-1649: SerializedBatch::clear() drops secondary buffer entirely — verify this is preferred over .clear().

Setting secondary_data = None deallocates the buffer. If batches are reused across iterations, this forces reallocation each time secondary data is present. If reuse is rare or secondary data size varies widely, None is fine (and matches the PR objective of avoiding retained memory). Just confirming this is intentional.

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@src/lib/unified_pipeline/base.rs` around lines 1645 - 1649, The clear()
implementation for SerializedBatch currently sets secondary_data = None which
drops/deallocates the buffer; if the goal is to retain the allocation for reuse
change this to call secondary_data.as_mut().map(|v| v.clear()) (or replace None
with Some(Vec::new()) only when needed) so the Vec's capacity is preserved
across iterations; if the PR intentionally wants to free memory keep
secondary_data = None but add a comment in SerializedBatch::clear() clarifying
the choice so reviewers know it is intentional.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Outside diff comments:
In `@src/commands/filter.rs`:
- Around line 93-95: The help text in the command's long_about still references
a removed flag (--sort-order); update the long_about string in
src/commands/filter.rs (the constant/assignment used to define long_about for
the Filter command) to remove or reword the sentence mentioning --sort-order so
it no longer points to a nonexistent option (alternatively, if you intended to
keep the flag, reintroduce the corresponding struct field and clap attribute in
the FilterArgs/struct that defines the options). Ensure the change touches the
long_about definition associated with the filter command and that the help text
now accurately reflects available flags.

---

Duplicate comments:
In `@src/commands/filter.rs`:
- Around line 203-212: The previous MemoryEstimate implementation double-counted
Vec overhead; update the impl for FilterProcessedBatchRaw so estimate_heap_size
computes the outer overhead as self.kept_records.capacity() *
size_of::<Vec<u8>>() and self.rejected_records.capacity() *
size_of::<Vec<u8>>(), and compute inner heap usage by summing each inner
Vec<u8>::capacity() (kept_inner and rejected_inner), then return the sum
(kept_outer + kept_inner + rejected_outer + rejected_inner) in the
estimate_heap_size method to avoid double-counting; locate and update the impl
MemoryEstimate for FilterProcessedBatchRaw and the estimate_heap_size function
accordingly.

In `@src/lib/unified_pipeline/bam.rs`:
- Around line 4108-4219: The secondary BAM may miss the 28-byte BGZF EOF block
because finish() used inside run_bam_pipeline may not write it; in
run_bam_pipeline_from_reader_with_secondary after calling run_bam_pipeline (and
where you append BGZF_EOF to the primary output on success), also open
secondary_output_path with append and write_all(&BGZF_EOF) on success (mirroring
the primary logic), using the same error mapping style (io::Error::new(e.kind(),
format!(...))) to report failures; reference
secondary_writer/secondary_output_path,
run_bam_pipeline_from_reader_with_secondary, run_bam_pipeline, and BGZF_EOF when
making the change.

---

Nitpick comments:
In `@src/commands/filter.rs`:
- Around line 486-523: The per-record closure process_fn currently allocates
empty Vecs for kept_records and rejected_records; to avoid a realloc in the
common single-read pass path, initialize kept_records with Vec::with_capacity(1)
(and optionally rejected_records with Vec::with_capacity(1) when track_rejects
is true) before pushing the record, ensuring FilterProcessedBatchRaw is
constructed the same way.
- Around line 341-343: The local binding reference is wrapped in Option but is
always Some because validate_file_exists ensures it exists; change the type from
Option<Arc<ReferenceReader>> to Arc<ReferenceReader> and construct it directly
with Arc::new(ReferenceReader::new(&self.reference)?); update any uses of
reference (and function signatures or structs that currently expect
Option<Arc<ReferenceReader>>) to accept Arc<ReferenceReader> instead (search for
the symbol reference and the type Option<Arc<ReferenceReader>>), removing the
needless Some(...) wrapper and any matches/unwraps that handled None.

In `@src/lib/unified_pipeline/bam.rs`:
- Around line 3118-3146: Extract the repeated "clear → secondary serialize →
primary serialize → compress → write → write secondary → log" sequence into a
single helper (e.g. process_and_write_group or process_group) that takes the
group, fns (used for process_fn, secondary_serialize_fn, serialize_fn), buffers
(buffers.secondary, buffers.serialized), compressor, output, secondary_writer
and progress; inside the helper call (fns.process_fn)(...), clear and run the
optional (fns.secondary_serialize_fn) into buffers.secondary, run
(fns.serialize_fn) into buffers.serialized to get record_count, then use
compressor.write_all/maybe_compress/write_blocks_to(output) and write_raw_bytes
to secondary_writer if buffers.secondary not empty, finally call
progress.log_if_needed(record_count) and return io::Result; replace the three
duplicated blocks in the main loop, final-batch loop and final-group handler
with calls to this helper to remove duplication.
- Around line 2836-2849: The code silently drops secondary_data when the
secondary_output mutex contains None; update the branch handling *sw_guard ==
None to surface this condition instead of no-op: detect when
state.output.secondary_output.lock() yields Some(None) and either call
state.set_error(...) with a descriptive error or return a failure tuple (false,
false) so callers know data was lost; reference the batch.secondary_data,
state.output.secondary_output, sw_guard/*sw*/, sw.write_raw_bytes(...) and
state.set_error(...) symbols when making the change.

In `@src/lib/unified_pipeline/base.rs`:
- Around line 1645-1649: The clear() implementation for SerializedBatch
currently sets secondary_data = None which drops/deallocates the buffer; if the
goal is to retain the allocation for reuse change this to call
secondary_data.as_mut().map(|v| v.clear()) (or replace None with
Some(Vec::new()) only when needed) so the Vec's capacity is preserved across
iterations; if the PR intentionally wants to free memory keep secondary_data =
None but add a comment in SerializedBatch::clear() clarifying the choice so
reviewers know it is intentional.

…racy, and add ordered reject output

Remove the single-threaded fallback in filter (~200 lines) that decoded
records into RecordBuf via noodles then re-encoded to raw bytes. The
unified pipeline with threads=1 is strictly faster and already validated
by test_threading_modes.

Fix MemoryEstimate implementations across 9 types that used .len()
instead of .capacity() and missed Vec element overhead, causing
backpressure to activate late and risking memory spikes.

Add ordered reject output for the filter command's --rejects option.
Previously rejected records were written via a Mutex in arbitrary thread
completion order. Now rejects are routed through the pipeline's secondary
output infrastructure (secondary_data on SerializedBatch and
CompressedBlockBatch, secondary_serialize_fn on PipelineFunctions) so
both kept and rejected BAM files maintain input order.

Additional improvements:
- Remove unused --sort-order CLI arg from filter (was parsed but never applied)
- Extract FilterProcessCaptures to reduce duplication between pipeline modes
- Use failed_reads directly instead of recomputing from total - passed
- Preserve primary pipeline error when secondary finalization also fails
- Pre-allocate secondary serialization buffer capacity to avoid per-batch reallocation
- Reset secondary_data in batch clear() methods
- Add write_raw_bytes to RawBamWriter for bulk secondary output
@nh13
nh13 force-pushed the nhomer/perf-filter-and-memory-estimate-fixes branch from 7420717 to 7fe3b0b Compare February 17, 2026 04:57
@nh13
nh13 temporarily deployed to github-actions February 17, 2026 04:57 — with GitHub Actions Inactive
@nh13 nh13 changed the title perf: remove filter single-threaded fallback and fix MemoryEstimate accuracy perf: remove filter single-threaded fallback, fix MemoryEstimate accuracy, and add ordered reject output Feb 17, 2026
@nh13
nh13 merged commit b10e946 into main Feb 17, 2026
7 checks passed
@nh13
nh13 deleted the nhomer/perf-filter-and-memory-estimate-fixes branch February 17, 2026 05:00
@nh13 nh13 mentioned this pull request Feb 18, 2026

This branch was previously deployed

1 inactive deployment
github-actions — 7fe3b0b1 Deployed Feb 17, 2026 by nh13 via coverage #386
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant