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):