Extract index modules into specialized workspace subcrates
This commit partitions the obikindex crate into multiple focused subcrates (obikfilter, obikmerge, obikquery, obikrebuild, obikselect, obikstats, obikdump, and obikidxcache) to reduce coupling and clarify module boundaries. It standardizes error handling across the workspace using OKIError and OKIResult, updates index APIs to support lazy, disk-backed partition access, and migrates NUMA system utilities to a new obisys crate. All modifications are structural, focusing on dependency graph expansion, import path updates, and API surface reorganization without altering core runtime behavior.
This commit is contained in:
@@ -0,0 +1,15 @@
|
||||
[package]
|
||||
name = "obikmerge"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
obikindex = { path = "../obikindex" }
|
||||
obikindexer = { path = "../obikindexer" }
|
||||
obikfilter = { path = "../obikfilter" }
|
||||
obicompactvec = { path = "../obicompactvec" }
|
||||
obiskio = { path = "../obiskio" }
|
||||
obikseq = { path = "../obikseq" }
|
||||
obipipeline = { path = "../obipipeline" }
|
||||
obisys = { path = "../obisys" }
|
||||
tracing = "0.1.44"
|
||||
@@ -0,0 +1,17 @@
|
||||
//! Merging multiple `obikindex::KmerIndex` sources into one: bootstrapping
|
||||
//! from the first source, then merging each remaining source partition by
|
||||
//! partition (de Bruijn graph union + column fill). A maintenance/
|
||||
//! transformation operation on already-built indexes, not part of the
|
||||
//! `Index { Partition { Layer } }` data model itself — kept out of
|
||||
//! `obikindex` so that crate doesn't grow this algorithm's own
|
||||
//! dependencies (`obipipeline`, `obidebruinj` via `obikindexer`).
|
||||
//!
|
||||
//! [`merge_layer`] holds the per-partition merge primitive
|
||||
//! (`merge_partition`, `SrcLayerData`) — the latter is also reused by
|
||||
//! `obikrebuild::rebuild_layer`, hence its `pub` visibility here.
|
||||
|
||||
mod merge;
|
||||
mod merge_layer;
|
||||
|
||||
pub use merge::*;
|
||||
pub use merge_layer::{MergeMode, SrcLayerData};
|
||||
@@ -0,0 +1,473 @@
|
||||
use std::collections::HashMap;
|
||||
use std::fs;
|
||||
use std::io;
|
||||
use std::path::Path;
|
||||
|
||||
use obikindex::IndexBuilder;
|
||||
|
||||
use obisys::{Reporter, Stage, progress_bar, spinner};
|
||||
use tracing::{debug, info};
|
||||
|
||||
use obikindex::layer::IndexMode;
|
||||
|
||||
use obikindex::{OKIError, OKIResult};
|
||||
use obikindex::KmerIndex;
|
||||
use obikindex::{GenomeInfo, IndexMeta};
|
||||
use obikindex::IndexState;
|
||||
use obisys::PartitionRunner;
|
||||
|
||||
pub use crate::merge_layer::MergeMode;
|
||||
|
||||
// ── per-partition diagnostic record ──────────────────────────────────────────
|
||||
|
||||
#[derive(Debug)]
|
||||
struct PartStat {
|
||||
id: usize,
|
||||
unitig_bytes: u64,
|
||||
g_len: usize,
|
||||
}
|
||||
|
||||
// ── main merge entry point ────────────────────────────────────────────────────
|
||||
|
||||
impl KmerIndex {
|
||||
pub fn merge<P: AsRef<Path>>(
|
||||
output: P,
|
||||
sources: &[&KmerIndex],
|
||||
mode: MergeMode,
|
||||
force: bool,
|
||||
rename_duplicates: bool,
|
||||
budget_fraction: f64,
|
||||
rep: &mut Reporter,
|
||||
) -> OKIResult<Self> {
|
||||
let output = output.as_ref();
|
||||
|
||||
if sources.is_empty() {
|
||||
return Err(OKIError::Io(io::Error::new(
|
||||
io::ErrorKind::InvalidInput,
|
||||
"merge requires at least one source index",
|
||||
)));
|
||||
}
|
||||
|
||||
// ── Validate config compatibility ─────────────────────────────────────
|
||||
let ref0 = sources[0];
|
||||
for src in sources {
|
||||
if src.state()? != IndexState::Indexed {
|
||||
return Err(OKIError::NotIndexed(src.root_path.clone()));
|
||||
}
|
||||
if src.kmer_size() != ref0.kmer_size()
|
||||
|| src.minimizer_size() != ref0.minimizer_size()
|
||||
|| src.n_partitions() != ref0.n_partitions()
|
||||
{
|
||||
return Err(OKIError::IncompatibleConfig);
|
||||
}
|
||||
if mode == MergeMode::Count && !src.meta.config.with_counts {
|
||||
return Err(OKIError::MismatchedMode);
|
||||
}
|
||||
}
|
||||
|
||||
// Read each source's genome list once — `IndexMeta::genomes` is a
|
||||
// fresh disk read every call, so cache it rather than re-reading it
|
||||
// repeatedly through the rest of this function.
|
||||
let src_genomes: Vec<Vec<GenomeInfo>> =
|
||||
sources.iter().map(|s| s.genomes()).collect::<OKIResult<_>>()?;
|
||||
|
||||
// ── Log source characteristics and choose base ────────────────────────
|
||||
let mode_str = if mode == MergeMode::Presence {
|
||||
"presence"
|
||||
} else {
|
||||
"count"
|
||||
};
|
||||
info!(
|
||||
"merge: {} source(s), smer-size={}, mode={}",
|
||||
sources.len(),
|
||||
sources[0].kmer_size(),
|
||||
mode_str,
|
||||
);
|
||||
for (i, src) in sources.iter().enumerate() {
|
||||
let genome_str = if src_genomes[i].len() == 1 {
|
||||
"mono-genome".to_string()
|
||||
} else {
|
||||
format!("{} genomes", src_genomes[i].len())
|
||||
};
|
||||
let trivial_str = if is_trivial(&src_genomes[i], mode) {
|
||||
" [trivial: no data approximation]"
|
||||
} else {
|
||||
""
|
||||
};
|
||||
info!(
|
||||
" [{}] {} — {}, {}, {}{}",
|
||||
i,
|
||||
src.root_path.display(),
|
||||
format_evidence(&src.meta.config.evidence),
|
||||
genome_str,
|
||||
mode_str,
|
||||
trivial_str,
|
||||
);
|
||||
}
|
||||
|
||||
let base_idx = choose_base(sources, &src_genomes, mode);
|
||||
let needs_approx = sources.iter().enumerate().any(|(i, src)| {
|
||||
!is_trivial(&src_genomes[i], mode)
|
||||
&& matches!(
|
||||
src.meta.config.evidence,
|
||||
IndexMode::Approx { .. } | IndexMode::Hybrid { .. }
|
||||
)
|
||||
});
|
||||
info!(
|
||||
"output evidence: {} ({}base: [{}] {})",
|
||||
format_evidence(&sources[base_idx].meta.config.evidence),
|
||||
if needs_approx {
|
||||
"forced approx — "
|
||||
} else {
|
||||
""
|
||||
},
|
||||
base_idx,
|
||||
sources[base_idx].root_path.display(),
|
||||
);
|
||||
|
||||
let mut ordered: Vec<&KmerIndex> = Vec::with_capacity(sources.len());
|
||||
let mut ordered_genomes: Vec<Vec<GenomeInfo>> = Vec::with_capacity(sources.len());
|
||||
ordered.push(sources[base_idx]);
|
||||
ordered_genomes.push(src_genomes[base_idx].clone());
|
||||
for (i, &src) in sources.iter().enumerate() {
|
||||
if i != base_idx {
|
||||
ordered.push(src);
|
||||
ordered_genomes.push(src_genomes[i].clone());
|
||||
}
|
||||
}
|
||||
let sources: &[&KmerIndex] = &ordered;
|
||||
let src_genomes = ordered_genomes;
|
||||
let evidence = sources[0].meta.config.evidence.clone();
|
||||
|
||||
// ── Compute final genome labels ────────────────────────────────────────
|
||||
let (source_labels, all_genomes) = compute_labels(&src_genomes, rename_duplicates)?;
|
||||
|
||||
// ── Prepare output directory ──────────────────────────────────────────
|
||||
KmerIndex::clear_output_for_create(output, force)?;
|
||||
|
||||
// ── Bootstrap: copy first source to output ────────────────────────────
|
||||
info!(
|
||||
"bootstrap: copying {} → {} ({} genome(s))",
|
||||
sources[0].root_path.display(),
|
||||
output.display(),
|
||||
src_genomes[0].len(),
|
||||
);
|
||||
let t = Stage::start("bootstrap");
|
||||
let pb = spinner("bootstrap");
|
||||
pb.set_message("copying index …");
|
||||
copy_dir_all(&sources[0].root_path, output)?;
|
||||
|
||||
let dst_meta = IndexMeta::open_at(output).map_err(OKIError::Io)?;
|
||||
let mut config = dst_meta.config.clone();
|
||||
config.with_counts = mode == MergeMode::Count;
|
||||
config.evidence = evidence.clone();
|
||||
dst_meta.rewrite_config(config, all_genomes).map_err(OKIError::Io)?;
|
||||
|
||||
if mode == MergeMode::Presence {
|
||||
remove_dirs_named(output, "counts")?;
|
||||
}
|
||||
pb.finish_and_clear();
|
||||
rep.push(t.stop());
|
||||
|
||||
// ── Rebuild spectrums ─────────────────────────────────────────────────
|
||||
info!("rebuilding spectrums for {} source(s)", sources.len());
|
||||
let t = Stage::start("spectrums");
|
||||
let pb = spinner("spectrums");
|
||||
pb.set_message("copying …");
|
||||
let spectrums_dir = output.join("spectrums");
|
||||
if spectrums_dir.exists() {
|
||||
fs::remove_dir_all(&spectrums_dir)?;
|
||||
}
|
||||
for ((src, new_labels), genomes) in sources.iter().zip(&source_labels).zip(&src_genomes) {
|
||||
let old_labels: Vec<String> = genomes.iter().map(|g| g.label.clone()).collect();
|
||||
copy_spectrums(&src.root_path, output, &old_labels, new_labels)?;
|
||||
}
|
||||
pb.finish_and_clear();
|
||||
rep.push(t.stop());
|
||||
|
||||
// ── Open destination ──────────────────────────────────────────────────
|
||||
let dst = KmerIndex::open(output)?;
|
||||
let n_partitions = dst.n_partitions();
|
||||
let n_dst_genomes = src_genomes[0].len();
|
||||
|
||||
// ── Merge partitions ──────────────────────────────────────────────────
|
||||
let remaining_sources: Vec<&KmerIndex> = sources[1..].to_vec();
|
||||
if !remaining_sources.is_empty() {
|
||||
let n_src_genomes: usize = src_genomes[1..].iter().map(|g| g.len()).sum();
|
||||
info!(
|
||||
"merging {} partition(s) × {} additional source genome(s) into {} destination genome(s)",
|
||||
n_partitions, n_src_genomes, n_dst_genomes,
|
||||
);
|
||||
let t = Stage::start("merge_partitions");
|
||||
let pb = progress_bar("merge", n_partitions as u64, "partitions");
|
||||
|
||||
let block_bits = dst.meta.config.block_bits;
|
||||
|
||||
// Pre-build source list once (avoid rebuilding per partition)
|
||||
let srcs: Vec<(&KmerIndex, usize)> = remaining_sources
|
||||
.iter()
|
||||
.zip(&src_genomes[1..])
|
||||
.map(|(s, g)| (*s, g.len()))
|
||||
.collect();
|
||||
|
||||
// Per-partition unitig byte sizes across remaining sources (stat() only)
|
||||
let partition_sizes: Vec<u64> = (0..n_partitions)
|
||||
.map(|i| {
|
||||
remaining_sources
|
||||
.iter()
|
||||
.map(|s| partition_unitig_bytes(s, i))
|
||||
.sum()
|
||||
})
|
||||
.collect();
|
||||
|
||||
// LFD sort: largest partition first
|
||||
let mut order: Vec<usize> = (0..n_partitions).collect();
|
||||
order.sort_unstable_by_key(|&i| std::cmp::Reverse(partition_sizes[i]));
|
||||
|
||||
let _ = budget_fraction; // kept in signature for CLI compatibility
|
||||
|
||||
// Shadow as references so closures can capture them by copy.
|
||||
let srcs = &srcs;
|
||||
let evidence = &evidence;
|
||||
|
||||
let runner = PartitionRunner::new();
|
||||
let mut part_stats: Vec<PartStat> = Vec::with_capacity(n_partitions);
|
||||
|
||||
runner
|
||||
.run(
|
||||
&order,
|
||||
|i| {
|
||||
dst.merge_partition(
|
||||
i,
|
||||
srcs,
|
||||
mode,
|
||||
n_dst_genomes,
|
||||
block_bits,
|
||||
evidence,
|
||||
)
|
||||
},
|
||||
|i, g_len, dur| {
|
||||
pb.inc(1);
|
||||
debug!(
|
||||
"partition {i}: done in {:.1}s — {} new kmers",
|
||||
dur.as_secs_f64(),
|
||||
g_len,
|
||||
);
|
||||
part_stats.push(PartStat {
|
||||
id: i,
|
||||
unitig_bytes: partition_sizes[i],
|
||||
g_len,
|
||||
});
|
||||
},
|
||||
)
|
||||
.map_err(OKIError::Partition)?;
|
||||
|
||||
pb.finish_and_clear();
|
||||
|
||||
// ── Diagnostic report ─────────────────────────────────────────────
|
||||
print_merge_partition_report(&part_stats, runner.max_workers());
|
||||
|
||||
rep.push(t.stop());
|
||||
}
|
||||
|
||||
// ── Pack matrices after merge ─────────────────────────────────────────
|
||||
{
|
||||
let t = Stage::start("pack");
|
||||
let pb = spinner("pack");
|
||||
pb.set_message("consolidating column files …");
|
||||
let dst2 = KmerIndex::open(output)?;
|
||||
dst2.pack_matrices(false)?;
|
||||
dst2.meta.mark_indexed().map_err(OKIError::Io)?;
|
||||
pb.finish_and_clear();
|
||||
rep.push(t.stop());
|
||||
}
|
||||
|
||||
KmerIndex::open(output)
|
||||
}
|
||||
}
|
||||
|
||||
// ── Diagnostic report ─────────────────────────────────────────────────────────
|
||||
|
||||
fn print_merge_partition_report(stats: &[PartStat], max_workers: usize) {
|
||||
let total_new: usize = stats.iter().map(|s| s.g_len).sum();
|
||||
let non_empty = stats.iter().filter(|s| s.unitig_bytes > 0).count();
|
||||
|
||||
if non_empty == 0 {
|
||||
info!("merge_partitions report: no data (all partitions empty)");
|
||||
return;
|
||||
}
|
||||
|
||||
info!("─── merge_partitions report ───");
|
||||
info!(
|
||||
" {} partition(s) processed, {} total new kmers",
|
||||
non_empty, total_new,
|
||||
);
|
||||
info!(" max workers: {max_workers}");
|
||||
|
||||
// Top 8 partitions by new-kmer count
|
||||
let mut by_new: Vec<&PartStat> = stats.iter().filter(|s| s.g_len > 0).collect();
|
||||
by_new.sort_by_key(|s| std::cmp::Reverse(s.g_len));
|
||||
if !by_new.is_empty() {
|
||||
info!(" top partitions by new kmers:");
|
||||
for s in by_new.iter().take(8) {
|
||||
info!(
|
||||
" partition {:4} : {}M new kmers ({} unitig bytes)",
|
||||
s.id,
|
||||
s.g_len / 1_000_000,
|
||||
fmt_bytes(s.unitig_bytes),
|
||||
);
|
||||
}
|
||||
}
|
||||
info!("───────────────────────────────");
|
||||
}
|
||||
|
||||
// ── helpers ───────────────────────────────────────────────────────────────────
|
||||
|
||||
fn fmt_bytes(b: u64) -> String {
|
||||
if b >= 1 << 30 {
|
||||
format!("{:.1} GB", b as f64 / (1u64 << 30) as f64)
|
||||
} else if b >= 1 << 20 {
|
||||
format!("{:.1} MB", b as f64 / (1u64 << 20) as f64)
|
||||
} else if b >= 1 << 10 {
|
||||
format!("{:.1} KB", b as f64 / (1u64 << 10) as f64)
|
||||
} else {
|
||||
format!("{b} B")
|
||||
}
|
||||
}
|
||||
|
||||
/// Sum of all unitigs.bin sizes across all layers of partition `i` in `src`.
|
||||
fn partition_unitig_bytes(src: &KmerIndex, i: usize) -> u64 {
|
||||
let mut total = 0u64;
|
||||
for l in 0.. {
|
||||
let p = src.layer_unitigs_path(i, l);
|
||||
if !p.exists() {
|
||||
break;
|
||||
}
|
||||
if let Ok(m) = std::fs::metadata(&p) {
|
||||
total += m.len();
|
||||
}
|
||||
}
|
||||
total
|
||||
}
|
||||
|
||||
fn compute_labels(
|
||||
src_genomes: &[Vec<GenomeInfo>],
|
||||
rename_duplicates: bool,
|
||||
) -> OKIResult<(Vec<Vec<String>>, Vec<GenomeInfo>)> {
|
||||
let mut seen: HashMap<String, usize> = HashMap::new();
|
||||
let mut source_labels: Vec<Vec<String>> = Vec::with_capacity(src_genomes.len());
|
||||
let mut all_genomes: Vec<GenomeInfo> = Vec::new();
|
||||
|
||||
for genomes in src_genomes {
|
||||
let mut labels = Vec::with_capacity(genomes.len());
|
||||
for genome in genomes {
|
||||
let label = &genome.label;
|
||||
let count = seen.entry(label.clone()).or_insert(0);
|
||||
let new_label = if *count == 0 {
|
||||
label.clone()
|
||||
} else if rename_duplicates {
|
||||
format!("{label}.{count}")
|
||||
} else {
|
||||
return Err(OKIError::DuplicateGenomeLabel(label.clone()));
|
||||
};
|
||||
*count += 1;
|
||||
labels.push(new_label.clone());
|
||||
all_genomes.push(GenomeInfo {
|
||||
label: new_label,
|
||||
meta: genome.meta.clone(),
|
||||
});
|
||||
}
|
||||
source_labels.push(labels);
|
||||
}
|
||||
|
||||
Ok((source_labels, all_genomes))
|
||||
}
|
||||
|
||||
fn copy_spectrums(
|
||||
src_root: &Path,
|
||||
dst_root: &Path,
|
||||
old_labels: &[String],
|
||||
new_labels: &[String],
|
||||
) -> io::Result<()> {
|
||||
let src_dir = src_root.join("spectrums");
|
||||
let dst_dir = dst_root.join("spectrums");
|
||||
fs::create_dir_all(&dst_dir)?;
|
||||
for (old, new) in old_labels.iter().zip(new_labels.iter()) {
|
||||
let src_file = src_dir.join(format!("{old}.json"));
|
||||
if src_file.exists() {
|
||||
fs::copy(&src_file, dst_dir.join(format!("{new}.json")))?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn remove_dirs_named(root: &Path, name: &str) -> io::Result<()> {
|
||||
for entry in fs::read_dir(root)? {
|
||||
let entry = entry?;
|
||||
let path = entry.path();
|
||||
if path.is_dir() {
|
||||
if path.file_name().and_then(|n| n.to_str()) == Some(name) {
|
||||
fs::remove_dir_all(&path)?;
|
||||
} else {
|
||||
remove_dirs_named(&path, name)?;
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn format_evidence(ev: &IndexMode) -> String {
|
||||
match ev {
|
||||
IndexMode::Exact => "exact".to_string(),
|
||||
IndexMode::Approx { b, z } => format!("approx (b={b}, z={z})"),
|
||||
IndexMode::Hybrid { b, z } => format!("hybrid (b={b}, z={z})"),
|
||||
}
|
||||
}
|
||||
|
||||
fn is_trivial(genomes: &[GenomeInfo], mode: MergeMode) -> bool {
|
||||
genomes.len() == 1 && mode == MergeMode::Presence
|
||||
}
|
||||
|
||||
fn index_unitig_size(src: &KmerIndex) -> u64 {
|
||||
let n = src.n_partitions();
|
||||
(0..n).map(|i| partition_unitig_bytes(src, i)).sum()
|
||||
}
|
||||
|
||||
fn choose_base(sources: &[&KmerIndex], src_genomes: &[Vec<GenomeInfo>], mode: MergeMode) -> usize {
|
||||
let needs_approx = sources.iter().enumerate().any(|(i, src)| {
|
||||
!is_trivial(&src_genomes[i], mode)
|
||||
&& matches!(
|
||||
src.meta.config.evidence,
|
||||
IndexMode::Approx { .. } | IndexMode::Hybrid { .. }
|
||||
)
|
||||
});
|
||||
|
||||
sources
|
||||
.iter()
|
||||
.enumerate()
|
||||
.filter(|(_, src)| {
|
||||
!needs_approx
|
||||
|| matches!(
|
||||
src.meta.config.evidence,
|
||||
IndexMode::Approx { .. } | IndexMode::Hybrid { .. }
|
||||
)
|
||||
})
|
||||
.max_by_key(|(_, src)| index_unitig_size(src))
|
||||
.map(|(i, _)| i)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn copy_dir_all(src: &Path, dst: &Path) -> io::Result<()> {
|
||||
fs::create_dir_all(dst)?;
|
||||
for entry in fs::read_dir(src)? {
|
||||
let entry = entry?;
|
||||
let src_path = entry.path();
|
||||
let dst_path = dst.join(entry.file_name());
|
||||
if src_path.is_dir() {
|
||||
copy_dir_all(&src_path, &dst_path)?;
|
||||
} else {
|
||||
fs::copy(&src_path, &dst_path)?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,595 @@
|
||||
//! Merging a source partition's new layer into a destination partition:
|
||||
//! de Bruijn graph union (pass 1) then column fill (pass 2).
|
||||
//!
|
||||
//! Submodules: [`src_layer`] (`SrcLayerData`, the opened-source-matrix
|
||||
//! lookup used by pass 2 here and by `rebuild_layer`). The `merge_partition`
|
||||
//! orchestration itself stays in this file — its ~400-line body is one
|
||||
//! tightly threaded pipeline (shared `Arc`/`Mutex` state across pass 1,
|
||||
//! builder setup, and pass 2), not a set of independently callable steps.
|
||||
|
||||
use std::fs;
|
||||
use std::io;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use obipipeline::{
|
||||
Pipeline, PipelineError, PipelineSender, SharedFlatFn, Stage, ThrottleGuard, WorkerPool,
|
||||
make_sink, make_source, make_transform, throttle,
|
||||
};
|
||||
use tracing::debug;
|
||||
|
||||
use obicompactvec::{PersistentBitMatrixBuilder, PersistentCompactIntMatrixBuilder};
|
||||
use obikseq::CanonicalKmer;
|
||||
use obikindex::layer::{IndexMode, TypedLayer, LayeredMap, MphfOnly};
|
||||
use obikindex::layer::utils::layer_dir;
|
||||
use obiskio::UnitigFileReader;
|
||||
use obikindex::{OKIError, OKIResult};
|
||||
|
||||
use obikindex::{ColBuilder, load_meta};
|
||||
use obikindexer::{build_graph, materialize_layer};
|
||||
use obikindex::KmerIndex;
|
||||
|
||||
mod src_layer;
|
||||
|
||||
pub use src_layer::SrcLayerData;
|
||||
|
||||
// ── MergeMode ─────────────────────────────────────────────────────────────────
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum MergeMode {
|
||||
Presence,
|
||||
Count,
|
||||
}
|
||||
|
||||
// ── MatrixBuilder ─────────────────────────────────────────────────────────────
|
||||
//
|
||||
// Wraps whichever matrix builder `mode` calls for, so the merge pipeline never
|
||||
// has to know the on-disk column naming (`col_NNNNNN.pbiv`/`.pciv`) or the
|
||||
// matrix `meta.json` schema itself — both stay private to obicompactvec.
|
||||
// `resume` reopens a matrix directory already closed by a previous builder
|
||||
// session (an existing destination layer), continuing from its current
|
||||
// `n_cols` instead of starting a fresh matrix at 0.
|
||||
|
||||
enum MatrixBuilder {
|
||||
Bit(PersistentBitMatrixBuilder),
|
||||
Int(PersistentCompactIntMatrixBuilder),
|
||||
}
|
||||
|
||||
impl MatrixBuilder {
|
||||
fn new(mode: MergeMode, n: usize, dir: &Path) -> io::Result<Self> {
|
||||
Ok(match mode {
|
||||
MergeMode::Presence => MatrixBuilder::Bit(PersistentBitMatrixBuilder::new(n, dir)?),
|
||||
MergeMode::Count => MatrixBuilder::Int(PersistentCompactIntMatrixBuilder::new(n, dir)?),
|
||||
})
|
||||
}
|
||||
|
||||
fn resume(mode: MergeMode, dir: &Path) -> io::Result<Self> {
|
||||
Ok(match mode {
|
||||
MergeMode::Presence => MatrixBuilder::Bit(PersistentBitMatrixBuilder::resume(dir)?),
|
||||
MergeMode::Count => MatrixBuilder::Int(PersistentCompactIntMatrixBuilder::resume(dir)?),
|
||||
})
|
||||
}
|
||||
|
||||
/// Add a column with no data written (all-zero/false) — for genome
|
||||
/// columns absent from this source (e.g. dst genomes in a new layer).
|
||||
fn add_absent_col(&mut self) -> io::Result<()> {
|
||||
match self {
|
||||
MatrixBuilder::Bit(b) => b.add_col()?.close(),
|
||||
MatrixBuilder::Int(b) => b.add_col()?.close(),
|
||||
}
|
||||
}
|
||||
|
||||
fn add_col(&mut self) -> io::Result<ColBuilder> {
|
||||
Ok(match self {
|
||||
MatrixBuilder::Bit(b) => ColBuilder::Bit(b.add_col()?),
|
||||
MatrixBuilder::Int(b) => ColBuilder::Int(b.add_col()?),
|
||||
})
|
||||
}
|
||||
|
||||
fn close(self) -> io::Result<()> {
|
||||
match self {
|
||||
MatrixBuilder::Bit(b) => b.close(),
|
||||
MatrixBuilder::Int(b) => b.close(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod matrix_builder_tests {
|
||||
use tempfile::tempdir;
|
||||
|
||||
use obicompactvec::{PersistentBitMatrix, PersistentCompactIntMatrix};
|
||||
|
||||
use super::{ColBuilder, MatrixBuilder, MergeMode};
|
||||
|
||||
/// Mirrors `merge_partition`'s "new layer" setup: absent (dst-genome)
|
||||
/// columns get no data, then source columns are filled — all through one
|
||||
/// continuous `MatrixBuilder` session, closed once at the end.
|
||||
#[test]
|
||||
fn new_layer_absent_then_source_columns_presence() {
|
||||
let dir = tempdir().unwrap();
|
||||
let data_dir = dir.path().join("presence");
|
||||
let mut mb = MatrixBuilder::new(MergeMode::Presence, 3, &data_dir).unwrap();
|
||||
|
||||
// Two absent (dst-genome) columns.
|
||||
mb.add_absent_col().unwrap();
|
||||
mb.add_absent_col().unwrap();
|
||||
|
||||
// One source column, filled like pass 2 would.
|
||||
let mut col = mb.add_col().unwrap();
|
||||
match &mut col {
|
||||
ColBuilder::Bit(b) => {
|
||||
b.set(0, true);
|
||||
b.set(1, false);
|
||||
b.set(2, true);
|
||||
}
|
||||
ColBuilder::Int(_) => unreachable!(),
|
||||
}
|
||||
col.close().unwrap();
|
||||
mb.close().unwrap();
|
||||
|
||||
let m = PersistentBitMatrix::open(dir.path()).unwrap();
|
||||
assert_eq!(m.n_cols(), 3);
|
||||
assert_eq!(&*m.row(0), &[false, false, true]);
|
||||
assert_eq!(&*m.row(1), &[false, false, false]);
|
||||
assert_eq!(&*m.row(2), &[false, false, true]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn new_layer_absent_then_source_columns_count() {
|
||||
let dir = tempdir().unwrap();
|
||||
let data_dir = dir.path().join("counts");
|
||||
let mut mb = MatrixBuilder::new(MergeMode::Count, 2, &data_dir).unwrap();
|
||||
|
||||
mb.add_absent_col().unwrap();
|
||||
|
||||
let mut col = mb.add_col().unwrap();
|
||||
match &mut col {
|
||||
ColBuilder::Int(b) => {
|
||||
b.set(0, 7);
|
||||
b.set(1, 42);
|
||||
}
|
||||
ColBuilder::Bit(_) => unreachable!(),
|
||||
}
|
||||
col.close().unwrap();
|
||||
mb.close().unwrap();
|
||||
|
||||
let m = PersistentCompactIntMatrix::open(dir.path()).unwrap();
|
||||
assert_eq!(m.n_cols(), 2);
|
||||
assert_eq!(&*m.row(0), &[0u32, 7]);
|
||||
assert_eq!(&*m.row(1), &[0u32, 42]);
|
||||
}
|
||||
|
||||
/// Mirrors `merge_partition`'s "existing dst layer" setup: an
|
||||
/// already-closed matrix (from a previous merge) gets more columns
|
||||
/// appended via `resume`, without callers ever seeing `col_path`/
|
||||
/// `MatrixMeta` — `n` itself is read back from the matrix directory.
|
||||
#[test]
|
||||
fn resume_appends_source_columns_to_existing_layer() {
|
||||
let dir = tempdir().unwrap();
|
||||
let data_dir = dir.path().join("presence");
|
||||
|
||||
// Previous merge: one dst-genome column already on disk.
|
||||
let mut mb0 = MatrixBuilder::new(MergeMode::Presence, 3, &data_dir).unwrap();
|
||||
mb0.add_absent_col().unwrap();
|
||||
mb0.close().unwrap();
|
||||
|
||||
// This merge: resume and append two more source columns.
|
||||
let mut mb = MatrixBuilder::resume(MergeMode::Presence, &data_dir).unwrap();
|
||||
for vals in [[true, false, true], [false, true, false]] {
|
||||
let mut col = mb.add_col().unwrap();
|
||||
match &mut col {
|
||||
ColBuilder::Bit(b) => {
|
||||
for (slot, v) in vals.into_iter().enumerate() {
|
||||
b.set(slot, v);
|
||||
}
|
||||
}
|
||||
ColBuilder::Int(_) => unreachable!(),
|
||||
}
|
||||
col.close().unwrap();
|
||||
}
|
||||
mb.close().unwrap();
|
||||
|
||||
let m = PersistentBitMatrix::open(dir.path()).unwrap();
|
||||
assert_eq!(m.n_cols(), 3);
|
||||
assert_eq!(&*m.row(0), &[false, true, false]);
|
||||
assert_eq!(&*m.row(1), &[false, false, true]);
|
||||
assert_eq!(&*m.row(2), &[false, true, false]);
|
||||
}
|
||||
}
|
||||
|
||||
// ── KmerPartition::merge_partition ────────────────────────────────────────────
|
||||
|
||||
impl KmerIndex {
|
||||
/// Merge `sources` into destination partition `i`.
|
||||
///
|
||||
/// Each entry in `sources` is `(partition, n_genomes)` where `n_genomes` is
|
||||
/// the number of genome columns that source contributes. A merged index
|
||||
/// contributes more than one. The total new columns added to the destination
|
||||
/// is `sum(n_genomes)`.
|
||||
///
|
||||
/// `n_dst_genomes` is the number of genome columns already in the destination
|
||||
/// matrices (copied from source_0 before this call).
|
||||
pub fn merge_partition(
|
||||
&self,
|
||||
i: usize,
|
||||
sources: &[(&KmerIndex, usize)],
|
||||
mode: MergeMode,
|
||||
n_dst_genomes: usize,
|
||||
block_bits: u8,
|
||||
evidence: &IndexMode,
|
||||
) -> OKIResult<usize> {
|
||||
let dst_index_dir = self.index_dir(i);
|
||||
if !dst_index_dir.exists() {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
load_meta(&dst_index_dir)?; // ensure meta.json exists before LayeredMap::open
|
||||
let dst_map = Arc::new(LayeredMap::<()>::open(&dst_index_dir)?);
|
||||
let n_dst_layers = dst_map.n_layers();
|
||||
let n_src_total: usize = sources.iter().map(|(_, n)| *n).sum();
|
||||
|
||||
// First merge in presence mode: init presence matrices on existing layers
|
||||
// (all slots true — every kmer in those layers belongs to genome_0).
|
||||
if n_dst_genomes == 1 && mode == MergeMode::Presence {
|
||||
for l in 0..n_dst_layers {
|
||||
TypedLayer::<()>::init_presence_matrix(
|
||||
&layer_dir(&dst_index_dir, l),
|
||||
dst_map.layer(l).n(),
|
||||
)?;
|
||||
}
|
||||
}
|
||||
|
||||
// ── Pass 1: pipeline — parallel file read + dst_map filter + graph fill ─
|
||||
//
|
||||
// Source : list of unitigs.bin paths (one per source × layer)
|
||||
// Flat : open file, emit Vec<CanonicalKmer> batches (BeeGFS parallel I/O)
|
||||
// Transform: filter via dst_map.query() — thread-safe, LayeredMap<()>: Sync
|
||||
// Sink : push new kmers into GraphDeBruijn (single thread, no locks needed)
|
||||
|
||||
// Collect file paths (propagates load_meta errors before the pipeline starts)
|
||||
let mut unitig_paths: Vec<PathBuf> = Vec::new();
|
||||
for (src, _) in sources.iter() {
|
||||
let src_index_dir = src.index_dir(i);
|
||||
if !src_index_dir.exists() {
|
||||
continue;
|
||||
}
|
||||
let src_meta = load_meta(&src_index_dir)?;
|
||||
for l in 0..src_meta.n_layers {
|
||||
let p = layer_dir(&src_index_dir, l).join("unitigs.bin");
|
||||
if p.exists() {
|
||||
unitig_paths.push(p);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let n_src_layers = unitig_paths.len();
|
||||
debug!("partition {i}: de Bruijn graph build start — {n_src_layers} source layer(s)");
|
||||
|
||||
const BATCH: usize = 4096;
|
||||
let n_workers = rayon::current_num_threads().min(16).max(4);
|
||||
// At most 2 files open simultaneously: keeps n_workers-2 workers free
|
||||
// for the Transform stage. Each open file monopolises one worker for the
|
||||
// full duration of its read, so this must stay well below n_workers.
|
||||
let max_open = 2;
|
||||
|
||||
let dst_filter = Arc::clone(&dst_map);
|
||||
|
||||
let g = build_graph(
|
||||
unitig_paths.into_iter(),
|
||||
move |path: PathBuf, emit: &mut dyn FnMut(Vec<CanonicalKmer>)| -> OKIResult<()> {
|
||||
let reader = UnitigFileReader::open_sequential(&path)?;
|
||||
let mut batch: Vec<CanonicalKmer> = Vec::with_capacity(BATCH);
|
||||
for (kmer, _, _) in reader.iter_indexed_canonical_kmers() {
|
||||
batch.push(kmer);
|
||||
if batch.len() == BATCH {
|
||||
emit(std::mem::replace(&mut batch, Vec::with_capacity(BATCH)));
|
||||
}
|
||||
}
|
||||
if !batch.is_empty() {
|
||||
emit(batch);
|
||||
}
|
||||
Ok(())
|
||||
},
|
||||
move |kmer| dst_filter.query(kmer).is_none(),
|
||||
n_workers,
|
||||
max_open,
|
||||
)?;
|
||||
|
||||
let any_new = g.len() > 0;
|
||||
debug!(
|
||||
"partition {i}: de Bruijn graph done — {} new kmers",
|
||||
g.len()
|
||||
);
|
||||
|
||||
// Build new layer from de Bruijn graph if there are new kmers.
|
||||
let new_layer_idx = n_dst_layers;
|
||||
let new_layer_dir = layer_dir(&dst_index_dir, new_layer_idx);
|
||||
|
||||
let n_new = if any_new {
|
||||
debug!("partition {i}: unitig traversal start — {} nodes", g.len());
|
||||
let n_nodes = materialize_layer(g, &new_layer_dir, block_bits, evidence)?;
|
||||
debug!("partition {i}: MPHF build done");
|
||||
n_nodes
|
||||
} else {
|
||||
drop(g);
|
||||
0
|
||||
};
|
||||
|
||||
let t_open = std::time::Instant::now();
|
||||
let new_mphf: Option<Arc<MphfOnly>> = if any_new {
|
||||
Some(Arc::new(MphfOnly::open(&new_layer_dir)?))
|
||||
} else {
|
||||
None
|
||||
};
|
||||
debug!(
|
||||
"partition {i}: MPHF open in {:.3}s",
|
||||
t_open.elapsed().as_secs_f64()
|
||||
);
|
||||
|
||||
// ── Prepare matrix directories for the new layer ──────────────────────
|
||||
// Absent columns (dst genomes) get an all-zero/false column. Source-genome
|
||||
// columns are created as mutable builders for pass 2. `new_mb` is kept
|
||||
// open (not `close`d) until pass 2 has filled every source column, so its
|
||||
// `n_cols` bookkeeping stays in sync with the columns actually created.
|
||||
let (new_src_builders, new_mb): (Vec<ColBuilder>, Option<MatrixBuilder>) = if any_new {
|
||||
let data_dir = match mode {
|
||||
MergeMode::Presence => new_layer_dir.join("presence"),
|
||||
MergeMode::Count => new_layer_dir.join("counts"),
|
||||
};
|
||||
fs::create_dir_all(&data_dir)?;
|
||||
let mut mb = MatrixBuilder::new(mode, n_new, &data_dir).map_err(OKIError::Io)?;
|
||||
for _ in 0..n_dst_genomes {
|
||||
mb.add_absent_col().map_err(OKIError::Io)?;
|
||||
}
|
||||
let cols = (0..n_src_total)
|
||||
.map(|_| mb.add_col().map_err(OKIError::Io))
|
||||
.collect::<OKIResult<Vec<_>>>()?;
|
||||
(cols, Some(mb))
|
||||
} else {
|
||||
(vec![], None)
|
||||
};
|
||||
|
||||
let t_builders = std::time::Instant::now();
|
||||
// Builders for existing layers: n_src_total per layer, resumed from
|
||||
// each layer's own matrix directory (already holding n_dst_genomes
|
||||
// columns from a previous merge). Columns land at
|
||||
// n_dst_genomes .. n_dst_genomes + n_src_total - 1.
|
||||
let mut exist_mbs: Vec<MatrixBuilder> = Vec::with_capacity(n_dst_layers);
|
||||
let mut exist_builders: Vec<Vec<ColBuilder>> = Vec::with_capacity(n_dst_layers);
|
||||
for l in 0..n_dst_layers {
|
||||
let layer_dir = layer_dir(&dst_index_dir, l);
|
||||
let data_dir = match mode {
|
||||
MergeMode::Presence => layer_dir.join("presence"),
|
||||
MergeMode::Count => layer_dir.join("counts"),
|
||||
};
|
||||
let mut mb = MatrixBuilder::resume(mode, &data_dir).map_err(OKIError::Io)?;
|
||||
let cols = (0..n_src_total)
|
||||
.map(|_| mb.add_col().map_err(OKIError::Io))
|
||||
.collect::<OKIResult<Vec<_>>>()?;
|
||||
exist_mbs.push(mb);
|
||||
exist_builders.push(cols);
|
||||
}
|
||||
|
||||
debug!(
|
||||
"partition {i}: builders ready in {:.3}s",
|
||||
t_builders.elapsed().as_secs_f64()
|
||||
);
|
||||
|
||||
// ── Pass 2: fill builders (pipeline) ─────────────────────────────────
|
||||
let t_pass2 = std::time::Instant::now();
|
||||
// Collect source items before the pipeline so load_meta errors propagate
|
||||
// via ? before any worker thread is spawned.
|
||||
let mut pass2_items: Vec<(usize, usize, PathBuf)> = Vec::new();
|
||||
{
|
||||
let mut col_offset = 0usize;
|
||||
for (src, src_n) in sources.iter() {
|
||||
let src_index_dir = src.index_dir(i);
|
||||
if !src_index_dir.exists() {
|
||||
col_offset += src_n;
|
||||
continue;
|
||||
}
|
||||
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);
|
||||
if src_layer_dir.join("unitigs.bin").exists() {
|
||||
pass2_items.push((col_offset, *src_n, src_layer_dir));
|
||||
}
|
||||
}
|
||||
col_offset += src_n;
|
||||
}
|
||||
}
|
||||
|
||||
enum Pass2Data {
|
||||
SrcLayer((usize, usize, PathBuf, ThrottleGuard)),
|
||||
RawBatch((usize, usize, Arc<SrcLayerData>, Vec<CanonicalKmer>)),
|
||||
WriteBatch(Vec<(Option<usize>, usize, usize, u32)>),
|
||||
}
|
||||
|
||||
let exist_locked: Vec<Vec<Arc<Mutex<ColBuilder>>>> = exist_builders
|
||||
.into_iter()
|
||||
.map(|layer| layer.into_iter().map(|b| Arc::new(Mutex::new(b))).collect())
|
||||
.collect();
|
||||
let new_locked: Vec<Arc<Mutex<ColBuilder>>> = new_src_builders
|
||||
.into_iter()
|
||||
.map(|b| Arc::new(Mutex::new(b)))
|
||||
.collect();
|
||||
let exist_sink: Vec<Vec<Arc<Mutex<ColBuilder>>>> = exist_locked
|
||||
.iter()
|
||||
.map(|layer| layer.iter().map(Arc::clone).collect())
|
||||
.collect();
|
||||
let new_sink: Vec<Arc<Mutex<ColBuilder>>> = new_locked.iter().map(Arc::clone).collect();
|
||||
let dst_map_t2 = Arc::clone(&dst_map);
|
||||
let new_mphf_t2 = new_mphf.clone();
|
||||
let pass2_err: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
|
||||
let err_cap2 = Arc::clone(&pass2_err);
|
||||
|
||||
let capacity = 2;
|
||||
let throttled_pass2 = throttle(pass2_items.into_iter(), max_open);
|
||||
|
||||
let pipeline2 = Pipeline::new(
|
||||
make_source!(
|
||||
Pass2Data,
|
||||
throttled_pass2.map(|t| {
|
||||
let (col_offset, src_n, src_layer_dir) = t.item;
|
||||
(col_offset, src_n, src_layer_dir, t.guard)
|
||||
}),
|
||||
SrcLayer
|
||||
),
|
||||
vec![
|
||||
Stage::Flat(Arc::new(
|
||||
move |data: Pass2Data,
|
||||
push: &PipelineSender<Result<Pass2Data, PipelineError>>,
|
||||
delta: &PipelineSender<isize>| {
|
||||
if let Pass2Data::SrcLayer((col_offset, src_n, src_layer_dir, _guard)) =
|
||||
data
|
||||
{
|
||||
// _guard dropped at end of block, releasing the slot.
|
||||
let reader = match UnitigFileReader::open_sequential(
|
||||
&src_layer_dir.join("unitigs.bin"),
|
||||
) {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
*err_cap2.lock().unwrap() = Some(e.to_string());
|
||||
delta.send(-1).ok();
|
||||
return;
|
||||
}
|
||||
};
|
||||
let src_data = match SrcLayerData::open(&src_layer_dir, mode) {
|
||||
Ok(d) => Arc::new(d),
|
||||
Err(e) => {
|
||||
*err_cap2.lock().unwrap() = Some(e.to_string());
|
||||
delta.send(-1).ok();
|
||||
return;
|
||||
}
|
||||
};
|
||||
const BATCH: usize = 4096;
|
||||
let mut batch: Vec<CanonicalKmer> = Vec::with_capacity(BATCH);
|
||||
let mut count: isize = 0;
|
||||
for (kmer, _, _) in reader.iter_indexed_canonical_kmers() {
|
||||
batch.push(kmer);
|
||||
if batch.len() == BATCH {
|
||||
let b =
|
||||
std::mem::replace(&mut batch, Vec::with_capacity(BATCH));
|
||||
push.send(Ok(Pass2Data::RawBatch((
|
||||
col_offset,
|
||||
src_n,
|
||||
Arc::clone(&src_data),
|
||||
b,
|
||||
))))
|
||||
.ok();
|
||||
count += 1;
|
||||
}
|
||||
}
|
||||
if !batch.is_empty() {
|
||||
push.send(Ok(Pass2Data::RawBatch((
|
||||
col_offset, src_n, src_data, batch,
|
||||
))))
|
||||
.ok();
|
||||
count += 1;
|
||||
}
|
||||
delta.send(count - 1).ok();
|
||||
}
|
||||
},
|
||||
) as SharedFlatFn<Pass2Data>),
|
||||
make_transform!(
|
||||
Pass2Data,
|
||||
{
|
||||
move |(col_offset, src_n, src_data, kmers): (
|
||||
usize,
|
||||
usize,
|
||||
Arc<SrcLayerData>,
|
||||
Vec<CanonicalKmer>,
|
||||
)|
|
||||
-> Vec<(Option<usize>, usize, usize, u32)> {
|
||||
let mut ops: Vec<(Option<usize>, usize, usize, u32)> = Vec::new();
|
||||
for kmer in kmers {
|
||||
let values = src_data.lookup(kmer, src_n);
|
||||
if let Some((dst_layer, hit)) = dst_map_t2.query(kmer) {
|
||||
for (g, val) in values.into_iter().enumerate() {
|
||||
ops.push((Some(dst_layer), col_offset + g, hit.slot, val));
|
||||
}
|
||||
} else if let Some(ref mphf) = new_mphf_t2 {
|
||||
let slot = mphf.index(kmer);
|
||||
for (g, val) in values.into_iter().enumerate() {
|
||||
ops.push((None, col_offset + g, slot, val));
|
||||
}
|
||||
}
|
||||
}
|
||||
ops
|
||||
}
|
||||
},
|
||||
RawBatch,
|
||||
WriteBatch
|
||||
),
|
||||
],
|
||||
make_sink!(
|
||||
Pass2Data,
|
||||
{
|
||||
move |ops: Vec<(Option<usize>, usize, usize, u32)>| {
|
||||
for (layer_opt, col, slot, val) in ops {
|
||||
match layer_opt {
|
||||
Some(l) => exist_sink[l][col].lock().unwrap().set_val(slot, val),
|
||||
None => new_sink[col].lock().unwrap().set_val(slot, val),
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
WriteBatch
|
||||
),
|
||||
);
|
||||
|
||||
WorkerPool::new(pipeline2, n_workers, capacity).run();
|
||||
debug!(
|
||||
"partition {i}: pass2 pipeline done in {:.3}s",
|
||||
t_pass2.elapsed().as_secs_f64()
|
||||
);
|
||||
|
||||
if let Some(msg) = Arc::try_unwrap(pass2_err)
|
||||
.unwrap_or_else(|_| panic!("pass2: pass2_err not uniquely owned"))
|
||||
.into_inner()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
{
|
||||
return Err(OKIError::InvalidData {
|
||||
context: "merge pass2",
|
||||
detail: msg,
|
||||
});
|
||||
}
|
||||
|
||||
let t_close = std::time::Instant::now();
|
||||
// ── Close builders and update metadata ────────────────────────────────
|
||||
for (mb, builders) in exist_mbs.into_iter().zip(exist_locked.into_iter()) {
|
||||
for b in builders {
|
||||
Arc::try_unwrap(b)
|
||||
.unwrap_or_else(|_| panic!("pass2: exist_builder not uniquely owned"))
|
||||
.into_inner()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.close()?;
|
||||
}
|
||||
mb.close().map_err(OKIError::Io)?;
|
||||
}
|
||||
|
||||
for b in new_locked {
|
||||
Arc::try_unwrap(b)
|
||||
.unwrap_or_else(|_| panic!("pass2: new_builder not uniquely owned"))
|
||||
.into_inner()
|
||||
.unwrap_or_else(|e| e.into_inner())
|
||||
.close()?;
|
||||
}
|
||||
if let Some(mb) = new_mb {
|
||||
mb.close().map_err(OKIError::Io)?;
|
||||
|
||||
let mut part_meta = self.partition_meta(i)?;
|
||||
part_meta.n_layers = new_layer_idx + 1;
|
||||
part_meta
|
||||
.save(&dst_index_dir)?;
|
||||
}
|
||||
|
||||
debug!(
|
||||
"partition {i}: builders closed in {:.3}s",
|
||||
t_close.elapsed().as_secs_f64()
|
||||
);
|
||||
|
||||
Ok(n_new)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
use std::path::Path;
|
||||
|
||||
use obicompactvec::{MatrixGroupOps, PersistentBitMatrix, PersistentCompactIntMatrix};
|
||||
use obikseq::CanonicalKmer;
|
||||
use obikindex::layer::MphfOnly;
|
||||
use obikindex::{OKIError, OKIResult};
|
||||
|
||||
|
||||
use super::MergeMode;
|
||||
|
||||
// ── SrcLayerData — opened source matrix for pass-2 lookup ─────────────────────
|
||||
|
||||
pub enum SrcLayerData {
|
||||
Presence(MphfOnly, PersistentBitMatrix),
|
||||
Count(MphfOnly, PersistentCompactIntMatrix),
|
||||
}
|
||||
|
||||
impl SrcLayerData {
|
||||
pub fn open(layer_dir: &Path, merge_mode: MergeMode) -> OKIResult<Self> {
|
||||
let counts_dir = layer_dir.join("counts");
|
||||
match merge_mode {
|
||||
MergeMode::Presence => {
|
||||
if counts_dir.exists() && !layer_dir.join("presence").exists() {
|
||||
let mphf = MphfOnly::open(layer_dir)?;
|
||||
let mat = PersistentCompactIntMatrix::open(layer_dir).map_err(OKIError::Io)?;
|
||||
Ok(SrcLayerData::Count(mphf, mat))
|
||||
} else {
|
||||
// presence dir exists, or neither exists → Implicit handled by open()
|
||||
let mphf = MphfOnly::open(layer_dir)?;
|
||||
let mat = PersistentBitMatrix::open(layer_dir).map_err(OKIError::Io)?;
|
||||
Ok(SrcLayerData::Presence(mphf, mat))
|
||||
}
|
||||
}
|
||||
MergeMode::Count => {
|
||||
let mphf = MphfOnly::open(layer_dir)?;
|
||||
if counts_dir.exists() {
|
||||
let mat = PersistentCompactIntMatrix::open(layer_dir).map_err(OKIError::Io)?;
|
||||
Ok(SrcLayerData::Count(mphf, mat))
|
||||
} else {
|
||||
// No counts → treat as implicit presence (all 1s)
|
||||
let mat = PersistentBitMatrix::open(layer_dir).map_err(OKIError::Io)?;
|
||||
Ok(SrcLayerData::Presence(mphf, mat))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Return one value per source genome for `kmer`.
|
||||
/// The caller guarantees `kmer` is in the source MPHF domain.
|
||||
#[inline]
|
||||
pub(crate) fn lookup(&self, kmer: CanonicalKmer, n_genomes: usize) -> Vec<u32> {
|
||||
let mut buf = vec![0u32; n_genomes];
|
||||
match self {
|
||||
SrcLayerData::Presence(mphf, mat) => mat.fill_row(mphf.index(kmer), &mut buf),
|
||||
SrcLayerData::Count(mphf, mat) => mat.fill_row(mphf.index(kmer), &mut buf),
|
||||
}
|
||||
buf
|
||||
}
|
||||
|
||||
pub fn n_slots(&self) -> usize {
|
||||
match self {
|
||||
SrcLayerData::Presence(_, mat) => mat.n(),
|
||||
SrcLayerData::Count(_, mat) => mat.n(),
|
||||
}
|
||||
}
|
||||
|
||||
/// MPHF lookup: returns the slot index for `kmer` (kmer must be in the domain).
|
||||
#[inline]
|
||||
pub fn slot(&self, kmer: CanonicalKmer) -> usize {
|
||||
match self {
|
||||
SrcLayerData::Presence(mphf, _) => mphf.index(kmer),
|
||||
SrcLayerData::Count(mphf, _) => mphf.index(kmer),
|
||||
}
|
||||
}
|
||||
|
||||
/// Row lookup by slot index, bypassing the MPHF.
|
||||
#[inline]
|
||||
pub fn fill_row_by_slot(&self, slot: usize, n_genomes: usize) -> Vec<u32> {
|
||||
let mut buf = vec![0u32; n_genomes];
|
||||
match self {
|
||||
SrcLayerData::Presence(_, mat) => mat.fill_row(slot, &mut buf),
|
||||
SrcLayerData::Count(_, mat) => mat.fill_row(slot, &mut buf),
|
||||
}
|
||||
buf
|
||||
}
|
||||
|
||||
/// Call `f` with a reference to the underlying matrix as `&dyn MatrixGroupOps`.
|
||||
pub fn with_matrix<R>(&self, f: impl FnOnce(&dyn MatrixGroupOps) -> R) -> R {
|
||||
match self {
|
||||
SrcLayerData::Presence(_, mat) => f(mat),
|
||||
SrcLayerData::Count(_, mat) => f(mat),
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user