refactor: extract partition iteration and unify group selection

Move partition iteration logic to obikdump, introducing a FilteredPartitionIter trait over IndexCache for batch-oriented scanning with configurable data retrieval and early termination. Consolidate ingroup and outgroup index storage in GroupQuorumFilter into a unified Selection struct driven by predicate matching. Update dependency manifests to include obikidxcache, rayon, and obikentropy, and remove the deprecated dump_layer module while adjusting public API re-exports.
This commit is contained in:
Eric Coissac
2026-08-26 09:33:37 +02:00
parent 881b1532b5
commit 16ade823d6
10 changed files with 256 additions and 291 deletions
+5 -2
View File
@@ -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"
+44 -31
View File
@@ -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<W: Write, F: Fn() + Send + Sync>(
fn dump<W: Write, F: Fn() + Send + Sync>(
&self,
out: &mut W,
force_presence: bool,
debug: bool,
head: Option<usize>,
filters: &[Box<dyn KmerFilter>],
on_partition: F,
) -> OKIResult<()>;
}
impl IndexDump for KmerIndex {
fn dump<W: Write, F: Fn() + Send + Sync>(
&self,
out: &mut W,
force_presence: bool,
@@ -30,8 +49,8 @@ impl KmerIndex {
filters: &[Box<dyn KmerFilter>],
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::<u8>::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)
+3
View File
@@ -6,3 +6,6 @@
//! reverse), same pattern as `obikindexer`/`obikquery`.
mod dump;
mod partition_iter;
pub use partition_iter::FilteredPartitionIter;
+126
View File
@@ -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<dyn KmerFilter>],
cb: impl FnMut(CanonicalKmer, Box<[u32]>) -> bool,
) -> OKIResult<bool>;
/// 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<dyn KmerFilter>],
cb: impl FnMut(usize, usize, CanonicalKmer, Box<[u32]>) -> bool,
) -> OKIResult<bool>;
}
impl FilteredPartitionIter for IndexCache<'_> {
fn iter_partition_kmers(
&self,
part: usize,
use_counts: bool,
n_genomes: usize,
filters: &[Box<dyn KmerFilter>],
mut cb: impl FnMut(CanonicalKmer, Box<[u32]>) -> bool,
) -> OKIResult<bool> {
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<dyn KmerFilter>],
mut cb: impl FnMut(usize, usize, CanonicalKmer, Box<[u32]>) -> bool,
) -> OKIResult<bool> {
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<dyn KmerFilter>],
cb: &mut dyn FnMut(CanonicalKmer, Box<[u32]>) -> bool,
) -> OKIResult<bool> {
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<Vec<u32>> = if read_counts {
let mut cols: Vec<Vec<u32>> = vec![Vec::new(); n_genomes];
layer.fill_sub_matrix(&slots, &mut cols);
cols
} else {
let mut bool_cols: Vec<Vec<bool>> = 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)
}