Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
587 changes: 34 additions & 553 deletions src/lib/commands/filter.rs

Large diffs are not rendered by default.

23 changes: 9 additions & 14 deletions src/lib/pipeline/chains/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1140,10 +1140,10 @@ impl<'a> ChainBuilder<'a> {
/// pass entirely; for a SAM-first sort that pass runs once per record in
/// `ParseSamChunk` and would otherwise be pure waste. This changes only the
/// discarded key, never the record bytes, so output is unchanged.
/// (The legacy single-threaded path uses `new_raw_no_cell`, which still
/// pays the combined aux pass; name-hash-only is strictly less work.) All
/// other first stages (group/dedup/consensus) need the full position/cell
/// key, so they fall through to [`Self::bam_group_key_config`].
/// (The full-key config `new_raw_no_cell` still pays the combined aux pass;
/// name-hash-only is strictly less work.) All other first stages
/// (group/dedup/consensus) need the full position/cell key, so they fall
/// through to [`Self::bam_group_key_config`].
///
/// Every arm below uses a DEFAULT `LibraryIndex`, never `from_header`:
/// `name_hash_only` never reads `library_index` (see `name_hash_key` in
Expand All @@ -1152,11 +1152,7 @@ impl<'a> ChainBuilder<'a> {
/// `library_index` at all), so resolving it from the header is pure waste
/// for these stages — and `LibraryIndex::from_header` panics on a header
/// with more than 65,535 distinct `@RG` libraries, a needless crash risk
/// this avoids. (Some of these stages' serial oracles independently build
/// a full-header group key of their own — e.g. filter-by-template's
/// `build_filter_pipeline_config` calls `LibraryIndex::from_header` — so
/// this isn't "a panic the oracle never has"; it is simply dead work this
/// arm has no reason to repeat.)
/// this avoids.
fn source_group_key_config(&self) -> fgumi_bam_io::GroupKeyConfig {
match self.spec.stages.first() {
Some(
Expand Down Expand Up @@ -4437,11 +4433,10 @@ impl<'a> ChainBuilder<'a> {

// FILT3-02: filtering is template-based, so coordinate-sorted input
// silently scatters mates and corrupts the both-primaries-pass logic.
// Reject it here exactly as `Filter::execute`'s legacy path does, so
// the two orchestrations of the filter stage cannot drift on accepted
// orders, error text, or logging (mirrors `add_group`'s
// `require_group_input_ordering` call, shared verbatim with
// `Group::execute`).
// Reject it on the chain — the only filter execution path — so filter
// rejects mis-ordered input before any record is processed (mirrors
// `add_group`'s `require_group_input_ordering` call, shared verbatim
// with `Group::execute`).
crate::commands::common::require_query_grouped(
&self.header,
&input_path.display().to_string(),
Expand Down
8 changes: 4 additions & 4 deletions src/lib/pipeline/chains/commands/filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -225,7 +225,7 @@ pub(crate) fn build_filter_step_single_no_rejects(
let mut record = decoded.into_raw_bytes();
let (bases_masked, pass) =
process_record_raw_call(&mut record, &captures).map_err(io::Error::other)?;
// Match the legacy path's fgbio-parity "Total bases masked" tally:
// Match fgbio's "Total bases masked" tally:
// count masked bases only in a retained primary read (0 for a
// rejected read / secondary / supplementary), not the raw count.
bases_masked_total += retained_primary_masked_bases(
Expand Down Expand Up @@ -302,7 +302,7 @@ pub(crate) fn build_filter_step_single_with_rejects(
let mut record = decoded.into_raw_bytes();
let (bases_masked, pass) = process_record_raw_call(&mut record, &captures)
.map_err(io::Error::other)?;
// Match the legacy path's fgbio-parity "Total bases masked" tally:
// Match fgbio's "Total bases masked" tally:
// count masked bases only in a retained primary read (0 for a
// rejected read / secondary / supplementary), not the raw count.
bases_masked_total +=
Expand Down Expand Up @@ -375,7 +375,7 @@ pub(crate) fn build_filter_step_template_no_rejects(
}

let template_pass = template_passes(&template_records, &pass_map);
// Match the legacy path's fgbio-parity "Total bases masked" tally:
// Match fgbio's "Total bases masked" tally:
// count masked bases only in retained primary reads of a retained
// template (0 for a dropped template), not the raw per-record sum.
bases_masked_total += retained_primary_masked_bases(
Expand Down Expand Up @@ -466,7 +466,7 @@ pub(crate) fn build_filter_step_template_with_rejects(
}

let template_pass = template_passes(&template_records, &pass_map);
// Match the legacy path's fgbio-parity "Total bases masked" tally:
// Match fgbio's "Total bases masked" tally:
// count masked bases only in retained primary reads of a retained
// template (0 for a dropped template), not the raw per-record sum.
bases_masked_total += retained_primary_masked_bases(
Expand Down
5 changes: 2 additions & 3 deletions src/lib/pipeline/steps/group/queryname.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,8 @@ use crate::template::Template;
/// Serial mutex acquisition; mirrors `GroupByMi::MAX_BATCHES_PER_LOCK`.
const MAX_BATCHES_PER_LOCK: usize = 8;

/// Default target batch count. Mirrors legacy's
/// `TemplateGrouper::new(1000)` (the production batch size in
/// `Filter::execute_threads_mode_template`).
/// Default target batch count: the production batch size used when filter
/// groups templates on the chain (`fgumi filter --filter-by-template=true`).
pub const DEFAULT_TARGET_BATCH_COUNT: usize = 1000;

/// `Serial + ByItemOrdinal` queryname grouper. Records arriving in
Expand Down
97 changes: 97 additions & 0 deletions src/lib/pipeline/steps/parse/decode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -441,6 +441,103 @@ mod tests {
}
}

/// Guards the `name_hash_only` decode skip that the chain filter path relies
/// on (`source_group_key_config` routes `Stage::Filter` — and `Correct` — to
/// `GroupKeyConfig::name_hash_only`). filter groups templates through
/// `GroupByQueryname`, which reads only `key.name_hash`, and its process step
/// operates on the raw record bytes (untouched by the key config), so the skip
/// is output-safe **iff** `name_hash_only` produces the same `name_hash` the
/// full key would for every record — even when the fields the full key also
/// computes (library index from `RG`, unclipped position) vary across records.
///
/// This is the in-repo replacement for the parity that
/// `test_filter_chain_matches_single_threaded_with_rg_and_cb_variation` used
/// to provide before the cutover made both of its runs take the same
/// `name_hash_only` chain: it pins the invariant directly against the two
/// decode-consumer key branches (`compute_group_key_from_raw` vs
/// `name_hash_key`) without needing `FGUMI_BASELINE_BIN`.
#[test]
fn name_hash_only_matches_full_key_grouping_when_rg_and_position_vary() {
use crate::sam::SamTag;
use fgumi_bam_io::{GroupKey, LibraryIndex};
use fgumi_raw_bam::SamBuilder;
use fgumi_raw_bam::flags::{FIRST_SEGMENT, LAST_SEGMENT, PAIRED};
use noodles::sam::alignment::record::data::field::Tag;
use noodles::sam::header::record::value::Map;
use noodles::sam::header::record::value::map::ReadGroup;
use noodles::sam::header::record::value::map::read_group::tag as rg_tag;

// Two read groups in distinct libraries, so the full key's `library_idx`
// genuinely differs by RG — the field name_hash_only skips computing.
let mut header = noodles::sam::Header::builder();
for (id, library) in [("RG1", "libA"), ("RG2", "libB")] {
let rg = Map::<ReadGroup>::builder()
.insert(rg_tag::LIBRARY, String::from(library))
.build()
.expect("read group builds");
header = header.add_read_group(bstr::BString::from(id), rg);
}
let lib = LibraryIndex::from_header(&header.build());
let cb = Some(Tag::from([b'C', b'B']));

// Two paired templates, each R1+R2 sharing a name; the templates carry
// DIFFERENT read groups and DIFFERENT mapped positions, so the full key's
// library_idx and pos1 differ across them (the test is non-vacuous only if
// the skipped fields actually vary).
let mate = |name: &[u8], rg: &[u8], pos: i32, first: bool| -> fgumi_raw_bam::RawRecord {
let mut b = SamBuilder::new();
b.read_name(name)
.flags(PAIRED | if first { FIRST_SEGMENT } else { LAST_SEGMENT })
.ref_id(0)
.pos(pos)
.mapq(60)
.cigar_ops(&[4u32 << 4]) // 4M
.sequence(b"ACGT")
.qualities(&[30u8; 4]);
b.add_string_tag(SamTag::RG, rg);
b.build()
};
let records = [
mate(b"tmpl-A", b"RG1", 100, true),
mate(b"tmpl-A", b"RG1", 200, false),
mate(b"tmpl-B", b"RG2", 300, true),
mate(b"tmpl-B", b"RG2", 400, false),
];

let full: Vec<_> =
records.iter().map(|r| compute_group_key_from_raw(r.as_ref(), &lib, cb)).collect();
let name_only: Vec<_> = records.iter().map(|r| name_hash_key(r.as_ref())).collect();

for (i, (full_key, name_key)) in full.iter().zip(&name_only).enumerate() {
assert_eq!(
name_key.name_hash, full_key.name_hash,
"record {i}: name_hash_only must reproduce the full key's name_hash despite \
differing RG/position, or filter's GroupByQueryname would group differently"
);
// Everything except name_hash is left at the config-independent default,
// so the RG/position the full key computes cannot leak into the grouping
// key the chain actually uses.
assert_eq!(
*name_key,
GroupKey { name_hash: full_key.name_hash, ..GroupKey::default() },
"record {i}: name_hash_only must leave all non-name fields at default"
);
}

// Mates of a template share a name_hash under BOTH configs, so template
// membership (and thus filter's both-primaries aggregation) is identical.
assert_eq!(name_only[0].name_hash, name_only[1].name_hash, "tmpl-A mates group together");
assert_eq!(name_only[2].name_hash, name_only[3].name_hash, "tmpl-B mates group together");
assert_ne!(name_only[0].name_hash, name_only[2].name_hash, "distinct templates stay apart");

// Non-vacuous: the full key really does populate the fields name_hash_only
// drops, and they differ across the two read groups / positions — so the
// name_hash parity above is a genuine skip guard, not a comparison of two
// all-default keys.
assert_ne!(full[0].library_idx, full[2].library_idx, "full key varies library_idx by RG");
assert_ne!(full[0].pos1, full[2].pos1, "full key varies pos1 by mapped position");
}

// ========================================================================
// Fail-closed validation of undersized / malformed records
// ========================================================================
Expand Down
1 change: 1 addition & 0 deletions tests/integration/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ mod test_extract_command;
mod test_fastq_command;
mod test_fastq_pipeline_memory_backpressure;
mod test_filter_command;
mod test_filter_cutover_parity;
mod test_group_command;
mod test_group_determinism;
mod test_input_source_matrix;
Expand Down
Loading
Loading