Repository navigation
fix(runall)!: CLI surface and validation parity; consensus depth/tie-break hardening (5/6) - #462
Conversation
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
WalkthroughConsensus counts now clamp instead of wrapping, CLI validation and help text are tightened, the fused pipeline bridge accepts record filters, UMI tie-break ordering is aligned across implementations, and runtime checks reject duplicate buffering and stalled process reaping. ChangesConsensus saturation hardening
Command surface and validation
Pipeline bridge and chain wiring
UMI assignment tie-break parity
Runtime and buffering safety
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Suggested labels
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Warning There were issues while running some tools. Please review the errors and either fix the tool's configuration or disable the tool if it's a critical failure. 🔧 markdownlint-cli2 (0.22.1)CHANGELOG.mdmarkdownlint-cli2 wrapper config was not available before execution Comment |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## feat-runall #462 +/- ##
==============================================
Coverage ? 94.00%
==============================================
Files ? 110
Lines ? 49665
Branches ? 0
==============================================
Hits ? 46687
Misses ? 2978
Partials ? 0 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
d7af841 to
8eed314
Compare
c82d616 to
ad680ab
Compare
8eed314 to
f3f0daa
Compare
ad680ab to
71a68ca
Compare
|
@coderabbitai review |
1 similar comment
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Actionable comments posted: 12
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/zipper.rs (1)
3774-3784: 📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick winAdd a rejection case for the removed flag.
This only locks in the new spelling. Add an
expect_errcase for--skip-pa-tagsso the breaking CLI contract cannot regress silently.Suggested test
#[rstest] // --skip-tc-tags (default false) #[case(&["zipper", "-u", "u.bam", "-r", "ref.fa", "-o", "out.bam"], false)] #[case(&["zipper", "-u", "u.bam", "-r", "ref.fa", "-o", "out.bam", "--skip-tc-tags"], true)] #[case(&["zipper", "-u", "u.bam", "-r", "ref.fa", "-o", "out.bam", "--skip-tc-tags", "true"], true)] #[case(&["zipper", "-u", "u.bam", "-r", "ref.fa", "-o", "out.bam", "--skip-tc-tags", "false"], false)] #[case(&["zipper", "-u", "u.bam", "-r", "ref.fa", "-o", "out.bam", "--skip-tc-tags=true"], true)] #[case(&["zipper", "-u", "u.bam", "-r", "ref.fa", "-o", "out.bam", "--skip-tc-tags=false"], false)] fn test_skip_tc_tags_parsing(#[case] args: &[&str], #[case] expected: bool) { let cmd = Zipper::try_parse_from(args).expect("failed to parse Zipper arguments"); assert_eq!(cmd.options.skip_tc_tags, expected); } + +#[test] +fn test_skip_pa_tags_is_rejected() { + Zipper::try_parse_from([ + "zipper", + "-u", + "u.bam", + "-r", + "ref.fa", + "-o", + "out.bam", + "--skip-pa-tags", + ]) + .expect_err("`--skip-pa-tags` should remain rejected"); +}🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/lib/commands/zipper.rs` around lines 3774 - 3784, Add a negative parsing test for the removed CLI flag so the contract stays locked down. In the `test_skip_tc_tags_parsing` area of `zipper::Zipper::try_parse_from`, add an `expect_err` case for `--skip-pa-tags` to verify it is rejected instead of accepted. Keep the existing `skip_tc_tags` parsing cases, but ensure the new assertion explicitly covers the old spelling and fails if it ever becomes valid again.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@CHANGELOG.md`:
- Line 10: The changelog entry for the `fgumi zipper` flag rename is missing the
correct scoped `runall` replacement, which can mislead users copying the old
name. Update the note around the `skip-pa-tags` to `skip-tc-tags` rename so it
explicitly lists both replacements: the direct `fgumi zipper` flag and the
`runall`-scoped `--zipper::skip-tc-tags`, using the existing `runall`, `fgumi
zipper`, and `--skip-pa-tags` references in this changelog section.
In `@src/lib/aligner.rs`:
- Around line 219-228: The final warning in aligner::kill/cleanup logic is still
triggered when Child::try_wait returns a benign error, causing false “left
unreaped” messages. Update the loop around try_wait so benign wait errors like
ECHILD are treated as already reaped (or tracked via an unknown_or_reaped flag)
before the log::warn! gate, and keep the warning only for cases where the
process truly appears stuck after SIGKILL.
In `@src/lib/commands/codec.rs`:
- Around line 277-280: The disagreement-rate validation in codec.rs currently
only checks bounds, so f64::NAN can slip through and bypass the later
duplex_error_rate gate. Update the validation in the codec command logic around
max_duplex_disagreement_rate to reject non-finite values as well as out-of-range
ones, and add a regression test covering f64::NAN to ensure it fails validation.
In `@src/lib/commands/runall.rs`:
- Around line 326-329: Update the generated CLI help in the runall command so it
consistently uses the consensus-stage model instead of algorithm-specific stage
names. In the help text around the stage/start-stop documentation in runall.rs,
replace any references that still describe `simplex`/`duplex`/`codec` as
start/stop values and change codec-specific gating from `--stop-after=codec` to
`--consensus codec`. Keep the wording aligned with the existing consensus-stage
help near the `runall` CLI docs and the related field help section so both
places describe `consensus` as the stage and `--consensus codec` as the
selector.
- Around line 1573-1582: Reject non-regular first FASTQ inputs before the
sampling pre-read in runall; the current temp_reader path assumes a seekable
regular file, but open_fastq_reader/File::open can still accept FIFO/device
paths and cause the later extract pass to block or hit EOF. Add a file-type
check on opts.inputs.first() using metadata().file_type().is_file() before
building temp_reader, and fail early with a clear error if the first input is
not a regular file.
In `@src/lib/pipeline/chains/builder.rs`:
- Around line 2370-2383: The incompatibility checks for group options are
duplicated in the chain builder and can drift from the standalone `fgumi group`
validation. Move the paired/no-UMI/min-UMI-length rule enforcement into the
shared validator used by `GroupReadsByUmi::execute`, then have the runall path
in `builder.rs` call that same validator instead of hard-coding the
`group.strategy`, `group.no_umi`, and `group.min_umi_length` checks locally.
In `@src/lib/pipeline/steps/source/pair_fastq.rs`:
- Around line 145-154: The duplicate-slot invariant in the pair_fastq source
path is currently enforced only by debug_assert! in the chunk insertion logic,
so release builds can silently overwrite an existing slot and corrupt read
identity. Update the code in the chunk handling path that uses slot.take() and
assigns *slot = Some(chunk) to enforce this invariant in production as well,
either by replacing debug_assert! with a hard assert! or by returning an error
before any overwrite occurs. Keep the check tied to the existing
chunk.chunk_serial and stream-specific slot handling so the failure happens
whenever a duplicate (serial, stream) pair is seen.
In `@src/lib/umi/parallel_assigner.rs`:
- Around line 1459-1529: The new parity tests in ParallelAdjacencyAssigner only
compare against AdjacencyUmiAssigner, so a shared regression could still pass
without matching the intended fgbio behavior. Add an fgbio-backed oracle for the
mixed-case equal-count tie-break cases in the proptest and targeted test near
proptest_adjacency_parallel_matches_sequential /
adjacency_tie_break_is_case_sensitive_raw_not_uppercased, or clearly document
the divergence contract if fgbio parity is intentionally unavailable. Keep the
existing sequential comparison, but extend coverage so the raw tie-break change
is validated against both the sequential path and the fgbio baseline.
In `@tests/integration/test_runall_parity.rs`:
- Around line 914-949: The parity helper currently only exercises the runall
path and matches against a fixed stderr substring, so it can miss regressions in
standalone validation. Update assert_runall_group_combo_rejects to also invoke
the standalone fgumi group path with the same extra_group_flags, using the
existing Stage::Group and fgumi helpers, and compare both failures rather than
relying on a hard-coded fragment. Keep the check focused on ensuring the
standalone group rejection and the fused runall rejection both fail consistently
and emit equivalent validation behavior.
- Around line 568-645: The new parity test only checks fused/staged equivalence
and does not independently verify that the both-unmapped family is actually
dropped. In new_002_group_to_duplex_allow_unmapped_parity, add a negative oracle
alongside assert_bams_record_equivalent: compare the outputs against a baseline
duplex fixture that excludes the unmapped pair, or assert directly that no
record from that extra family appears in either runall_out or staged_out. Use
the existing fgumi invocations and assert_bams_record_equivalent as the main
flow, but strengthen the test so it fails if the unmapped family leaks even when
both paths match.
- Around line 951-965: The two invalid group combo tests are duplicating the
same helper pattern, so replace them with a single parameterized `rstest` case
in `rejects_group_paired_with_min_umi_length`/`rejects_group_no_umi_with_paired`
area using `assert_runall_group_combo_rejects` as the shared assertion target.
Convert the current wrapper tests into one `#[rstest]` table that supplies the
args slice and expected error message for each invalid `--group::strategy
paired` combination, so future cases can be added in one matrix instead of
separate test functions.
In `@tests/integration/test_streaming_input.rs`:
- Around line 183-188: The stdin integration test currently only checks exit
status and file creation, which is too weak for the identity case. Update the
test around the streamed-input flow in test_streaming_input.rs so it reads both
sorted_bam and output_bam after running the command, then compare their BAM
records for identical content and order. Keep the existing existence-check
assertions, but add record-level validation in the same test path to ensure the
piped input output is truly equivalent, not just non-empty.
---
Outside diff comments:
In `@src/lib/commands/zipper.rs`:
- Around line 3774-3784: Add a negative parsing test for the removed CLI flag so
the contract stays locked down. In the `test_skip_tc_tags_parsing` area of
`zipper::Zipper::try_parse_from`, add an `expect_err` case for `--skip-pa-tags`
to verify it is rejected instead of accepted. Keep the existing `skip_tc_tags`
parsing cases, but ensure the new assertion explicitly covers the old spelling
and fails if it ever becomes valid again.
🪄 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: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 6a68193c-c1e9-4ab4-8495-1bd6c0201597
📒 Files selected for processing (23)
CHANGELOG.mdcrates/fgumi-consensus/src/base_builder.rscrates/fgumi-consensus/src/codec_caller.rscrates/fgumi-umi/src/assigner.rssrc/lib/aligner.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/dedup.rssrc/lib/commands/downsample.rssrc/lib/commands/extract.rssrc/lib/commands/filter.rssrc/lib/commands/runall.rssrc/lib/commands/zipper.rssrc/lib/pipeline/chains/builder.rssrc/lib/pipeline/chains/commands/correct.rssrc/lib/pipeline/chains/validate.rssrc/lib/pipeline/steps/source/pair_fastq.rssrc/lib/umi/parallel_assigner.rssrc/main.rstests/integration/helpers/bam_generator.rstests/integration/test_runall_parity.rstests/integration/test_streaming_input.rstests/integration/test_zipper_command.rs
💤 Files with no reviewable changes (1)
- src/lib/pipeline/chains/commands/correct.rs
f3f0daa to
5246192
Compare
71a68ca to
a92836c
Compare
5246192 to
5a7573e
Compare
a92836c to
46df46f
Compare
5a7573e to
c9c3502
Compare
46df46f to
6cf8aba
Compare
c9c3502 to
40cd503
Compare
6cf8aba to
1f4afc6
Compare
40cd503 to
c904a02
Compare
1f4afc6 to
2d1a84e
Compare
|
@coderabbitai review |
1 similar comment
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Actionable comments posted: 4
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
src/lib/pipeline/steps/align_and_merge.rs (1)
1699-1708: 🩺 Stability & Availability | 🔴 Critical | ⚡ Quick winFailed-kill teardown can still hang on thread joins.
If
aligner.kill()returnsfalse, the child may still hold stdout open, so the reader can remain blocked inread()and the writer can remain blocked ontoken_tx.send; the unconditional joins then reintroduce the cancel/error-path hang. Track the kill result and only join finished threads on the failed-kill path, otherwise detach/log.Suggested fix
- if let Some(mut aligner) = self.aligner.take() { + let mut aligner_reaped = true; + if let Some(mut aligner) = self.aligner.take() { // Bounded teardown: only issue the unbounded `wait()` once the child // has actually been reaped. If `kill()` timed out (stuck child, e.g. // uninterruptible I/O), calling `wait()` would block teardown // forever — defeating the 1s kill deadline. In that case we skip // `wait()` and let `AlignerProcess::Drop` (itself bounded) clean up. - if aligner.kill() { + aligner_reaped = aligner.kill(); + if aligner_reaped { let _ = aligner.wait(); } } if let Some(h) = self.writer_thread.take() { - let _ = h.join(); + if aligner_reaped || h.is_finished() { + let _ = h.join(); + } else { + log::warn!("aligner writer thread left detached after failed bounded kill"); + } } if let Some(h) = self.reader_thread.take() { - let _ = h.join(); + if aligner_reaped || h.is_finished() { + let _ = h.join(); + } else { + log::warn!("aligner reader thread left detached after failed bounded kill"); + } }As per path instructions, “This is the hand-rolled concurrent step pipeline; its bugs are deadlocks, lost output, and unbounded memory, not style.”
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/lib/pipeline/steps/align_and_merge.rs` around lines 1699 - 1708, The teardown in the align_and_merge pipeline can hang when `aligner.kill()` fails because `writer_thread.join()` and `reader_thread.join()` are still called even though the child may keep stdout open and block those threads. Update the shutdown logic around `aligner.kill()`, `self.writer_thread.take()`, and `self.reader_thread.take()` to track the kill result and skip blocking joins on the failed-kill path; instead detach, log, or otherwise avoid waiting for threads that may never finish.Source: Path instructions
crates/fgumi-consensus/src/codec_caller.rs (1)
1346-1362: 🩺 Stability & Availability | 🟠 Major | ⚡ Quick winUse the fallible joined-name builder for the no-UMI fallback too.
The UMI path now returns a typed error, but
prefix:counterstill goes through non-falliblebuild_record; an over-long user-controlled prefix can still panic. Usetry_build_record_joined_name(prefix, counter)there as well.Based on learnings, over-long consensus joined-name failures should propagate as typed errors rather than panic or be silently dropped.
Proposed fix
} else { - let read_name = format!("{}:{}", self.read_name_prefix, self.consensus_counter); - self.bam_builder.build_record( - read_name.as_bytes(), + let suffix = self.consensus_counter.to_string(); + self.bam_builder.try_build_record_joined_name( + self.read_name_prefix.as_bytes(), + suffix.as_bytes(), flag, &consensus.bases, &consensus.quals, - ); + )?; }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/fgumi-consensus/src/codec_caller.rs` around lines 1346 - 1362, The no-UMI fallback in codec_caller::CodecCaller still uses the infallible build_record path with a formatted prefix:counter name, so an over-long read_name_prefix can panic instead of returning a typed error. Update that branch to use the same fallible joined-name builder as the UMI path, namely try_build_record_joined_name with the prefix and counter-derived name, and propagate the Result so consensus record creation fails cleanly rather than panicking.Source: Learnings
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/fgumi-consensus/src/codec_caller.rs`:
- Around line 1419-1430: Keep the scalar strand summary tags/rates in sync with
the clamped per-base arrays in codec_caller::append strand-tag logic. The
current AD/AM/AE/BD/BM/BE values are still computed from the original u16
depths/errors while AD_BASES/AE_BASES/BD_BASES/BE_BASES are saturated via
saturate_u16_to_i16, so values above i16::MAX can make downstream parity checks
fail. Update the scalar summaries to derive from the same clamped values used
for the serialized arrays, or add a generated fgbio-parity test in
codec_caller.rs that documents and verifies any intentional divergence.
In `@src/lib/commands/runall.rs`:
- Around line 426-432: Update the `runall` help text for the `--ref` option so
the “not consumed” note only applies to the pure `correct`-only path; in the
`runall` option docstring near `validate_align_and_merge`/`execute`, replace the
current `--start-from correct` wording with a version that either scopes it to
`--stop-after correct` or explicitly says `--ref` may still be needed by later
stages in a `correct → align/zipper/...` chain.
In `@tests/integration/test_streaming_input.rs`:
- Around line 159-182: The streaming input test currently spawns an external cat
process just to forward the BAM into stdin, which adds an unnecessary dependency
and leaves the child unmanaged. Update the downsample test setup in the
streaming input integration test to open sorted_bam directly and pass it into
Command::new(env!("CARGO_BIN_EXE_fgumi")) via Stdio::from(File::open(...)) (or
equivalent stdin piping), removing cat_child and its stdout handoff entirely.
In `@tests/integration/test_zipper_command.rs`:
- Around line 250-257: The zipper integration test is too permissive and no
longer enforces the intended mapped-first preflight contract. In the test case
in test_zipper_command.rs, tighten the stderr assertion to require the Mapped
BAM (--input) existence error only, since zipper::run in
src/lib/commands/zipper.rs checks --input before --unmapped; keep the failure
check, but remove the alternate Unmapped BAM branch so regressions in mapped
validation are caught.
---
Outside diff comments:
In `@crates/fgumi-consensus/src/codec_caller.rs`:
- Around line 1346-1362: The no-UMI fallback in codec_caller::CodecCaller still
uses the infallible build_record path with a formatted prefix:counter name, so
an over-long read_name_prefix can panic instead of returning a typed error.
Update that branch to use the same fallible joined-name builder as the UMI path,
namely try_build_record_joined_name with the prefix and counter-derived name,
and propagate the Result so consensus record creation fails cleanly rather than
panicking.
In `@src/lib/pipeline/steps/align_and_merge.rs`:
- Around line 1699-1708: The teardown in the align_and_merge pipeline can hang
when `aligner.kill()` fails because `writer_thread.join()` and
`reader_thread.join()` are still called even though the child may keep stdout
open and block those threads. Update the shutdown logic around `aligner.kill()`,
`self.writer_thread.take()`, and `self.reader_thread.take()` to track the kill
result and skip blocking joins on the failed-kill path; instead detach, log, or
otherwise avoid waiting for threads that may never finish.
🪄 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: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 627f7027-7358-4326-b889-3d62ac0e4ac4
📒 Files selected for processing (25)
CHANGELOG.mdcrates/fgumi-consensus/src/base_builder.rscrates/fgumi-consensus/src/codec_caller.rscrates/fgumi-consensus/src/vanilla_caller.rscrates/fgumi-umi/src/assigner.rssrc/lib/aligner.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/dedup.rssrc/lib/commands/downsample.rssrc/lib/commands/extract.rssrc/lib/commands/filter.rssrc/lib/commands/runall.rssrc/lib/commands/zipper.rssrc/lib/pipeline/chains/builder.rssrc/lib/pipeline/chains/commands/correct.rssrc/lib/pipeline/chains/validate.rssrc/lib/pipeline/steps/align_and_merge.rssrc/lib/pipeline/steps/source/pair_fastq.rssrc/lib/umi/parallel_assigner.rssrc/main.rstests/integration/helpers/bam_generator.rstests/integration/test_runall_parity.rstests/integration/test_streaming_input.rstests/integration/test_zipper_command.rs
💤 Files with no reviewable changes (1)
- src/lib/pipeline/chains/commands/correct.rs
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
|
@coderabbitai review |
✅ Action performedReview finished.
|
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 (1)
src/lib/pipeline/steps/align_and_merge.rs (1)
1699-1708: 🩺 Stability & Availability | 🟠 Major | ⚡ Quick winSkip pipe-thread joins when
kill()times out.If Line 1699 returns
false, the child may still hold stdout/stdin open, so Lines 1704-1708 can block forever joining reader/writer threads. Trackaligner_reaped; only join these pipe threads when the child was reaped, otherwise drop theJoinHandles and log that teardown was detached.Proposed fix
- if let Some(mut aligner) = self.aligner.take() { + let aligner_reaped = if let Some(mut aligner) = self.aligner.take() { // Bounded teardown: only issue the unbounded `wait()` once the child // has actually been reaped. If `kill()` timed out (stuck child, e.g. // uninterruptible I/O), calling `wait()` would block teardown // forever — defeating the 1s kill deadline. In that case we skip // `wait()` and let `AlignerProcess::Drop` (itself bounded) clean up. if aligner.kill() { let _ = aligner.wait(); + true + } else { + false } - } + } else { + true + }; - if let Some(h) = self.writer_thread.take() { - let _ = h.join(); - } - if let Some(h) = self.reader_thread.take() { - let _ = h.join(); + if aligner_reaped { + if let Some(h) = self.writer_thread.take() { + let _ = h.join(); + } + if let Some(h) = self.reader_thread.take() { + let _ = h.join(); + } + } else { + log::warn!("aligner kill timed out; detaching aligner pipe threads during teardown"); + let _ = self.writer_thread.take(); + let _ = self.reader_thread.take(); }As per path instructions, “treat any change to the reorder stage, queue push/pop, backpressure, or producer gating as liveness- and memory-critical” and flag pipeline deadlocks.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/lib/pipeline/steps/align_and_merge.rs` around lines 1699 - 1708, The teardown in the align-and-merge stage can deadlock because the pipe-thread joins run even when `aligner.kill()` times out and the child may still be holding stdout/stdin open. Update the cleanup logic in `align_and_merge` to track an `aligner_reaped` flag from the `kill()`/`wait()` path, and only call `self.writer_thread.take()` and `self.reader_thread.take()` with `join()` when the child was successfully reaped; otherwise drop the `JoinHandle`s and emit a log that teardown was detached. Use the existing `aligner.kill()`, `aligner.wait()`, `writer_thread`, and `reader_thread` flow as the anchor points for the fix.Source: Path instructions
♻️ Duplicate comments (1)
crates/fgumi-consensus/src/codec_caller.rs (1)
1419-1430: 🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick winKeep the scalar strand summaries on the same clamp path as the serialized arrays.
AD/AM/AEandBD/BM/BEare still computed from the unclampedu16vectors just above, whilead/ae/bd/benow clamp ati16::MAX. Above32767, one record can therefore emitmax(AD_BASES) != ADand strand error rates that no longer match the serialized arrays. Either derive the scalar summaries from the clamped arrays too, or add an end-to-end fgbio parity test and document the divergence. As per path instructions,crates/{fgumi-consensus,fgumi-umi,fgumi-metrics}/**/*.rschanges need identity coverage against the fgbio baseline or a documented divergence.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/fgumi-consensus/src/codec_caller.rs` around lines 1419 - 1430, The scalar strand summary tags are still being computed from the original unclamped depth/error values, so they can diverge from the serialized i16 arrays once values exceed i16::MAX. Update the logic in codec_caller.rs around the strand summary and array tag construction so AD/AM/AE and BD/BM/BE are derived from the same clamped values used by append_i16_array_tag, or add an explicit fgbio parity test plus documentation if the divergence is intentional. Use the existing saturate_u16_to_i16 path and the strand summary variables in this block to keep the scalar and array outputs aligned.Source: Path instructions
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@src/lib/umi/parallel_assigner.rs`:
- Around line 461-483: The parallel UMI assignment fallback is still creating
fresh molecule IDs for invalid UMIs instead of matching the sequential path’s
MoleculeId::None behavior. Update the remap logic around the BitEnc filtering
and the final per-read fallback in parallel_assigner so invalid, non-encodable
UMIs are cached as invalid reads and consistently mapped to MoleculeId::None,
including the all-invalid case and any mixed batch fallback. Keep the behavior
aligned with AdjacencyUmiAssigner / the sequential assigner semantics so
parallel and single-threaded results stay identical.
In `@tests/integration/test_runall_parity.rs`:
- Around line 613-645: The staged parity test only checks the final duplex
output, so it can miss a regression where group --allow-unmapped drops
fam_unmapped before duplex runs. In test_runall_parity.rs, strengthen the staged
path around fgumi(group_args) by asserting group_bam still contains the expected
RX=TTAA-GGCC record before invoking fgumi(duplex_args), using the existing
assert_bams_record_equivalent/assert helpers or a direct BAM scan. This keeps
the oracle independent and verifies both “retained by group” and “dropped by
duplex” using the group_bam and staged_out flow.
---
Outside diff comments:
In `@src/lib/pipeline/steps/align_and_merge.rs`:
- Around line 1699-1708: The teardown in the align-and-merge stage can deadlock
because the pipe-thread joins run even when `aligner.kill()` times out and the
child may still be holding stdout/stdin open. Update the cleanup logic in
`align_and_merge` to track an `aligner_reaped` flag from the `kill()`/`wait()`
path, and only call `self.writer_thread.take()` and `self.reader_thread.take()`
with `join()` when the child was successfully reaped; otherwise drop the
`JoinHandle`s and emit a log that teardown was detached. Use the existing
`aligner.kill()`, `aligner.wait()`, `writer_thread`, and `reader_thread` flow as
the anchor points for the fix.
---
Duplicate comments:
In `@crates/fgumi-consensus/src/codec_caller.rs`:
- Around line 1419-1430: The scalar strand summary tags are still being computed
from the original unclamped depth/error values, so they can diverge from the
serialized i16 arrays once values exceed i16::MAX. Update the logic in
codec_caller.rs around the strand summary and array tag construction so AD/AM/AE
and BD/BM/BE are derived from the same clamped values used by
append_i16_array_tag, or add an explicit fgbio parity test plus documentation if
the divergence is intentional. Use the existing saturate_u16_to_i16 path and the
strand summary variables in this block to keep the scalar and array outputs
aligned.
🪄 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: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 1e34c9a6-5ed3-46a2-8113-0ab8643f6e64
📒 Files selected for processing (25)
CHANGELOG.mdcrates/fgumi-consensus/src/base_builder.rscrates/fgumi-consensus/src/codec_caller.rscrates/fgumi-consensus/src/vanilla_caller.rscrates/fgumi-umi/src/assigner.rssrc/lib/aligner.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/dedup.rssrc/lib/commands/downsample.rssrc/lib/commands/extract.rssrc/lib/commands/filter.rssrc/lib/commands/runall.rssrc/lib/commands/zipper.rssrc/lib/pipeline/chains/builder.rssrc/lib/pipeline/chains/commands/correct.rssrc/lib/pipeline/chains/validate.rssrc/lib/pipeline/steps/align_and_merge.rssrc/lib/pipeline/steps/source/pair_fastq.rssrc/lib/umi/parallel_assigner.rssrc/main.rstests/integration/helpers/bam_generator.rstests/integration/test_runall_parity.rstests/integration/test_streaming_input.rstests/integration/test_zipper_command.rs
💤 Files with no reviewable changes (1)
- src/lib/pipeline/chains/commands/correct.rs
b2861c8 to
1f0a9dd
Compare
S5c1-001..008, S5c2-001..012, S5b1-003..006, S5b2-001/003, S5d-003/005/007, NEW-002: thread the duplex record_filter into the fused bridge; delete the dead --raw-tag/--assign-tag flags and rename --skip-pa-tags to --skip-tc-tags (breaking); enforce codec/group/zipper validation parity between standalone and runall; aligner kill diagnostics and display_order fixes.
X5-004: PairRawFastq::buffer_chunk added chunk.data.len() to pending_total_bytes then wrote the slot unconditionally. If a slot were ever overwritten, the dropped chunk's bytes would never be subtracted, drifting the total that feeds the catastrophic-desync guard and the backpressure decision. Each reader mints one chunk per (chunk_serial, stream), so overwrite is impossible today — but the accounting had no guard. Add a debug_assert pinning the invariant and, on the currently-impossible overwrite path, subtract the dropped chunk's bytes first so the total cannot silently drift if a future reader change (retry/replay/re-chunk) violates it. Test asserts the guard trips on a duplicate slot.
…(S7-002/FU-004) The sequential AdjacencyUmiAssigner tie-broke equal-count UMIs by the raw first-seen string (matching fgbio), while the parallel ParallelAdjacencyAssigner keyed on the uppercased string. For mixed-case input the two paths ordered equal-count roots oppositely and captured a shared neighbor into different groups. Converge the parallel path onto the raw first-seen UMI string (tracked alongside the uppercased count key; node identity and counting stay case-insensitive via BitEnc), keeping the sequential and paired paths on their case-sensitive raw keys to match fgbio. A non-vacuous direction test (aAAA/CAAA/GAAA) pins that GAAA groups with CAAA under raw bytes — and fails under an uppercased tie-break — plus a mixed-case parallel-vs-sequential parity proptest.
…ow (S7-001/S7-005/FU-002) The CODEC duplex consensus merge summed per-position depths and error counts as u16 (`da + db`, `ea + eb`, ...), which overflows at >65535 reads/strand/position (debug panic, release wraparound). The per-base aD/bD/ae/be arrays were emitted via a raw truncating `d as i16` cast, producing silently-negative depths for values in 32768..=65535. The per-base observation accumulator (`observations[idx] += 1`) and its `contributions()` sum were likewise unchecked u16. Compute the duplex merge sums in u32 so the intermediate math is exact, then saturate once at the u16 storage boundary. Replace the truncating per-base i16 cast with a saturating `i16::try_from(...).unwrap_or(i16::MAX)` (i16::MAX is the genuine SAM short-array ceiling, matching fgbio's Short). Saturate the per-base observation counter and contributions sum. Each saturation point logs a once-per-process warning so an adversarial/corrupt input that collapses a huge read set into one MI group is observable rather than silently mis-tagged. This is latent on real data (requires >32k-65k reads on one strand of one molecule at a single position) and pre-existing since the initial commit; severity minor (defensive hardening). Adds tests: in-range depths/errors pass through exactly, while values past u16::MAX/i16::MAX saturate without wrapping or going negative.
c37c1eb to
79ad42f
Compare
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
src/lib/pipeline/steps/align_and_merge.rs (1)
1693-1708: 🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy liftSkip thread joins when bounded kill times out.
If
aligner.kill()returnsfalse, the child can still hold stdout open, leavingreader_threadblocked inread()andwriter_threadblocked behind the token path. The unconditional joins here can still hang teardown forever even thoughwait()is skipped. Only join on a reaped child, or detach both handles on timeout.Suggested direction
- if let Some(mut aligner) = self.aligner.take() { + let reaped = if let Some(mut aligner) = self.aligner.take() { // Bounded teardown: only issue the unbounded `wait()` once the child // has actually been reaped. If `kill()` timed out (stuck child, e.g. // uninterruptible I/O), calling `wait()` would block teardown // forever — defeating the 1s kill deadline. In that case we skip // `wait()` and let `AlignerProcess::Drop` (itself bounded) clean up. if aligner.kill() { let _ = aligner.wait(); + true + } else { + false } - } + } else { + true + }; - if let Some(h) = self.writer_thread.take() { - let _ = h.join(); - } - if let Some(h) = self.reader_thread.take() { - let _ = h.join(); + if reaped { + if let Some(h) = self.writer_thread.take() { + let _ = h.join(); + } + if let Some(h) = self.reader_thread.take() { + let _ = h.join(); + } + } else { + let _ = self.writer_thread.take(); + let _ = self.reader_thread.take(); }As per path instructions,
src/lib/{unified_pipeline,pipeline}/**/*.rs: “This is the hand-rolled concurrent step pipeline; its bugs are deadlocks...” and cancel/liveness bugs should be flagged.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/lib/pipeline/steps/align_and_merge.rs` around lines 1693 - 1708, The teardown in align_and_merge still blocks because writer_thread and reader_thread are always joined even when aligner.kill() times out. Update the shutdown logic in the align_and_merge step so the joins only happen after a successful kill()/wait() path, and on timeout either skip joining or explicitly detach the thread handles. Use the existing aligner.kill(), aligner.wait(), writer_thread.take(), and reader_thread.take() flow to keep teardown bounded and avoid deadlocks when the child keeps stdout open.Source: Path instructions
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Outside diff comments:
In `@src/lib/pipeline/steps/align_and_merge.rs`:
- Around line 1693-1708: The teardown in align_and_merge still blocks because
writer_thread and reader_thread are always joined even when aligner.kill() times
out. Update the shutdown logic in the align_and_merge step so the joins only
happen after a successful kill()/wait() path, and on timeout either skip joining
or explicitly detach the thread handles. Use the existing aligner.kill(),
aligner.wait(), writer_thread.take(), and reader_thread.take() flow to keep
teardown bounded and avoid deadlocks when the child keeps stdout open.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro
Run ID: 7c601abb-75d6-4209-9058-18d834656cef
📒 Files selected for processing (25)
CHANGELOG.mdcrates/fgumi-consensus/src/base_builder.rscrates/fgumi-consensus/src/codec_caller.rscrates/fgumi-consensus/src/vanilla_caller.rscrates/fgumi-umi/src/assigner.rssrc/lib/aligner.rssrc/lib/commands/codec.rssrc/lib/commands/correct.rssrc/lib/commands/dedup.rssrc/lib/commands/downsample.rssrc/lib/commands/extract.rssrc/lib/commands/filter.rssrc/lib/commands/runall.rssrc/lib/commands/zipper.rssrc/lib/pipeline/chains/builder.rssrc/lib/pipeline/chains/commands/correct.rssrc/lib/pipeline/chains/validate.rssrc/lib/pipeline/steps/align_and_merge.rssrc/lib/pipeline/steps/source/pair_fastq.rssrc/lib/umi/parallel_assigner.rssrc/main.rstests/integration/helpers/bam_generator.rstests/integration/test_runall_parity.rstests/integration/test_streaming_input.rstests/integration/test_zipper_command.rs
💤 Files with no reviewable changes (1)
- src/lib/pipeline/chains/commands/correct.rs
Stack: 5 / 6 · base:
nh/runall-fix-4-per-crateThe runall command surface plus two consensus follow-up hardening items.
fix(runall)!— thread the duplexrecord_filterinto the fused bridge; delete the dead--raw-tag/--assign-tagflags and rename--skip-pa-tagsto--skip-tc-tags(both breaking); enforce codec/group/zipper validation parity between the standalone and runall paths; aligner kill diagnostics anddisplay_orderfixes.u32and saturate once at theu16storage boundary with a diagnostic, fixing silent wrap for counts in 32768..65535.Summary by CodeRabbit
runallno longer accepts--raw-tag/-tand--assign-tag/-T.zipper(andrunall) renames--skip-pa-tagsto--skip-tc-tags.--consensusnow reports an error instead of panicking..dictsidecars.