From 904d85f33bbbc80694a2bd666ed4370c927e04b2 Mon Sep 17 00:00:00 2001 From: Eric Coissac Date: Wed, 26 Aug 2026 09:08:41 +0200 Subject: [PATCH] Implement pack CLI command for matrix packing and layer compaction Adds the `pack` subcommand to `obikmer2` for persisting presence and count matrices in sparse or dense formats. Introduces the `compact_layer` module in `obikrebuild` to merge multiple partition layers into a single layer in-place via atomic operations. Updates dependencies across `obikdump`, `obikmer2`, and `obikrebuild` to resolve new local crate references. Shifts the architecture from multi-layer accumulation to in-place compaction, removing legacy rebuild modules. --- src/Cargo.lock | 7 +- src/obikdump/Cargo.toml | 2 + src/obikindex/src/index/kmer_index.rs | 33 +++- src/obikmer2/Cargo.toml | 1 + src/obikmer2/src/cmd/mod.rs | 1 + src/obikmer2/src/cmd/pack/mod.rs | 71 +++++++ src/obikmer2/src/main.rs | 3 + src/obikrebuild/Cargo.toml | 18 +- src/obikrebuild/src/compact_layer.rs | 113 +++++++++++ src/obikrebuild/src/lib.rs | 36 ++-- src/obikrebuild/src/rebuild.rs | 79 -------- src/obikrebuild/src/rebuild_layer.rs | 270 -------------------------- 12 files changed, 252 insertions(+), 382 deletions(-) create mode 100644 src/obikmer2/src/cmd/pack/mod.rs create mode 100644 src/obikrebuild/src/compact_layer.rs delete mode 100644 src/obikrebuild/src/rebuild.rs delete mode 100644 src/obikrebuild/src/rebuild_layer.rs diff --git a/src/Cargo.lock b/src/Cargo.lock index 9d76b1b7..89441f67 100644 --- a/src/Cargo.lock +++ b/src/Cargo.lock @@ -1523,12 +1523,14 @@ version = "0.1.0" name = "obikdump" version = "0.1.0" dependencies = [ + "obicompactvec", "obidebruinj", "obifastwrite", "obikfilter", "obikidxcache", "obikindex", "obikseq", + "obisys", "rayon", ] @@ -1674,6 +1676,7 @@ dependencies = [ "obikindex", "obikindexer", "obikmerge", + "obikrebuild", "obikselect", "obikseq", "obipipeline", @@ -1736,14 +1739,12 @@ dependencies = [ name = "obikrebuild" version = "0.1.0" dependencies = [ - "obicompactvec", "obidebruinj", "obikfilter", + "obikidxcache", "obikindex", "obikindexer", - "obikmerge", "obikseq", - "obiskio", "obisys", "tracing", ] diff --git a/src/obikdump/Cargo.toml b/src/obikdump/Cargo.toml index 802ede05..173519fa 100644 --- a/src/obikdump/Cargo.toml +++ b/src/obikdump/Cargo.toml @@ -10,4 +10,6 @@ obikidxcache = { path = "../obikidxcache" } obikseq = { path = "../obikseq" } obidebruinj = { path = "../obidebruinj" } obifastwrite = { path = "../obifastwrite" } +obicompactvec = { path = "../obicompactvec" } +obisys = { path = "../obisys" } rayon = "1" diff --git a/src/obikindex/src/index/kmer_index.rs b/src/obikindex/src/index/kmer_index.rs index 4956fab5..06c0d6f2 100644 --- a/src/obikindex/src/index/kmer_index.rs +++ b/src/obikindex/src/index/kmer_index.rs @@ -245,14 +245,25 @@ impl KmerIndex { /// Reduces per-query file-open overhead from O(n_genomes) to O(1) per partition. /// Column files are kept in place; packed files take priority when opening. /// - /// If `sparse` is set, presence matrices go one step further, from the - /// dense `.pbmx` form into `obicompactvec::PersistentSparseBitMatrix`'s - /// on-disk format (see `DevDocMD/architecture/siblings.md`) — count - /// matrices are unaffected, sparse count matrices aren't implemented - /// (see the sparse-matrix design plan's "Explicitly deferred"). - /// TODO: To be moved somewhere else - pub(crate) fn pack_matrices(&self, sparse: bool) -> OKIResult<()> { - use obicompactvec::{pack_bit_matrix, pack_compact_int_matrix, pack_sparse_bit_matrix}; + /// If `sparse` is set, both matrix kinds go one step further, from their + /// dense single-file form into `obicompactvec`'s sparse, deduplicated + /// on-disk formats (`PersistentSparseBitMatrix` for presence, + /// `PersistentSparseCompactIntMatrix` for counts) — smaller and faster + /// for single-row access on real, sparse data; column-oriented access + /// (`--metric` distance matrices) is much slower on the sparse format. + /// + /// Public — and not moved out of this crate, despite covering the same + /// "maintenance operation on an already-built index" territory as + /// `obikrebuild` — because `finalize_indexed` calls it internally as the + /// shared tail of every index-building algorithm (`select`/`filter`/ + /// `merge`), and `obikindex` cannot depend on `obikrebuild` (wrong + /// direction). Standalone re-packing (`obikmer2 pack`) calls this same + /// method directly. + pub fn pack_matrices(&self, sparse: bool) -> OKIResult<()> { + use obicompactvec::{ + pack_bit_matrix, pack_compact_int_matrix, pack_sparse_bit_matrix, + pack_sparse_compact_int_matrix, + }; let n = self.n_partitions(); let order: Vec = (0..n).collect(); @@ -277,7 +288,11 @@ impl KmerIndex { } } if counts_dir.exists() { - pack_compact_int_matrix(&counts_dir).map_err(OKIError::Io)?; + if sparse { + pack_sparse_compact_int_matrix(&counts_dir).map_err(OKIError::Io)?; + } else { + pack_compact_int_matrix(&counts_dir).map_err(OKIError::Io)?; + } } } Ok(()) diff --git a/src/obikmer2/Cargo.toml b/src/obikmer2/Cargo.toml index a8445e7d..9cc1f074 100644 --- a/src/obikmer2/Cargo.toml +++ b/src/obikmer2/Cargo.toml @@ -19,6 +19,7 @@ obikmerge = { path = "../obikmerge" } obikfilter = { path = "../obikfilter" } obikselect = { path = "../obikselect" } obikdump = { path = "../obikdump" } +obikrebuild = { path = "../obikrebuild" } obifastwrite = { path = "../obifastwrite" } obiskbuilder = { path = "../obiskbuilder" } clap = { version = "4", features = ["derive"] } diff --git a/src/obikmer2/src/cmd/mod.rs b/src/obikmer2/src/cmd/mod.rs index 58081aa7..0baca15c 100644 --- a/src/obikmer2/src/cmd/mod.rs +++ b/src/obikmer2/src/cmd/mod.rs @@ -4,6 +4,7 @@ pub mod estimate; pub mod filter; pub mod index; pub mod merge; +pub mod pack; mod predicate; pub mod select; pub mod superkmer; diff --git a/src/obikmer2/src/cmd/pack/mod.rs b/src/obikmer2/src/cmd/pack/mod.rs new file mode 100644 index 00000000..1dbcac3a --- /dev/null +++ b/src/obikmer2/src/cmd/pack/mod.rs @@ -0,0 +1,71 @@ +use std::path::PathBuf; + +use clap::Args; +use obikindex::KmerIndex; +use obikrebuild::IndexCompact; +use obisys::{Reporter, Stage, progress_bar}; +use tracing::info; + +#[derive(Args)] +pub struct PackArgs { + /// Index directory to pack + pub index: PathBuf, + + /// Compact every partition's accumulated layers into one before packing + /// — undoes the multi-layer stopgap `merge` leaves behind. + #[arg(long, default_value_t = false)] + pub compact_layers: bool, + + /// Pack presence and count matrices into the dense on-disk format instead + /// of the default sparse, deduplicated one. Dense is faster for + /// column-oriented access (`--metric` distance matrices); sparse is + /// smaller and faster for single-row access on real, sparse data. + #[arg(long, default_value_t = false)] + pub dense: bool, +} + +pub fn run(args: PackArgs) { + // Modifies the index in place; acquired before opening so a concurrent + // writer can't slip in between the open and the pack below. + let _lock = obisys::DirLock::acquire(&args.index).unwrap_or_else(|e| { + eprintln!("error locking index directory {}: {e}", args.index.display()); + std::process::exit(1); + }); + + let idx = KmerIndex::open(&args.index).unwrap_or_else(|e| { + eprintln!("error opening index: {e}"); + std::process::exit(1); + }); + + let n_genomes = idx.meta().genomes().unwrap_or_else(|e| { + eprintln!("error reading index metadata: {e}"); + std::process::exit(1); + }).len(); + info!( + "pack: {} partition(s), {} genome(s)", + idx.n_partitions(), + n_genomes, + ); + + let mut rep = Reporter::new(); + + if args.compact_layers { + let t = Stage::start("compact layers"); + let pb = progress_bar("compact", idx.n_partitions() as u64, "partitions"); + idx.compact_layers(|| pb.inc(1)).unwrap_or_else(|e| { + eprintln!("compact error: {e}"); + std::process::exit(1); + }); + pb.finish_and_clear(); + rep.push(t.stop()); + } + + let t = Stage::start("pack"); + idx.pack_matrices(!args.dense).unwrap_or_else(|e| { + eprintln!("pack error: {e}"); + std::process::exit(1); + }); + rep.push(t.stop()); + + rep.print(); +} diff --git a/src/obikmer2/src/main.rs b/src/obikmer2/src/main.rs index a8153de1..7dcb2d4e 100644 --- a/src/obikmer2/src/main.rs +++ b/src/obikmer2/src/main.rs @@ -27,6 +27,8 @@ enum Commands { Dump(cmd::dump::DumpArgs), /// Assemble an index's kmers into unitigs and write them as FASTA Unitig(cmd::unitig::UnitigArgs), + /// Pack an index's matrices into single-file, sparse format by default (--dense to opt out), in place + Pack(cmd::pack::PackArgs), /// Estimate approximate-evidence false-positive rates for given parameters Estimate(cmd::estimate::EstimateArgs), /// Read/write genome metadata (CSV) on an already-built index @@ -50,6 +52,7 @@ fn main() { Commands::Select(args) => cmd::select::run(args), Commands::Dump(args) => cmd::dump::run(args), Commands::Unitig(args) => cmd::unitig::run(args), + Commands::Pack(args) => cmd::pack::run(args), Commands::Estimate(args) => cmd::estimate::run(args), Commands::Annotate(args) => cmd::annotate::run(args), } diff --git a/src/obikrebuild/Cargo.toml b/src/obikrebuild/Cargo.toml index 50c39c6c..3be610a9 100644 --- a/src/obikrebuild/Cargo.toml +++ b/src/obikrebuild/Cargo.toml @@ -4,13 +4,11 @@ version = "0.1.0" edition = "2024" [dependencies] -obikindex = { path = "../obikindex" } -obikindexer = { path = "../obikindexer" } -obikfilter = { path = "../obikfilter" } -obikmerge = { path = "../obikmerge" } -obicompactvec = { path = "../obicompactvec" } -obiskio = { path = "../obiskio" } -obikseq = { path = "../obikseq" } -obidebruinj = { path = "../obidebruinj" } -obisys = { path = "../obisys" } -tracing = "0.1.44" +obikindex = { path = "../obikindex" } +obikindexer = { path = "../obikindexer" } +obikfilter = { path = "../obikfilter" } +obikidxcache = { path = "../obikidxcache" } +obikseq = { path = "../obikseq" } +obidebruinj = { path = "../obidebruinj" } +obisys = { path = "../obisys" } +tracing = "0.1.44" diff --git a/src/obikrebuild/src/compact_layer.rs b/src/obikrebuild/src/compact_layer.rs new file mode 100644 index 00000000..5d603dba --- /dev/null +++ b/src/obikrebuild/src/compact_layer.rs @@ -0,0 +1,113 @@ +//! Compacting a partition's multiple layers into one, in place. The +//! multi-layer structure is a stopgap that lets `obikmerge` add genomes +//! without a full rebuild; this collapses that accumulated history back +//! into a single layer. Distinct from `obikfilter::Filter`: no k-mer is +//! removed here, and every layer of a partition merges into *one*, rather +//! than one output layer per source layer. `KmerIndex` is foreign, so this +//! is an extension trait. + +use std::collections::HashMap; +use std::fs; + +use obidebruinj::GraphDeBruijn; +use obikfilter::FilteredPartitionIter; +use obikidxcache::index_cache::IndexCache; +use obikindex::layer::{IndexMode, KmerLayer, TypedLayer}; +use obikindex::{KmerIndex, OKIError, OKIResult}; +use obikindexer::write_graph_as_unitigs; +use obikseq::CanonicalKmer; +use obisys::PartitionRunner; + +pub trait IndexCompact { + /// Compact every partition with more than one layer into a single + /// layer, in place. Partitions already at one layer are left untouched. + fn compact_layers(&self, on_partition: impl Fn() + Send + Sync) -> OKIResult<()>; +} + +impl IndexCompact for KmerIndex { + fn compact_layers(&self, on_partition: impl Fn() + Send + Sync) -> OKIResult<()> { + let n_genomes = self.meta().genomes().map_err(OKIError::Io)?.len(); + let use_counts = self.meta().config.with_counts; + let block_bits = self.block_bits(); + let mode = self.evidence_mode(); + let n_partitions = self.n_partitions(); + let order: Vec = (0..n_partitions).collect(); + + PartitionRunner::new().run( + &order, + |i| compact_partition(self, i, use_counts, n_genomes, block_bits, mode), + |_, _, _| on_partition(), + )?; + + Ok(()) + } +} + +/// Merge every layer of partition `i` into a single one, replacing the +/// existing layers in place. No-op if the partition already has one layer +/// (or none). +fn compact_partition( + idx: &KmerIndex, + i: usize, + use_counts: bool, + n_genomes: usize, + block_bits: u8, + mode: &IndexMode, +) -> OKIResult<()> { + let partition = idx.partition(i)?; + let n_layers = partition.n_layers(); + if n_layers <= 1 { + return Ok(()); + } + + // ── Read: every layer's kmers, combined — no filter, nothing removed ── + let mut survivors: HashMap> = HashMap::new(); + { + let cache = IndexCache::new(idx, Some(vec![i])); + cache.iter_partition_kmers(i, use_counts, n_genomes, &[], |kmer, row| { + survivors.insert(kmer, row); + true + })?; + } + + // ── Build the combined layer in a scratch directory ──────────────────── + let tmp_dir = partition.dir().join("layer_compact_tmp"); + if tmp_dir.exists() { + fs::remove_dir_all(&tmp_dir).map_err(OKIError::Io)?; + } + fs::create_dir_all(&tmp_dir).map_err(OKIError::Io)?; + + let mut g = GraphDeBruijn::new(); + for &kmer in survivors.keys() { + g.push(kmer); + } + write_graph_as_unitigs(g, &tmp_dir)?; + + TypedLayer::<()>::build_with_matrix( + &tmp_dir, + block_bits, + mode, + !use_counts, + n_genomes, + |kmer| survivors.get(&kmer).cloned().expect("kmer pushed into the graph must be a survivor"), + )?; + + // ── Swap: drop every existing layer, install the combined one as layer 0 ─ + for l in 0..n_layers { + let layer_dir = partition.layer_dir(l)?; + fs::remove_dir_all(&layer_dir).map_err(OKIError::Io)?; + } + let layer0_dir = KmerLayer::new(&partition, 0) + .create() + .map_err(OKIError::Io)? + .dir() + .to_path_buf(); + for entry in fs::read_dir(&tmp_dir).map_err(OKIError::Io)? { + let entry = entry.map_err(OKIError::Io)?; + let dest = layer0_dir.join(entry.file_name()); + fs::rename(entry.path(), dest).map_err(OKIError::Io)?; + } + fs::remove_dir_all(&tmp_dir).map_err(OKIError::Io)?; + + Ok(()) +} diff --git a/src/obikrebuild/src/lib.rs b/src/obikrebuild/src/lib.rs index 71c8a41f..3e6658e8 100644 --- a/src/obikrebuild/src/lib.rs +++ b/src/obikrebuild/src/lib.rs @@ -1,12 +1,26 @@ -//! Rebuilding an `obikindex::KmerIndex`: compacting a source index into a -//! new single-layer index while applying k-mer filters (`rebuild`), and -//! reindexing an existing index's evidence representation in place -//! (`reindex`, `Exact`/`Approx`/`Hybrid`). Both are maintenance/ -//! transformation operations on an already-built index, not part of the -//! `Index { Partition { Layer } }` data model itself — kept out of -//! `obikindex` so that crate doesn't grow every algorithm's own -//! dependencies (`obikfilter`, `obikmerge`, `obikindexer`). +//! Maintenance/transformation operations on an already-built +//! `obikindex::KmerIndex`, not part of the `Index { Partition { Layer } }` +//! data model itself — kept out of `obikindex` so that crate doesn't grow +//! every algorithm's own dependencies (`obikfilter`, `obikidxcache`, +//! `obikindexer`). +//! +//! [`compact_layer`] collapses a partition's accumulated layers (the +//! stopgap `obikmerge` uses to add genomes without a full rebuild) back +//! into one, in place. `reindex` (reindexing an existing index's evidence +//! representation in place — `Exact`/`Approx`/`Hybrid`) is temporarily +//! disabled: it predates this session's refactor and still reaches +//! `obikindex` internals (`root_path`/`meta` fields, `index_dir`/ +//! `layer_dir` methods) that are no longer accessible from outside that +//! crate. Fixing it is separate, future work. +//! +//! Removing k-mers (filtering) and repacking matrix files are *not* here: +//! the former lives in `obikfilter::Filter` (a fresh index at a new +//! output, not in place); the latter is `obikindex::KmerIndex:: +//! pack_matrices` itself, since `finalize_indexed` — the shared tail of +//! every index-building algorithm — calls it internally and `obikindex` +//! cannot depend back on this crate. -mod rebuild; -mod rebuild_layer; -mod reindex; +mod compact_layer; +// mod reindex; // disabled — see module doc above. + +pub use compact_layer::IndexCompact; diff --git a/src/obikrebuild/src/rebuild.rs b/src/obikrebuild/src/rebuild.rs deleted file mode 100644 index c7945328..00000000 --- a/src/obikrebuild/src/rebuild.rs +++ /dev/null @@ -1,79 +0,0 @@ -use std::path::Path; - -use obikindex::IndexBuilder; -use obikfilter::KmerFilter; -use obikmerge::MergeMode; -use obisys::{Reporter, Stage, progress_bar}; -use tracing::info; - -use obikindex::{OKIError, OKIResult}; -use obikindex::KmerIndex; -use obikindex::IndexState; -use obisys::PartitionRunner; - -impl KmerIndex { - /// Rebuild `src` into a new compact single-layer index at `output`. - /// - /// Only k-mers whose per-genome row passes every filter in `filters` are - /// written. If `filters` is empty every k-mer is kept (pure compaction). - /// - /// `mode` controls whether the output stores counts or presence/absence. - /// A count source may be rebuilt in presence mode; a presence source - /// cannot be rebuilt in count mode. - pub fn rebuild>( - output: P, - src: &KmerIndex, - filters: &[Box], - mode: MergeMode, - force: bool, - rep: &mut Reporter, - ) -> OKIResult { - let output = output.as_ref(); - - if src.state()? != IndexState::Indexed { - return Err(OKIError::NotIndexed(src.root_path.clone())); - } - - if mode == MergeMode::Count && !src.meta.config.with_counts { - return Err(OKIError::InvalidInput( - "cannot rebuild in count mode from a presence-only source index".into(), - )); - } - - KmerIndex::clear_output_for_create(output, force)?; - - // ── Create output directory + metadata ──────────────────────────────── - let mut config = src.meta.config.clone(); - config.with_counts = mode == MergeMode::Count; - let genomes = src.genomes()?; - - let n_genomes = genomes.len(); - let n_partitions = src.n_partitions(); - let block_bits = config.block_bits; - - // ── Create an empty destination KmerPartition ───────────────────────── - let dst_partition = KmerIndex::create_skeleton(output, config, genomes)?; - - info!( - "rebuild: {} partition(s), {} genome(s), mode={:?}", - n_partitions, n_genomes, mode, - ); - - let t = Stage::start("rebuild"); - let pb = progress_bar("rebuild", n_partitions as u64, "partitions"); - - let order: Vec = (0..n_partitions).collect(); - let runner = PartitionRunner::new(); - runner.run( - &order, - |i| dst_partition.rebuild_partition(src, i, filters, mode, n_genomes, block_bits), - |_, _, _| { pb.inc(1); }, - ).map_err(OKIError::Partition)?; - - pb.finish_and_clear(); - - rep.push(t.stop()); - - KmerIndex::finalize_indexed(output, rep, false) - } -} diff --git a/src/obikrebuild/src/rebuild_layer.rs b/src/obikrebuild/src/rebuild_layer.rs deleted file mode 100644 index a3e61d7b..00000000 --- a/src/obikrebuild/src/rebuild_layer.rs +++ /dev/null @@ -1,270 +0,0 @@ -use std::path::Path; - -use obicompactvec::{ - FilterMask, PersistentBitMatrixBuilder, PersistentBitVecBuilder, - PersistentCompactIntMatrixBuilder, PersistentCompactIntVecBuilder, eval_filter_mask, -}; -use obidebruinj::GraphDeBruijn; -use obikseq::CanonicalKmer; -use obikindex::layer::meta::PartitionMeta; -use obikindex::layer::{IndexMode, MphfLayer}; -use obikindex::layer::utils::layer_dir; -use obiskio::UnitigFileReader; -use obikindex::{OKIError, OKIResult}; - -use obikindex::load_meta; -use obikfilter::KmerFilter; -use obikindexer::materialize_layer; -use obikmerge::{MergeMode, SrcLayerData}; -use obikindex::KmerIndex; - -// ── Builders — pair matrix builder + column builders for one mode ───────────── - -enum Builders { - Presence(PersistentBitMatrixBuilder, Vec), - Count( - PersistentCompactIntMatrixBuilder, - Vec, - ), -} - -impl Builders { - fn new(mode: MergeMode, n: usize, dir: &Path, n_genomes: usize) -> OKIResult { - match mode { - MergeMode::Presence => { - let mut mat = PersistentBitMatrixBuilder::new(n, dir).map_err(OKIError::Io)?; - let mut cols = Vec::with_capacity(n_genomes); - for _ in 0..n_genomes { - cols.push(mat.add_col().map_err(OKIError::Io)?); - } - Ok(Builders::Presence(mat, cols)) - } - MergeMode::Count => { - let mut mat = - PersistentCompactIntMatrixBuilder::new(n, dir).map_err(OKIError::Io)?; - let mut cols = Vec::with_capacity(n_genomes); - for _ in 0..n_genomes { - cols.push(mat.add_col().map_err(OKIError::Io)?); - } - Ok(Builders::Count(mat, cols)) - } - } - } - - fn set_val(&mut self, col: usize, slot: usize, value: u32) { - match self { - Builders::Presence(_, cols) => cols[col].set(slot, value > 0), - Builders::Count(_, cols) => cols[col].set(slot, value), - } - } - - fn close(self) -> OKIResult<()> { - match self { - Builders::Presence(mat, cols) => { - for b in cols { - b.close().map_err(OKIError::Io)?; - } - mat.close().map_err(OKIError::Io) - } - Builders::Count(mat, cols) => { - for b in cols { - b.close().map_err(OKIError::Io)?; - } - mat.close().map_err(OKIError::Io) - } - } - } -} - -// ── try_compute_combined_mask ───────────────────────────────────────────────── - -/// Build a per-slot `TempBitVec` mask from `filters` using column operations -/// on the source matrix — no per-kmer MPHF lookup or row read needed. -/// -/// Returns `Some(mask)` when every filter in `filters` can express itself as -/// a [`FilterMask`] expression. Returns `None` when any filter requires -/// row-level inspection (fall back to `passes_all`). -fn try_compute_combined_mask( - filters: &[Box], - src_data: &SrcLayerData, - n_genomes: usize, -) -> OKIResult> { - if filters.is_empty() { - return Ok(None); - } - let mut exprs: Vec = Vec::with_capacity(filters.len()); - for f in filters { - match f.column_mask_expr(n_genomes) { - Some(expr) => exprs.push(expr), - None => return Ok(None), - } - } - let combined = FilterMask::And(exprs); - let n = src_data.n_slots(); - let mask = src_data - .with_matrix(|mat| eval_filter_mask(&combined, mat, n)) - .map_err(OKIError::Io)?; - Ok(Some(mask)) -} - -// ── iter_src_kmers_masked (pass 1) ──────────────────────────────────────────── - -/// Iterate all passing kmers in `src_index_dir`, yielding only the kmer value. -/// -/// When all filters can be expressed as column operations, a per-slot mask is -/// computed once per layer and used for O(1) slot-check per kmer instead of a -/// full row read. Falls back to row-level `passes_all` otherwise. -fn iter_src_kmers_masked( - src_index_dir: &Path, - mode: MergeMode, - n_genomes: usize, - filters: &[Box], - mut cb: impl FnMut(CanonicalKmer), -) -> OKIResult<()> { - let src_meta = load_meta(src_index_dir)?; - for l in 0..src_meta.n_layers { - let src_layer_dir = layer_dir(src_index_dir, l); - let unitigs_path = src_layer_dir.join("unitigs.bin"); - if !unitigs_path.exists() { - continue; - } - - let src_data = SrcLayerData::open(&src_layer_dir, mode)?; - let mask = try_compute_combined_mask(filters, &src_data, n_genomes)?; - let reader = UnitigFileReader::open_sequential(&unitigs_path)?; - - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - let slot = src_data.slot(kmer); - let passes = match &mask { - Some(m) => m.get(slot), - None => { - let row = src_data.fill_row_by_slot(slot, n_genomes); - filters.iter().all(|f| f.passes(kmer, &row, n_genomes)) - } - }; - if passes { - cb(kmer); - } - } - } - Ok(()) -} - -// ── iter_src_layers (pass 2) ────────────────────────────────────────────────── - -/// Iterate all passing kmers in `src_index_dir`, yielding `(kmer, row)`. -/// -/// When the slot mask is available, skips the row read for filtered-out slots. -fn iter_src_layers( - src_index_dir: &Path, - mode: MergeMode, - n_genomes: usize, - filters: &[Box], - mut cb: impl FnMut(CanonicalKmer, Box<[u32]>), -) -> OKIResult<()> { - let src_meta = load_meta(src_index_dir)?; - for l in 0..src_meta.n_layers { - let src_layer_dir = layer_dir(src_index_dir, l); - let unitigs_path = src_layer_dir.join("unitigs.bin"); - if !unitigs_path.exists() { - continue; - } - - let src_data = SrcLayerData::open(&src_layer_dir, mode)?; - let mask = try_compute_combined_mask(filters, &src_data, n_genomes)?; - let reader = UnitigFileReader::open_sequential(&unitigs_path)?; - - for (kmer, _, _) in reader.iter_indexed_canonical_kmers() { - let slot = src_data.slot(kmer); - if let Some(ref m) = mask { - if !m.get(slot) { - continue; - } - let row = src_data.fill_row_by_slot(slot, n_genomes); - cb(kmer, row.into_boxed_slice()); - } else { - let row = src_data.fill_row_by_slot(slot, n_genomes); - if filters.iter().all(|f| f.passes(kmer, &row, n_genomes)) { - cb(kmer, row.into_boxed_slice()); - } - } - } - } - Ok(()) -} - -// ── KmerPartition::rebuild_partition ───────────────────────────────────────── - -impl KmerIndex { - /// Rebuild partition `i` from `src` into `self` (an empty destination partition). - /// - /// Only k-mers whose per-genome row passes all `filters` are written. - /// The output is a single-layer index — regardless of how many layers the - /// source has. - /// - /// `n_genomes` is the number of genome columns in the source (and destination). - pub fn rebuild_partition( - &self, - src: &KmerIndex, - i: usize, - filters: &[Box], - mode: MergeMode, - n_genomes: usize, - block_bits: u8, - ) -> OKIResult<()> { - let src_index_dir = src.index_dir(i); - if !src_index_dir.exists() { - return Ok(()); - } - - let src_meta = load_meta(&src_index_dir)?; - if src_meta.n_layers == 0 { - return Ok(()); - } - - // ── Pass 1: collect filtered kmers into de Bruijn graph ─────────────── - let mut g = GraphDeBruijn::new(); - iter_src_kmers_masked(&src_index_dir, mode, n_genomes, filters, |kmer| { - g.push(kmer); - })?; - - if g.len() == 0 { - return Ok(()); - } - - // ── Build MPHF in dst layer_0 ───────────────────────────────────────── - let dst_index_dir = self.index_dir(i); - let dst_layer_dir = self.layer_dir(i, 0); - - let n_new = materialize_layer(g, &dst_layer_dir, block_bits, &IndexMode::Exact)?; - let dst_mphf = MphfLayer::open(&dst_layer_dir)?; - - // ── Prepare matrix builders (one column per genome) ─────────────────── - let data_dir = match mode { - MergeMode::Presence => dst_layer_dir.join("presence"), - MergeMode::Count => dst_layer_dir.join("counts"), - }; - std::fs::create_dir_all(&data_dir)?; - let mut builders = Builders::new(mode, n_new, &data_dir, n_genomes)?; - - // ── Pass 2: fill builders ───────────────────────────────────────────── - iter_src_layers(&src_index_dir, mode, n_genomes, filters, |kmer, row| { - if let Some(slot) = dst_mphf.find(kmer) { - for (col, &value) in row.iter().enumerate() { - builders.set_val(col, slot, value); - } - } - })?; - - // ── Close builders and write metadata ───────────────────────────────── - builders.close()?; - - PartitionMeta { - n_layers: 1, - mode: IndexMode::Exact, - } - .save(&dst_index_dir)?; - - Ok(()) - } -}