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.
This commit is contained in:
Generated
+4
-3
@@ -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",
|
||||
]
|
||||
|
||||
@@ -10,4 +10,6 @@ obikidxcache = { path = "../obikidxcache" }
|
||||
obikseq = { path = "../obikseq" }
|
||||
obidebruinj = { path = "../obidebruinj" }
|
||||
obifastwrite = { path = "../obifastwrite" }
|
||||
obicompactvec = { path = "../obicompactvec" }
|
||||
obisys = { path = "../obisys" }
|
||||
rayon = "1"
|
||||
|
||||
@@ -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<usize> = (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(())
|
||||
|
||||
@@ -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"] }
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
@@ -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),
|
||||
}
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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<usize> = (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<CanonicalKmer, Box<[u32]>> = 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(())
|
||||
}
|
||||
+25
-11
@@ -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;
|
||||
|
||||
@@ -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<P: AsRef<Path>>(
|
||||
output: P,
|
||||
src: &KmerIndex,
|
||||
filters: &[Box<dyn KmerFilter>],
|
||||
mode: MergeMode,
|
||||
force: bool,
|
||||
rep: &mut Reporter,
|
||||
) -> OKIResult<Self> {
|
||||
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<usize> = (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)
|
||||
}
|
||||
}
|
||||
@@ -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<PersistentBitVecBuilder>),
|
||||
Count(
|
||||
PersistentCompactIntMatrixBuilder,
|
||||
Vec<PersistentCompactIntVecBuilder>,
|
||||
),
|
||||
}
|
||||
|
||||
impl Builders {
|
||||
fn new(mode: MergeMode, n: usize, dir: &Path, n_genomes: usize) -> OKIResult<Self> {
|
||||
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<dyn KmerFilter>],
|
||||
src_data: &SrcLayerData,
|
||||
n_genomes: usize,
|
||||
) -> OKIResult<Option<obicompactvec::TempBitVec>> {
|
||||
if filters.is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
let mut exprs: Vec<FilterMask> = 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<dyn KmerFilter>],
|
||||
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<dyn KmerFilter>],
|
||||
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<dyn KmerFilter>],
|
||||
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(())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user