From fd744325a0495121ecd61644db7f76a96ef946a4 Mon Sep 17 00:00:00 2001 From: Nils Homer Date: Mon, 28 Sep 2026 09:26:08 -0700 Subject: [PATCH] fix(runall): honor --extract::no-check-crc on BGZF FASTQ input The all-file BGZF FASTQ parallel-decode split (build_bgzf_fastq_split) took its CRC32 policy from ChainSpec::verify_crc, while the fused FASTQ readers resolved theirs from the extract options' check_crc / no_check_crc. runall set verify_crc to a hardcoded true for a FASTQ source, believing it inert, so under --extract::no-check-crc the split aborted with a CRC mismatch while the run logged "CRC verify: off"; standalone extract only worked because it resolved the field a second time itself. Resolve the split's policy in the builder from the same extract options the fused readers use, so no FASTQ command can leave the two decode fronts out of sync. ChainSpec::verify_crc is now genuinely inert for a FASTQ source, and extract's own copy of the resolution is removed. The default policy (verify a file, skip stdin) is unchanged. --- src/lib/commands/extract.rs | 80 ++-------------- src/lib/commands/runall.rs | 4 +- src/lib/pipeline/chains/builder.rs | 47 ++++++++-- src/lib/pipeline/chains/spec.rs | 7 +- tests/integration/helpers/fastq.rs | 23 +++++ tests/integration/test_runall_command.rs | 113 +++++++++++++++++++++++ 6 files changed, 193 insertions(+), 81 deletions(-) diff --git a/src/lib/commands/extract.rs b/src/lib/commands/extract.rs index 3c03a1b42..5dceb848a 100644 --- a/src/lib/commands/extract.rs +++ b/src/lib/commands/extract.rs @@ -133,27 +133,13 @@ pub(crate) fn open_fastq_reader( // CRC-verification policy, resolved per input path: an explicit flag wins, // otherwise verify a file and trust (skip) stdin. Only consulted for the // single-threaded BGZF path below. - let verify_crc = if check_crc { - true - } else if no_check_crc { - false - } else { - !is_stdin_path(path) - }; + let verify_crc = crate::commands::common::resolve_check_crc(check_crc, no_check_crc, path); // Mirror the BAM path's `CRC verify:` line (see `RunOptions::log_effective_check_crc`) // so the effective policy — including the default-off-for-stdin skip, a change from // fgumi's previous always-verify default — is visible in every run's log rather than a // silent decision. Emitted below only for BGZF input, the only format that carries a // per-block CRC32 to verify or skip. - let crc_reason = if check_crc { - " (--check-crc)" - } else if no_check_crc { - " (--no-check-crc)" - } else if is_stdin_path(path) { - " (trusted stdin)" - } else { - "" - }; + let crc_reason = crate::commands::common::check_crc_reason(check_crc, no_check_crc, path); // stdin cannot be sniffed by path and then re-opened — the bytes read to // classify it are gone. Peek them off the stream and chain them back in @@ -744,22 +730,14 @@ impl Extract { /// [`SourceSpec::Fastqs`] — and the sink is a BAM. The chain opens its own /// readers and detects the quality encoding in `ChainBuilder::open_source`. /// - /// Split out from [`Self::execute_chain`] so the flag-derived spec fields — - /// notably the resolved CRC policy — can be unit-tested without running the - /// pipeline, mirroring `Sort::build_sort_chain_spec`. - /// - /// `read_streams` is a seekable-BAM knob and inert for a FASTQ source, but - /// `verify_crc` is NOT inert: the all-file BGZF parallel-decode split - /// (`ChainBuilder::build_bgzf_fastq_split` → `FastqDecompress`) reads it as - /// its per-block CRC32 policy. It is resolved here through the shared - /// [`resolve_check_crc`] so `--check-crc` / `--no-check-crc` reach that - /// decoder instead of it silently skipping verification. + /// `read_streams` and `verify_crc` are inert for a FASTQ source: the CRC + /// policy reaches both FASTQ decode fronts through the extract options' + /// `check_crc`/`no_check_crc` (see `ChainSpec::verify_crc`). /// /// [`ChainSpec`]: crate::pipeline::chains::ChainSpec /// [`ChainSpec::single_stage`]: crate::pipeline::chains::ChainSpec::single_stage /// [`SourceSpec::Fastqs`]: crate::pipeline::chains::SourceSpec::Fastqs /// [`SourceSpec::InterleavedFastq`]: crate::pipeline::chains::SourceSpec::InterleavedFastq - /// [`resolve_check_crc`]: crate::commands::common::resolve_check_crc fn build_extract_chain_spec( &self, command_line: &str, @@ -776,19 +754,6 @@ impl Extract { let stage_opts = StageOptionsBag { extract: Some(self.to_extract_options()), ..Default::default() }; - // `verify_crc` is consumed only by the all-file BGZF parallel-decode - // split; the fused reader path derives its own per-path policy directly - // from `check_crc`/`no_check_crc`. The split's eligibility gate - // guarantees every input is a non-stdin BGZF file, so - // `resolve_check_crc`'s stdin-trust branch never applies on that path — - // resolve against the first input (a real file whenever the split runs; - // the value is inert when the fused path handles the source). - let verify_crc = crate::commands::common::resolve_check_crc( - self.check_crc, - self.no_check_crc, - &self.inputs[0], - ); - Ok(ChainSpec { stages: vec![Stage::Extract], source, @@ -800,7 +765,9 @@ impl Extract { queue_memory: self.queue_memory.clone(), async_reader: self.async_reader, read_streams: fgumi_bam_io::ReadStreams::Fixed(1), - verify_crc, + // Inert for a FASTQ source: both decode fronts resolve the CRC + // policy from `stage_opts.extract` (see `ChainSpec::verify_crc`). + verify_crc: true, command_line: command_line.to_string(), }) } @@ -2028,37 +1995,6 @@ mod tests { assert_crc_error_msg(&format!("{check_err:#}")); } - /// `--check-crc` / `--no-check-crc` reach the extract `ChainSpec.verify_crc` - /// through the shared [`resolve_check_crc`] policy, so the all-file BGZF - /// parallel-decode split (`ChainBuilder::build_bgzf_fastq_split` → - /// `FastqDecompress`) honors the flags instead of silently skipping CRC - /// verification. The field was previously hardcoded `false`, disabling CRC - /// on that path regardless of the flags. Mirrors - /// `Sort::build_sort_chain_spec_resolves_verify_crc`. - #[rstest] - #[case::file_default_verifies(false, false, true)] - #[case::file_no_check_crc_skips(false, true, false)] - #[case::file_check_crc_verifies(true, false, true)] - fn build_extract_chain_spec_resolves_verify_crc( - #[case] check_crc: bool, - #[case] no_check_crc: bool, - #[case] expected: bool, - ) { - // A file path (never stdin) — the split's eligibility gate the field - // feeds only fires on real files, so the file-default policy applies. - let extract = bgzf_crc_extract( - PathBuf::from("in.fq.gz"), - PathBuf::from("out.bam"), - ThreadingOptions::none(), - check_crc, - no_check_crc, - ); - let spec = extract - .build_extract_chain_spec("fgumi extract (test)") - .expect("spec build should succeed"); - assert_eq!(spec.verify_crc, expected, "file input CRC policy on the FASTQ split"); - } - /// End-to-end guard that the BGZF FASTQ split decoder itself honors the CRC /// policy, not just that the spec carries it. `test_bgzf_fastq_honors_check_crc` /// cannot prove this: its small input is fully read during quality-encoding diff --git a/src/lib/commands/runall.rs b/src/lib/commands/runall.rs index b98840b45..b65c4be43 100644 --- a/src/lib/commands/runall.rs +++ b/src/lib/commands/runall.rs @@ -1824,7 +1824,9 @@ impl Command for RunAll { // literal before that borrow is done with. let verify_crc = match &self.input { Some(p) => BamIoOptions::new(p, &self.output).effective_check_crc(), - // FASTQ source ignores verify_crc; value is inert. + // FASTQ source: inert. Both FASTQ decode fronts resolve their CRC + // policy from `--extract::check-crc` / `--extract::no-check-crc` + // (see `ChainBuilder::open_fastq_source` / `build_bgzf_fastq_split`). None => true, }; let source = self.derive_source_spec()?; diff --git a/src/lib/pipeline/chains/builder.rs b/src/lib/pipeline/chains/builder.rs index a54e00737..cb347687f 100644 --- a/src/lib/pipeline/chains/builder.rs +++ b/src/lib/pipeline/chains/builder.rs @@ -118,6 +118,10 @@ pub(crate) enum PendingSource { /// decoded (one continuous DEFLATE stream). Paths line up 1:1 with /// `readers` by stream index. bgzf_paths: Option>, + /// CRC32 policy for the BGZF split (`bgzf_paths`), resolved in + /// `open_fastq_source` from the same extract options the fused readers + /// use, so the two FASTQ decode fronts cannot disagree. + split_verify_crc: bool, }, } @@ -872,10 +876,30 @@ impl<'a> ChainBuilder<'a> { } else { None }; + // The split's CRC policy, from the same extract options the fused + // readers above used — NOT `spec.verify_crc`, which is the BAM/SAM + // source's knob (a command that left it hardcoded, as runall did, made + // the split ignore `--no-check-crc`). The eligibility gate guarantees + // every split path is a non-stdin BGZF file, so the first input's + // resolution (flag, else verify-a-file) holds for every stream; the + // value is unused when `bgzf_paths` is `None`. + let split_verify_crc = inputs.first().is_some_and(|first| { + crate::commands::common::resolve_check_crc( + extract_opts.check_crc, + extract_opts.no_check_crc, + first, + ) + }); Ok(( header, - PendingSource::Fastq { readers, encoding, force_round_robin: interleaved, bgzf_paths }, + PendingSource::Fastq { + readers, + encoding, + force_round_robin: interleaved, + bgzf_paths, + split_verify_crc, + }, )) } @@ -901,6 +925,7 @@ impl<'a> ChainBuilder<'a> { fn build_bgzf_fastq_split( &mut self, paths: &[std::path::PathBuf], + verify_crc: bool, num_threads: usize, batch_records: usize, byte_limit: u64, @@ -913,9 +938,7 @@ impl<'a> ChainBuilder<'a> { use crate::pipeline::steps::source::find_fastq_boundaries::FindFastqBoundaries; use crate::pipeline::steps::source::read_fastq::FastqOrdinalSequence; - // CRC policy mirrors the BAM decode path: honor the command's - // --check-crc / --no-check-crc (falling back to verify). - let verify_crc = self.spec.verify_crc; + // `verify_crc` is resolved by `open_fastq_source` (see `split_verify_crc`). // Honor --async-reader here too: the split re-opens each file raw, so // (unlike the fused path, which wraps in `open_fastq_reader`) it must // apply the prefetch wrap itself, or the flag would be a silent no-op @@ -1242,7 +1265,13 @@ impl<'a> ChainBuilder<'a> { self.current_tail = Some(unmapped_tail); self.paired_tail = Some(mapped_tail); } - PendingSource::Fastq { readers, encoding: _, force_round_robin, bgzf_paths } => { + PendingSource::Fastq { + readers, + encoding: _, + force_round_robin, + bgzf_paths, + split_verify_crc, + } => { use crate::pipeline::core::step::Affinity; use crate::pipeline::steps::source::read_fastq::{ FastqOrdinalSequence, ReadFastqInputs, @@ -1320,7 +1349,13 @@ impl<'a> ChainBuilder<'a> { // block-parallel decoded. let _ = force_round_robin; // per-stream output-queue backpressure bounds drift for all K let tail = if let Some(paths) = bgzf_paths { - self.build_bgzf_fastq_split(&paths, num_threads, batch_records, byte_limit)? + self.build_bgzf_fastq_split( + &paths, + split_verify_crc, + num_threads, + batch_records, + byte_limit, + )? } else { // Fused path: build one single-stream reader per stream, each // with its OWN FastqOrdinalSequence (each is its own edge into diff --git a/src/lib/pipeline/chains/spec.rs b/src/lib/pipeline/chains/spec.rs index 51e77bddc..b5d1187b2 100644 --- a/src/lib/pipeline/chains/spec.rs +++ b/src/lib/pipeline/chains/spec.rs @@ -42,8 +42,11 @@ pub struct ChainSpec { /// chain reproduces the non-chain path's CRC behavior. Honored on both BAM /// decode fronts: the `BgzfDecompress` step and the sort arena front's /// `InflateToArena` (the only path standalone `fgumi sort` takes). Inert for - /// the SAM source (no BGZF) and for the FASTQ source (which carries its own - /// policy). + /// the SAM source (no BGZF) and for FASTQ sources: both FASTQ decode fronts — + /// the fused readers (`ChainBuilder::open_fastq_source`) and the all-file + /// BGZF split (`ChainBuilder::build_bgzf_fastq_split`) — resolve their policy + /// from the extract options' `check_crc`/`no_check_crc`, so a FASTQ command + /// cannot leave the split out of sync with the flags it parsed. pub verify_crc: bool, /// For `@PG` line injection into the output header. pub command_line: String, diff --git a/tests/integration/helpers/fastq.rs b/tests/integration/helpers/fastq.rs index 2c110e858..fda5c8dda 100644 --- a/tests/integration/helpers/fastq.rs +++ b/tests/integration/helpers/fastq.rs @@ -17,3 +17,26 @@ pub fn write_gzip_fastq(path: &Path, records: &[(&str, &str, &str)]) { } encoder.finish().expect("finish gzip fastq"); } + +/// Write `fastq` (raw FASTQ text) BGZF-compressed to `path`, then flip one +/// byte of the CRC32 footer of the last *data* block. Callers pass enough +/// records to span several blocks, so the corruption can sit past any reader's +/// sampling window. +pub fn write_bgzf_fastq_with_corrupt_last_crc(path: &Path, fastq: &[u8]) { + let mut bytes = Vec::new(); + { + let mut writer = noodles::bgzf::io::Writer::new(&mut bytes); + writer.write_all(fastq).expect("write bgzf fastq"); + writer.finish().expect("finish bgzf fastq"); + } + let blocks = { + let mut cursor: &[u8] = &bytes; + fgumi_bgzf::read_raw_blocks(&mut cursor, 1_000_000).expect("read bgzf blocks") + }; + // `read_raw_blocks` skips the empty EOF marker, so the returned blocks are + // all data blocks and their total length ends at the last one's footer. + assert!(blocks.len() >= 2, "input must span >= 2 data blocks; got {}", blocks.len()); + let data_end: usize = blocks.iter().map(fgumi_bgzf::RawBgzfBlock::len).sum(); + bytes[data_end - fgumi_bgzf::BGZF_FOOTER_SIZE] ^= 0x01; + fs::write(path, bytes).expect("write corrupted bgzf fastq"); +} diff --git a/tests/integration/test_runall_command.rs b/tests/integration/test_runall_command.rs index 2aa8ca669..325f6e7a1 100644 --- a/tests/integration/test_runall_command.rs +++ b/tests/integration/test_runall_command.rs @@ -39,6 +39,7 @@ use crate::helpers::bam_generator::{ create_minimal_header, create_test_reference, create_umi_family_at_pos, write_bam, }; use crate::helpers::cutover::decompressed_records_without_pg; +use crate::helpers::fastq::write_bgzf_fastq_with_corrupt_last_crc; use crate::helpers::read_bam_output; use crate::helpers::{aligner_binary, build_aligner_index, write_gzip_fastq}; @@ -1991,6 +1992,118 @@ fn correct_rejects_without_correct_stage_warns() { assert!(!rejects.exists(), "no correct stage ran, so no rejects file may be written"); } +/// `runall --start-from extract` must honor `--extract::no-check-crc` on an +/// all-BGZF FASTQ file input, as standalone `fgumi extract --no-check-crc` +/// does. The chain's `verify_crc` was hardcoded `true` for a FASTQ source, so +/// the BGZF split decoder aborted with a CRC mismatch even while the run +/// logged `CRC verify: off`. The default (file ⇒ verify) must still reject. +#[test] +fn extract_honors_no_check_crc_on_bgzf_fastq() { + const NUM_RECORDS: usize = 350_000; + let tmp = TempDir::new().unwrap(); + let fastq = tmp.path().join("reads.fq.gz"); + let mut records = Vec::new(); + for i in 0..NUM_RECORDS { + writeln!(records, "@q{i}\nACGTACGTAC\n+\nIIIIIIIIII").unwrap(); + } + write_bgzf_fastq_with_corrupt_last_crc(&fastq, &records); + + let runall = |out: &Path, extra: &[&str]| { + let mut args = vec![ + "runall", + "--start-from", + "extract", + "--stop-after", + "extract", + "--extract::inputs", + p(&fastq), + "--extract::read-structures", + "+T", + "--extract::sample", + "s1", + "--extract::library", + "lib1", + "-o", + p(out), + ]; + args.extend_from_slice(extra); + fgumi(args) + }; + + let skipped = tmp.path().join("skipped.bam"); + let output = runall(&skipped, &["--extract::no-check-crc"]); + assert!( + output.status.success(), + "--extract::no-check-crc must accept a corrupted BGZF CRC32: {}", + String::from_utf8_lossy(&output.stderr) + ); + let standalone = tmp.path().join("standalone.bam"); + run_ok( + [ + "extract", + "--inputs", + p(&fastq), + "--read-structures", + "+T", + "--sample", + "s1", + "--library", + "lib1", + "--no-check-crc", + "-o", + p(&standalone), + ], + "standalone extract --no-check-crc", + ); + let (_, records) = read_bam_output(&skipped); + assert_eq!(records.len(), NUM_RECORDS, "every record must be extracted"); + // The standalone oracle runs the same decoder, so also pin identity and + // order against the fixture itself — including the corrupted last block. + for (i, record) in records.iter().enumerate() { + let name = record.name().map(ToString::to_string).unwrap_or_default(); + assert_eq!(name, format!("q{i}"), "record {i} name/order"); + } + assert_bams_record_equivalent_nonempty(&skipped, &standalone); + + let verified = tmp.path().join("verified.bam"); + let output = runall(&verified, &[]); + let stderr = String::from_utf8_lossy(&output.stderr).to_lowercase(); + assert!(!output.status.success(), "the default must reject a corrupted BGZF CRC32"); + // Match the mismatch error itself: a bare "crc" would also match the + // run's own `CRC verify: on` log line and pass on any unrelated failure. + assert!( + stderr.contains("crc32 mismatch") || stderr.contains("checksum mismatch"), + "the default must fail on the CRC mismatch, not some other error: {stderr}" + ); + // The corruption sits past quality-encoding detection's window, so only the + // BGZF split decoder can reach it — pin that it is the step enforcing the + // default, not some earlier reader. + assert!( + stderr.contains("step \"fastqdecompress\" failed"), + "the BGZF split decoder must be what rejects the corruption: {stderr}" + ); +} + +/// A BAM-source start without `-i` must still be told `--input` is missing, +/// not be misrouted into FASTQ-source (extract) option handling. +#[test] +fn bam_start_without_input_reports_missing_input() { + let tmp = TempDir::new().unwrap(); + assert_rejected_with( + [ + "runall", + "--start-from", + "sort", + "--stop-after", + "sort", + "-o", + p(&tmp.path().join("out.bam")), + ], + "--input is required with --start-from sort", + "runall --start-from sort without -i", + ); +} + // ══════════════════════════ Extract→Extract (interleaved, no aligner) ══════════════════════════ /// Validates A2 (`--extract::interleaved`'s interleaved-vs-count check):