diff --git a/src/Cargo.lock b/src/Cargo.lock index 40c259ba..1e1fcdfc 100644 --- a/src/Cargo.lock +++ b/src/Cargo.lock @@ -1524,7 +1524,10 @@ name = "obikdump" version = "0.1.0" dependencies = [ "obikfilter", + "obikidxcache", "obikindex", + "obikseq", + "rayon", ] [[package]] @@ -1539,6 +1542,7 @@ name = "obikfilter" version = "0.1.0" dependencies = [ "obicompactvec", + "obikentropy", "obikindex", "obikseq", "obiskio", diff --git a/src/obikdump/Cargo.toml b/src/obikdump/Cargo.toml index b802741d..2d7ecc60 100644 --- a/src/obikdump/Cargo.toml +++ b/src/obikdump/Cargo.toml @@ -4,5 +4,8 @@ version = "0.1.0" edition = "2024" [dependencies] -obikindex = { path = "../obikindex" } -obikfilter = { path = "../obikfilter" } +obikindex = { path = "../obikindex" } +obikfilter = { path = "../obikfilter" } +obikidxcache = { path = "../obikidxcache" } +obikseq = { path = "../obikseq" } +rayon = "1" diff --git a/src/obikdump/src/dump.rs b/src/obikdump/src/dump.rs index 560ef95f..c3ae6b05 100644 --- a/src/obikdump/src/dump.rs +++ b/src/obikdump/src/dump.rs @@ -5,9 +5,14 @@ use rayon::prelude::*; use obikindex::{OKIError, OKIResult}; use obikindex::KmerIndex; +use obikidxcache::index_cache::IndexCache; use obikfilter::KmerFilter; -impl KmerIndex { +use crate::partition_iter::FilteredPartitionIter; + +/// Raw content export of a `KmerIndex` — `KmerIndex` is a foreign type +/// (`obikindex`), so this is an extension trait rather than an inherent `impl`. +pub trait IndexDump { /// Write a CSV table of all indexed kmers to `out`. /// /// Columns: `kmer`, then one column per genome (in index order). @@ -17,11 +22,25 @@ impl KmerIndex { /// the output uses 0/1 presence columns. /// /// Partitions are scanned in parallel; each partition buffers its output locally - /// before the main thread writes the chunks in partition order. + /// before the main thread writes the chunks in partition order. Each partition's + /// layers are cached (`IndexCache`) only for the scan of that one partition — + /// `self` is a complete, read-only source index, never a destination. /// /// The caller must have set the global kmer length (`obikseq::set_k`) before /// calling this method. - pub fn dump( + fn dump( + &self, + out: &mut W, + force_presence: bool, + debug: bool, + head: Option, + filters: &[Box], + on_partition: F, + ) -> OKIResult<()>; +} + +impl IndexDump for KmerIndex { + fn dump( &self, out: &mut W, force_presence: bool, @@ -30,8 +49,8 @@ impl KmerIndex { filters: &[Box], on_partition: F, ) -> OKIResult<()> { - let genomes = self.meta.genomes().map_err(OKIError::Io)?; - let use_counts = self.meta.config.with_counts && !force_presence; + let genomes = self.meta().genomes().map_err(OKIError::Io)?; + let use_counts = self.meta().config.with_counts && !force_presence; let n_genomes = genomes.len().max(1); let kmer_size = self.kmer_size(); @@ -68,20 +87,17 @@ impl KmerIndex { Ok(_) => { write_row(buf, row, prefix); true } } }; + let cache = IndexCache::new(self, Some(vec![i])); if debug { - self - .iter_partition_kmers_located(i, use_counts, n_genomes, filters, |part, layer, kmer, row| { - let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); - try_write(&mut buf, &row, &format!("{part},{layer},{seq}")) - }) - .map_err(OKIError::Partition)?; + cache.iter_partition_kmers_located(i, use_counts, n_genomes, filters, |part, layer, kmer, row| { + let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); + try_write(&mut buf, &row, &format!("{part},{layer},{seq}")) + })?; } else { - self - .iter_partition_kmers(i, use_counts, n_genomes, filters, |kmer, row| { - let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); - try_write(&mut buf, &row, &seq) - }) - .map_err(OKIError::Partition)?; + cache.iter_partition_kmers(i, use_counts, n_genomes, filters, |kmer, row| { + let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); + try_write(&mut buf, &row, &seq) + })?; } on_partition(); Ok(buf) @@ -90,22 +106,19 @@ impl KmerIndex { // ── Unbounded: no atomic, no contention ─────────────────────────── (0..n).into_par_iter().map(|i| { let mut buf = Vec::::new(); + let cache = IndexCache::new(self, Some(vec![i])); if debug { - self - .iter_partition_kmers_located(i, use_counts, n_genomes, filters, |part, layer, kmer, row| { - let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); - write_row(&mut buf, &row, &format!("{part},{layer},{seq}")); - true - }) - .map_err(OKIError::Partition)?; + cache.iter_partition_kmers_located(i, use_counts, n_genomes, filters, |part, layer, kmer, row| { + let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); + write_row(&mut buf, &row, &format!("{part},{layer},{seq}")); + true + })?; } else { - self - .iter_partition_kmers(i, use_counts, n_genomes, filters, |kmer, row| { - let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); - write_row(&mut buf, &row, &seq); - true - }) - .map_err(OKIError::Partition)?; + cache.iter_partition_kmers(i, use_counts, n_genomes, filters, |kmer, row| { + let seq = String::from_utf8(kmer.to_ascii()).unwrap_or_else(|_| "?".repeat(kmer_size)); + write_row(&mut buf, &row, &seq); + true + })?; } on_partition(); Ok(buf) diff --git a/src/obikdump/src/lib.rs b/src/obikdump/src/lib.rs index f953444f..45a7124d 100644 --- a/src/obikdump/src/lib.rs +++ b/src/obikdump/src/lib.rs @@ -6,3 +6,6 @@ //! reverse), same pattern as `obikindexer`/`obikquery`. mod dump; +mod partition_iter; + +pub use partition_iter::FilteredPartitionIter; diff --git a/src/obikdump/src/partition_iter.rs b/src/obikdump/src/partition_iter.rs new file mode 100644 index 00000000..b419c3e9 --- /dev/null +++ b/src/obikdump/src/partition_iter.rs @@ -0,0 +1,126 @@ +//! Filtered, batch-oriented iteration over an already-cached index's +//! partitions/layers — the read side of `obikfilter`'s `KmerFilter`s. +//! `IndexCache` is a foreign type (`obikidxcache`), so this is an extension +//! trait rather than an inherent `impl`. +//! +//! Only meaningful on a *complete* source index: `IndexCache` panics if a +//! layer is missing, which a finished index never has. Never use this on a +//! destination index still being built (see `obikmerge::partition_merge`, +//! which follows the same source-only-cache rule). + +use obikindex::OKIResult; +use obikindex::layer::{KmerLayer, LayerContent}; +use obikidxcache::index_cache::IndexCache; +use obikseq::CanonicalKmer; + +use obikfilter::{KmerFilter, passes_all}; + +/// Kmers pulled per batch from a layer before filtering — keeps matrix reads +/// grouped by (partition, layer) for locality instead of hopping row to row +/// across the index. Same convention as `obikphylo::siblings::build`. +const BATCH_SIZE: usize = 32768; + +pub trait FilteredPartitionIter { + /// Iterate all indexed kmers in partition `part`, calling `cb(kmer, row)` for each + /// kmer that passes every filter in `filters`. + /// + /// `use_counts = true` → reads count columns (u32 values per genome), only + /// meaningful for `Count` layers. `use_counts = false` → reads presence + /// columns, converted to 0/1 u32 (works for both `Count` and `Presence` + /// layers — counts collapse to presence). + /// + /// Returns `Ok(true)` if all kmers were visited, `Ok(false)` if the callback halted. + fn iter_partition_kmers( + &self, + part: usize, + use_counts: bool, + n_genomes: usize, + filters: &[Box], + cb: impl FnMut(CanonicalKmer, Box<[u32]>) -> bool, + ) -> OKIResult; + + /// Like [`iter_partition_kmers`](Self::iter_partition_kmers) but the callback + /// also receives `(partition, layer)` indices, enabling debug output that + /// identifies where each kmer was stored. + fn iter_partition_kmers_located( + &self, + part: usize, + use_counts: bool, + n_genomes: usize, + filters: &[Box], + cb: impl FnMut(usize, usize, CanonicalKmer, Box<[u32]>) -> bool, + ) -> OKIResult; +} + +impl FilteredPartitionIter for IndexCache<'_> { + fn iter_partition_kmers( + &self, + part: usize, + use_counts: bool, + n_genomes: usize, + filters: &[Box], + mut cb: impl FnMut(CanonicalKmer, Box<[u32]>) -> bool, + ) -> OKIResult { + for l in 0..self.n_layer(part).unwrap_or(0) { + let layer = self.get_layer(part, l).expect("layer within n_layer(part)"); + if !iter_layer_kmers(layer, use_counts, n_genomes, filters, &mut |kmer, row| cb(kmer, row))? { + return Ok(false); + } + } + Ok(true) + } + + fn iter_partition_kmers_located( + &self, + part: usize, + use_counts: bool, + n_genomes: usize, + filters: &[Box], + mut cb: impl FnMut(usize, usize, CanonicalKmer, Box<[u32]>) -> bool, + ) -> OKIResult { + for l in 0..self.n_layer(part).unwrap_or(0) { + let layer = self.get_layer(part, l).expect("layer within n_layer(part)"); + if !iter_layer_kmers(layer, use_counts, n_genomes, filters, &mut |kmer, row| cb(part, l, kmer, row))? { + return Ok(false); + } + } + Ok(true) + } +} + +/// Batch-and-transpose one layer's kmers into per-kmer filtered rows. +/// Returns `Ok(false)` if `cb` asked to stop early. +fn iter_layer_kmers( + layer: &KmerLayer, + use_counts: bool, + n_genomes: usize, + filters: &[Box], + cb: &mut dyn FnMut(CanonicalKmer, Box<[u32]>) -> bool, +) -> OKIResult { + let read_counts = use_counts && matches!(layer.content(), LayerContent::Count); + + for kmers in layer.iter_kmers_batch(BATCH_SIZE) { + // Kmers come straight from this layer's own iterator, so every one + // is a guaranteed member — a raw hash is enough, no membership + // recheck needed (see `KmerLayer::hash_batch`'s own doc). + let slots = layer.hash_batch(&kmers); + + let cols: Vec> = if read_counts { + let mut cols: Vec> = vec![Vec::new(); n_genomes]; + layer.fill_sub_matrix(&slots, &mut cols); + cols + } else { + let mut bool_cols: Vec> = vec![Vec::new(); n_genomes]; + layer.fill_sub_matrix_carries(&slots, &mut bool_cols); + bool_cols.iter().map(|c| c.iter().map(|&b| b as u32).collect()).collect() + }; + + for (i, kmer) in kmers.into_iter().enumerate() { + let row: Box<[u32]> = cols.iter().map(|c| c[i]).collect(); + if passes_all(filters, kmer, &row, n_genomes) && !cb(kmer, row) { + return Ok(false); + } + } + } + Ok(true) +} diff --git a/src/obikfilter/Cargo.toml b/src/obikfilter/Cargo.toml index f1bde2f8..750e9f25 100644 --- a/src/obikfilter/Cargo.toml +++ b/src/obikfilter/Cargo.toml @@ -9,3 +9,4 @@ obicompactvec = { path = "../obicompactvec" } obikseq = { path = "../obikseq" } obiskio = { path = "../obiskio" } obitaxonomy = { path = "../obitaxonomy" } +obikentropy = { path = "../obikentropy" } diff --git a/src/obikfilter/src/dump_layer.rs b/src/obikfilter/src/dump_layer.rs deleted file mode 100644 index 47507806..00000000 --- a/src/obikfilter/src/dump_layer.rs +++ /dev/null @@ -1,192 +0,0 @@ -use obikindex::layer::MphfLayer; -use obicompactvec::{PersistentBitMatrix, PersistentCompactIntMatrix}; -use obikseq::CanonicalKmer; -use obiskio::UnitigFileReader; -use obikindex::{OKIError, OKIResult}; - -use crate::filter::{KmerFilter, passes_all}; -use obikindex::KmerIndex; - -impl KmerIndex { - /// Iterate all indexed kmers in partition `part`, calling `cb(kmer, row)` for each - /// kmer that passes every filter in `filters`. - /// - /// `use_counts = true` → reads count columns (u32 values per genome). - /// `use_counts = false` → reads presence columns, converted to 0/1 u32. - /// - /// If no data matrix exists for a layer (pure set-membership, single genome), - /// a row of `n_genomes` ones is emitted for every kmer in that layer — unless - /// the filter rejects it, in which case the whole layer is skipped. - /// Like [`iter_partition_kmers`] but the callback returns `false` to stop early. - /// Returns `Ok(true)` if all kmers were visited, `Ok(false)` if the callback halted. - pub fn iter_partition_kmers( - &self, - part: usize, - use_counts: bool, - n_genomes: usize, - filters: &[Box], - mut cb: impl FnMut(CanonicalKmer, Box<[u32]>) -> bool, - ) -> OKIResult { - let index_dir = self.index_dir(part); - if !index_dir.exists() { - return Ok(true); - } - - let mut l = 0; - loop { - let layer_dir = self.layer_dir(part, l)?; - if !layer_dir.exists() { - break; - } - l += 1; - let mphf = MphfLayer::open(&layer_dir)?; - let reader = UnitigFileReader::open_sequential(&layer_dir.join("unitigs.bin"))?; - - let counts_dir = layer_dir.join("counts"); - let presence_dir = layer_dir.join("presence"); - - let cont = if use_counts && counts_dir.exists() { - let mat = PersistentCompactIntMatrix::open(&layer_dir).map_err(OKIError::Io)?; - let mut cont = true; - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - if let Some(slot) = mphf.find(kmer) { - let row = mat.row(slot); - if passes_all(filters, kmer, &row, n_genomes) { - cont = cb(kmer, row); - if !cont { - break; - } - } - } - } - cont - } else if !use_counts && presence_dir.exists() { - let mat = PersistentBitMatrix::open(&layer_dir).map_err(OKIError::Io)?; - let mut cont = true; - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - if let Some(slot) = mphf.find(kmer) { - let row: Box<[u32]> = mat.row(slot).iter().map(|&b| b as u32).collect(); - if passes_all(filters, kmer, &row, n_genomes) { - cont = cb(kmer, row); - if !cont { - break; - } - } - } - } - cont - } else { - // No data matrix: implicit presence — all values are 1. `row` - // is identical for every kmer, but a filter can still depend - // on the kmer's own sequence (e.g. MinComplexity), so this - // cannot be evaluated once for the whole layer — filters must - // still be tested per kmer. - let all_present: Box<[u32]> = vec![1u32; n_genomes].into(); - let mut cont = true; - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - if mphf.find(kmer).is_some() - && passes_all(filters, kmer, &all_present, n_genomes) - { - cont = cb(kmer, all_present.clone()); - if !cont { - break; - } - } - } - cont - }; - - if !cont { - return Ok(false); - } - } - - Ok(true) - } - - /// Like [`iter_partition_kmers`] but the callback also receives `(partition, layer)` - /// indices, enabling debug output that identifies where each kmer was stored. - /// Returns `Ok(true)` if all kmers were visited, `Ok(false)` if the callback halted. - pub fn iter_partition_kmers_located( - &self, - part: usize, - use_counts: bool, - n_genomes: usize, - filters: &[Box], - mut cb: impl FnMut(usize, usize, CanonicalKmer, Box<[u32]>) -> bool, - ) -> OKIResult { - let index_dir = self.index_dir(part); - if !index_dir.exists() { - return Ok(true); - } - - let mut layer = 0; - loop { - let layer_dir = self.layer_dir(part, layer); - if !layer_dir.exists() { - break; - } - let mphf = MphfLayer::open(&layer_dir)?; - let reader = UnitigFileReader::open_sequential(&layer_dir.join("unitigs.bin"))?; - - let counts_dir = layer_dir.join("counts"); - let presence_dir = layer_dir.join("presence"); - - let cont = if use_counts && counts_dir.exists() { - let mat = PersistentCompactIntMatrix::open(&layer_dir).map_err(OKIError::Io)?; - let mut cont = true; - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - if let Some(slot) = mphf.find(kmer) { - let row = mat.row(slot); - if passes_all(filters, kmer, &row, n_genomes) { - cont = cb(part, layer, kmer, row); - if !cont { - break; - } - } - } - } - cont - } else if !use_counts && presence_dir.exists() { - let mat = PersistentBitMatrix::open(&layer_dir).map_err(OKIError::Io)?; - let mut cont = true; - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - if let Some(slot) = mphf.find(kmer) { - let row: Box<[u32]> = mat.row(slot).iter().map(|&b| b as u32).collect(); - if passes_all(filters, kmer, &row, n_genomes) { - cont = cb(part, layer, kmer, row); - if !cont { - break; - } - } - } - } - cont - } else { - // Same as iter_partition_kmers: row is constant but a filter - // may still depend on the kmer's own sequence, so this must - // be tested per kmer, not once for the whole layer. - let all_present: Box<[u32]> = vec![1u32; n_genomes].into(); - let mut cont = true; - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - if mphf.find(kmer).is_some() - && passes_all(filters, kmer, &all_present, n_genomes) - { - cont = cb(part, layer, kmer, all_present.clone()); - if !cont { - break; - } - } - } - cont - }; - - if !cont { - return Ok(false); - } - layer += 1; - } - - Ok(true) - } -} diff --git a/src/obikfilter/src/filter.rs b/src/obikfilter/src/filter.rs index 9f54a695..ffb8b745 100644 --- a/src/obikfilter/src/filter.rs +++ b/src/obikfilter/src/filter.rs @@ -1,6 +1,8 @@ use obicompactvec::FilterMask; use obikseq::CanonicalKmer; +use crate::predicate::Selection; + /// Trait for kmer filters. /// /// `kmer` is the k-mer's own canonical sequence, reconstructed from the @@ -173,14 +175,12 @@ impl KmerFilter for MaxTotalCount { // ── Group-based quorum filter ───────────────────────────────────────────────── -/// Quorum filter operating on pre-classified genome groups. +/// Quorum filter operating on a pre-classified genome [`Selection`]. /// -/// `ingroup_idx` / `outgroup_idx` are column indices into the per-genome row. -/// When `ingroup_idx` is empty, no ingroup quorum is checked. -/// When `outgroup_idx` is empty, no outgroup quorum is checked. +/// `selection.ingroup_idx` / `selection.outgroup_idx` are column indices into +/// the per-genome row. When empty, the corresponding quorum is not checked. pub struct GroupQuorumFilter { - pub ingroup_idx: Vec, - pub outgroup_idx: Vec, + pub selection: Selection, pub threshold: u32, pub min_count: usize, pub max_count: usize, @@ -225,22 +225,22 @@ impl GroupQuorumFilter { impl KmerFilter for GroupQuorumFilter { fn passes(&self, _kmer: CanonicalKmer, row: &[u32], _n_genomes: usize) -> bool { - if !self.ingroup_idx.is_empty() { - let n = self.ingroup_idx.iter() + if !self.selection.ingroup_idx.is_empty() { + let n = self.selection.ingroup_idx.iter() .filter(|&&i| row.get(i).copied().unwrap_or(0) > self.threshold) .count(); - let denom = self.ingroup_idx.len(); + let denom = self.selection.ingroup_idx.len(); if n < self.min_count { return false; } if n > self.max_count { return false; } let frac = n as f64 / denom as f64; if frac < self.min_frac { return false; } if frac > self.max_frac { return false; } } - if !self.outgroup_idx.is_empty() { - let n = self.outgroup_idx.iter() + if !self.selection.outgroup_idx.is_empty() { + let n = self.selection.outgroup_idx.iter() .filter(|&&i| row.get(i).copied().unwrap_or(0) > self.threshold) .count(); - let denom = self.outgroup_idx.len(); + let denom = self.selection.outgroup_idx.len(); if n < self.min_outgroup_count { return false; } if n > self.max_outgroup_count { return false; } let frac = n as f64 / denom as f64; @@ -253,17 +253,17 @@ impl KmerFilter for GroupQuorumFilter { fn column_mask_expr(&self, _n_genomes: usize) -> Option { let t = self.threshold.checked_add(1)?; let mut parts: Vec = Vec::new(); - if !self.ingroup_idx.is_empty() { + if !self.selection.ingroup_idx.is_empty() { Self::group_mask_parts( - &self.ingroup_idx, t, + &self.selection.ingroup_idx, t, self.min_count, self.max_count, self.min_frac, self.max_frac, &mut parts, ); } - if !self.outgroup_idx.is_empty() { + if !self.selection.outgroup_idx.is_empty() { Self::group_mask_parts( - &self.outgroup_idx, t, + &self.selection.outgroup_idx, t, self.min_outgroup_count, self.max_outgroup_count, self.min_outgroup_frac, self.max_outgroup_frac, &mut parts, diff --git a/src/obikfilter/src/lib.rs b/src/obikfilter/src/lib.rs index 26a95c23..a6b5736f 100644 --- a/src/obikfilter/src/lib.rs +++ b/src/obikfilter/src/lib.rs @@ -4,19 +4,18 @@ //! operates on already-retained k-mers. //! //! [`filter`] (the `KmerFilter` trait + its implementations) depends only -//! on `obicompactvec`/`obikseq`, not on `obikindex`. [`dump_layer`] is the -//! extension over `obikindex::KmerIndex` that actually iterates a -//! partition's k-mers through those filters (`iter_partition_kmers`, -//! `iter_partition_kmers_located`) — every k-mer, filtered or not, goes -//! through this same path (`passes_all` on an empty filter list is always -//! `true`), so this crate depends one-way on `obikindex`, not the reverse. +//! on `obicompactvec`/`obikseq`, not on `obikindex`. The partition/layer +//! iteration that actually runs these filters over an index +//! (`iter_partition_kmers`, `iter_partition_kmers_located`) lives in +//! `obikdump` instead (`FilteredPartitionIter`) — it needs `obikidxcache`'s +//! `IndexCache` to read a *complete* source index, which this crate has no +//! reason to depend on. mod filter; -mod dump_layer; mod predicate; pub use filter::{ GroupQuorumFilter, KmerFilter, MaxGenomeCount, MaxGenomeFraction, MaxTotalCount, MinComplexity, MinGenomeCount, MinGenomeFraction, MinTotalCount, passes_all, }; -pub use predicate::{GroupFilterParams, MetaPred}; +pub use predicate::{GenomeSelector, GroupFilterParams, MetaPred, Selection}; diff --git a/src/obikfilter/src/predicate.rs b/src/obikfilter/src/predicate.rs index db1be2d8..065c7cff 100644 --- a/src/obikfilter/src/predicate.rs +++ b/src/obikfilter/src/predicate.rs @@ -68,14 +68,6 @@ impl MetaPred { } } -impl GenomeInfo { - /// Evaluate a single metadata predicate against this genome. - /// Returns `None` when the predicate's key is absent (NA propagation). - pub fn matches(&self, pred: &MetaPred) -> Option { - pred.eval(&self.meta) - } -} - // ── Path matching ───────────────────────────────────────────────────────────── /// True if the stored taxonomy `value` matches `pattern`. @@ -153,19 +145,49 @@ pub struct GroupFilterParams { pub max_outgroup_frac: Option, } -impl IndexMeta { - /// Returns indices of genomes matching `pred_str` (single predicate). - pub fn matching_genome_indices(&self, pred_str: &str) -> Result, String> { - let pred = MetaPred::parse(pred_str)?; - let genomes = self.genomes().map_err(|e| e.to_string())?; - Ok(genomes.iter().enumerate() - .filter_map(|(i, g)| { - if g.matches(&pred) == Some(true) { Some(i) } else { std::option::Option::None } - }) - .collect()) +// ── Genome selector ────────────────────────────────────────────────────────── + +/// Result of running a [`GenomeSelector`] against an index's metadata. +pub struct Selection { + pub ingroup_idx: Vec, + pub outgroup_idx: Vec, +} + +pub struct GenomeSelector { + pub(crate) ingroup: Vec, + pub(crate) outgroup: Vec, +} + +impl GenomeSelector { + /// Parse ingroup (AND'd) and outgroup (OR'd) predicate strings. + pub fn parse(ingroup: &[String], outgroup: &[String]) -> Result { + let ingroup = ingroup.iter().map(|s| MetaPred::parse(s)).collect::, _>>()?; + let outgroup = outgroup.iter().map(|s| MetaPred::parse(s)).collect::, _>>()?; + Ok(Self { ingroup, outgroup }) } - /// Build a `GroupQuorumFilter` from parsed predicates, evaluated against `self.genomes`. + /// Classify `meta`'s genomes into ingroup/outgroup indices. + /// + /// - No predicates at all: every genome is (implicitly) ingroup. + /// - Otherwise: ingroup wins on overlap; uncategorized genomes are dropped. + pub fn run(&self, meta: &IndexMeta) -> Result { + let genomes = meta.genomes().map_err(|e| e.to_string())?; + + if self.ingroup.is_empty() && self.outgroup.is_empty() { + return Ok(Selection { ingroup_idx: (0..genomes.len()).collect(), outgroup_idx: vec![] }); + } + + let members = classify(&genomes, &self.ingroup, &self.outgroup); + let ingroup_idx: Vec = members.iter().enumerate() + .filter(|(_, m)| matches!(m, Membership::Ingroup)) + .map(|(i, _)| i).collect(); + let outgroup_idx: Vec = members.iter().enumerate() + .filter(|(_, m)| matches!(m, Membership::Outgroup)) + .map(|(i, _)| i).collect(); + Ok(Selection { ingroup_idx, outgroup_idx }) + } + + /// Build a `GroupQuorumFilter` from this selector's classification of `meta`. /// /// - No groups defined: `ingroup_idx` = all genomes (implicit ingroup). /// - `ingroup` predicates only: outgroup indices are empty. @@ -173,34 +195,20 @@ impl IndexMeta { /// - Both defined: ingroup wins on overlap; uncategorized genomes are ignored. pub fn build_group_filter( &self, - ingroup_preds: &[MetaPred], - outgroup_preds: &[MetaPred], - p: GroupFilterParams, + meta: &IndexMeta, + p: GroupFilterParams, ) -> Result { - let genomes = self.genomes().map_err(|e| e.to_string())?; - let (ingroup_idx, outgroup_idx) = if ingroup_preds.is_empty() && outgroup_preds.is_empty() { - ((0..genomes.len()).collect(), vec![]) - } else { - let members = classify(&genomes, ingroup_preds, outgroup_preds); - let in_idx: Vec = members.iter().enumerate() - .filter(|(_, m)| matches!(m, Membership::Ingroup)) - .map(|(i, _)| i).collect(); - let out_idx: Vec = members.iter().enumerate() - .filter(|(_, m)| matches!(m, Membership::Outgroup)) - .map(|(i, _)| i).collect(); - (in_idx, out_idx) - }; - - let in_size = ingroup_idx.len(); - let out_size = outgroup_idx.len(); + let selection = self.run(meta)?; + let in_size = selection.ingroup_idx.len(); + let out_size = selection.outgroup_idx.len(); let ingroup_quorum_explicit = p.min_count.is_some() || p.max_count.is_some() || p.min_frac.is_some() || p.max_frac.is_some(); let outgroup_quorum_explicit = p.min_outgroup_count.is_some() || p.max_outgroup_count.is_some() || p.min_outgroup_frac.is_some() || p.max_outgroup_frac.is_some(); - let default_min_frac = if !ingroup_preds.is_empty() && !ingroup_quorum_explicit { 1.0 } else { 0.0 }; - let default_max_outgroup_count = if !outgroup_preds.is_empty() && !outgroup_quorum_explicit { 0 } else { out_size }; + let default_min_frac = if !self.ingroup.is_empty() && !ingroup_quorum_explicit { 1.0 } else { 0.0 }; + let default_max_outgroup_count = if !self.outgroup.is_empty() && !outgroup_quorum_explicit { 0 } else { out_size }; // Resolve a signed count: negative means an offset from the group size // (e.g. -1 = all but one), floored at 1 so the negative form always keeps @@ -238,8 +246,7 @@ impl IndexMeta { } Ok(GroupQuorumFilter { - ingroup_idx, - outgroup_idx, + selection, threshold: p.threshold, min_count, max_count, @@ -252,3 +259,4 @@ impl IndexMeta { }) } } +