Repository navigation
feat(retag): add --threads and route through the unified pipeline - #871
Conversation
|
Note Reviews pausedUse the following commands to manage reviews:
Use the checkboxes below for quick actions:
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (1)
Included review availability: 0 reviews are currently available. Your included PR review attempts over the past 7 days set your current allowance at 1 review per hour. WalkthroughRisk: threaded execution can alter record ordering, headers, and metric aggregation. Fix: select serial or unified pipeline execution and validate parity with threaded tests. ChangesRetag threading
Estimated code review effort: 4 (Complex) | ~45 minutes Merge Risk: ⚪ Minimal · up to The change adds optional threaded execution while preserving the existing serial behavior and reports targeted validation for output, metrics, ordering, compression, warnings, and invalid thread counts. No actionable merge-blocking risk remains beyond normal checks and review. Sequence Diagram(s)sequenceDiagram
participant Retag
participant BAMReader
participant PipelineWorkers
participant BAMWriter
participant Metrics
Retag->>BAMReader: read grouped records
BAMReader->>PipelineWorkers: submit raw records
PipelineWorkers->>PipelineWorkers: apply retag operations and aggregate counts
PipelineWorkers->>BAMWriter: emit records in input order
Retag->>Metrics: report record totals and operation counts
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
Comment |
|
@coderabbitai pause |
✅ Action performedReviews paused. |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #871 +/- ##
==========================================
- Coverage 94.69% 94.68% -0.01%
==========================================
Files 194 194
Lines 119898 120103 +205
==========================================
+ Hits 113532 113720 +188
- Misses 6366 6383 +17 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
retag was single-threaded with a hardcoded 1 reader / 1 writer serial loop. Add --threads N to route records through the unified 7-step pipeline, matching filter and clip, along with the --scheduler and --max-memory/--memory-reserve/--memory-per-thread knobs that ride with it. Without --threads the original serial fast path is unchanged. retag is a pure per-record transform, so the pipeline path uses SingleRawRecordGrouper (one record per unit, no template batching) and applies the ops in process_fn. Per-operation counts accumulate into a PerThreadAccumulator and are summed after the pipeline drains; because OpCounts are u64 and addition is commutative, the aggregated counts and the --metrics TSV are identical to the serial path at any thread count. This keeps --metrics available under --threads (a deliberate divergence from clip, whose per-read metrics are not summable). Also set group_key_config to the raw-no-cell config to skip a wasted per-record cell-barcode scan, and call log_effective_check_crc() so retag finally emits the per-run 'CRC verify:' line its --check-crc doc already promised. Output order is preserved by the pipeline, so threaded output matches serial record-for-record.
|
@coderabbitai review |
✅ Action performedReview finished.
|
What
Adds
--threads Ntofgumi retagand routes it through the existing unified 7-step BAM pipeline, matching howfilterandclipalready work. Without--threads, the original single-threaded serial loop runs unchanged.retagis a pure per-record SAM aux-tag rewriter (copy/move/delete ops applied left-to-right). It previously hardcoded a 1-reader / 1-writer serial loop with no way to opt into parallelism, while every other BAM-processing command exposed--threads.How
execute()keeps all up-front guards (io.validate, output-collision, input-aliasing via canonicalize, and nowlog_effective_check_crc) before the mode branch, so both paths enforce them identically.--threads→ the serial loop, extracted verbatim intorun_single_threaded.--threads N→run_threaded, usingSingleRawRecordGrouper(one record per group — retag needs no template batching) and applying the ops in the pipeline'sprocess_fn.OpCountsareu64counters accumulated into aPerThreadAccumulator<Vec<OpCounts>>and summed after the pipeline drains. Because addition is commutative and the counters areu64, the aggregated counts — and the--metricsTSV — are identical between serial and threaded modes at any thread count.--metricsdeliberately keeps working under--threads(unlikeclip, whose per-read metrics are not summable and are rejected under--threads).group_key_configto the raw-no-cell config to skip a per-record cell-barcode scan.The flattened
ThreadingOptions/SchedulerOptions/QueueMemoryOptions(i.e.--threads,--scheduler,--max-memory/--memory-reserve/--memory-per-thread) mirrorfilter/clip.--compression-leveland the CRC/async flags already existed.Also a small drive-by fix:
retagnow callslog_effective_check_crc(), finally emitting the per-runCRC verify:line its--check-crcdoc already promised.Correctness
Output order is preserved by the pipeline, so threaded output matches serial record-for-record. Verified via the real binary (serial vs
--threads 4): decoded records +--metricsTSV byte-identical; output headers differ only in the@PG CL:field (which correctly records the different argv);--compression-levelhonored in threaded mode; the zero-match warning fires in threaded mode;--threads 0is rejected cleanly.Tests
All in
retag.rs:threaded_output_matches_single_threaded(--threads 1and--threads 4): asserts decoded-record equality, header equality, and byte-identical--metricsTSV against the serial path, over a 50-record corpus spanning the classes where a decode divergence could hide — mapped/paired reads with real CIGAR +MC/RGtags (some secondary), records that pre-carryBXsodst_overwrittenis exercised, and a zero-matchdelete.threaded_preserves_input_record_order(--threads 4).threads_flag_coexists_with_positional_operations(clap parse:--threadsalongside the positionalSRC::op::DSTargs).cargo ci-test,cargo ci-fmt,cargo ci-lintall pass.Risk: output changes for threaded
retag, pinned by serial-equivalent metrics, ordering, validation, and compression tests;unsafechanges: none; memory and thread policy changes:--threads, scheduler, and queue-memory options add bounded threaded execution.Fix: use the threaded path only when
--threadsis set. Keep the existing serial path unchanged.--threads Ntofgumi retag.u64counts.