Start of a huge refactoring
This commit is contained in:
@@ -6,10 +6,14 @@ edition = "2024"
|
||||
[dependencies]
|
||||
obikindex = { path = "../obikindex" }
|
||||
obikindexer = { path = "../obikindexer" }
|
||||
obikfilter = { path = "../obikfilter" }
|
||||
obikalgorithm = { path = "../obikalgorithm" }
|
||||
obicompactvec = { path = "../obicompactvec" }
|
||||
obiskio = { path = "../obiskio" }
|
||||
obikseq = { path = "../obikseq" }
|
||||
obipipeline = { path = "../obipipeline" }
|
||||
obisys = { path = "../obisys" }
|
||||
rayon = "1"
|
||||
tracing = "0.1.44"
|
||||
|
||||
[dev-dependencies]
|
||||
tempfile = "3"
|
||||
|
||||
@@ -1,17 +1,19 @@
|
||||
//! 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/
|
||||
//! from a chosen base 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.
|
||||
//! Structured on the model of `obikindexer::algorithms`: [`Merge`]
|
||||
//! implements `obikalgorithm::Algorithm` (two-phase `new` + setters, then
|
||||
//! `run`). Pass 2 of [`partition_merge::merge_partition`] reads each source
|
||||
//! layer through `obikindex::layer::KmerLayer` directly, batched
|
||||
//! column-first via `KmerLayer::fill_sub_matrix` — no merge-local matrix
|
||||
//! wrapper needed.
|
||||
|
||||
mod merge;
|
||||
mod merge_layer;
|
||||
mod partition_merge;
|
||||
|
||||
pub use merge::*;
|
||||
pub use merge_layer::{MergeMode, SrcLayerData};
|
||||
pub use merge::{Merge, MergeMode};
|
||||
|
||||
+156
-108
@@ -1,22 +1,25 @@
|
||||
use std::collections::HashMap;
|
||||
use std::fs;
|
||||
use std::io;
|
||||
use std::path::Path;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
use obikindex::IndexBuilder;
|
||||
use obikalgorithm::Algorithm;
|
||||
use obikindex::layer::IndexMode;
|
||||
use obikindex::partition::PARTITIONS_SUBDIR;
|
||||
use obikindex::{GenomeInfo, IndexBuilder, IndexState, KmerIndex};
|
||||
|
||||
use obisys::{Reporter, Stage, progress_bar, spinner};
|
||||
use obisys::{PartitionRunner, Progress, Reporter, Stage};
|
||||
use tracing::{debug, info};
|
||||
|
||||
use obikindex::layer::IndexMode;
|
||||
use crate::partition_merge::merge_partition;
|
||||
|
||||
use obikindex::{OKIError, OKIResult};
|
||||
use obikindex::KmerIndex;
|
||||
use obikindex::{GenomeInfo, IndexMeta};
|
||||
use obikindex::IndexState;
|
||||
use obisys::PartitionRunner;
|
||||
// ── MergeMode ─────────────────────────────────────────────────────────────────
|
||||
|
||||
pub use crate::merge_layer::MergeMode;
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum MergeMode {
|
||||
Presence,
|
||||
Count,
|
||||
}
|
||||
|
||||
// ── per-partition diagnostic record ──────────────────────────────────────────
|
||||
|
||||
@@ -27,58 +30,117 @@ struct PartStat {
|
||||
g_len: usize,
|
||||
}
|
||||
|
||||
// ── main merge entry point ────────────────────────────────────────────────────
|
||||
// ── Merge — two-phase construction, `Algorithm` shape ─────────────────────────
|
||||
|
||||
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();
|
||||
/// Merges `sources` into a new index at `output`. Bootstraps from whichever
|
||||
/// source can host the merge unapproximated (see [`choose_base`]) via
|
||||
/// `KmerIndex::create` + `set_genomes` — `config` is written once, at
|
||||
/// creation, never patched afterward, same invariant every other
|
||||
/// construction path (`create_skeleton`) already holds. Then merges each
|
||||
/// remaining source into the bootstrap destination, partition by partition,
|
||||
/// via [`merge_partition`].
|
||||
pub struct Merge<'a> {
|
||||
sources: &'a [&'a KmerIndex],
|
||||
output: PathBuf,
|
||||
mode: MergeMode,
|
||||
force: bool,
|
||||
rename_duplicates: bool,
|
||||
reporter: Reporter,
|
||||
on_progress: Option<Box<dyn FnMut(Progress) + Send + 'a>>,
|
||||
}
|
||||
|
||||
impl<'a> Merge<'a> {
|
||||
pub fn new(sources: &'a [&'a KmerIndex], output: impl Into<PathBuf>, mode: MergeMode) -> Self {
|
||||
Self {
|
||||
sources,
|
||||
output: output.into(),
|
||||
mode,
|
||||
force: false,
|
||||
rename_duplicates: false,
|
||||
reporter: Reporter::new(),
|
||||
on_progress: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Remove a pre-existing index at `output` before creating the merged
|
||||
/// one (default: `false`, fails if one already exists).
|
||||
pub fn force(mut self, v: bool) -> Self {
|
||||
self.force = v;
|
||||
self
|
||||
}
|
||||
|
||||
/// Disambiguate duplicate genome labels across sources by suffixing
|
||||
/// `.1`, `.2`, ... instead of failing (default: `false`).
|
||||
pub fn rename_duplicates(mut self, v: bool) -> Self {
|
||||
self.rename_duplicates = v;
|
||||
self
|
||||
}
|
||||
|
||||
/// Progress callback, called once per completed partition.
|
||||
pub fn on_progress(mut self, cb: impl FnMut(Progress) + Send + 'a) -> Self {
|
||||
self.on_progress = Some(Box::new(cb));
|
||||
self
|
||||
}
|
||||
|
||||
/// Per-stage timing report accumulated during `run` (bootstrap,
|
||||
/// spectrums, merge_partitions, finalize) — call after `run` returns.
|
||||
pub fn reporter(&self) -> &Reporter {
|
||||
&self.reporter
|
||||
}
|
||||
}
|
||||
|
||||
impl Algorithm for Merge<'_> {
|
||||
type Output = KmerIndex;
|
||||
|
||||
fn run(&mut self) -> obikalgorithm::Result<KmerIndex> {
|
||||
let output = self.output.clone();
|
||||
let sources = self.sources;
|
||||
|
||||
if sources.is_empty() {
|
||||
return Err(OKIError::Io(io::Error::new(
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidInput,
|
||||
"merge requires at least one source index",
|
||||
)));
|
||||
)
|
||||
.into());
|
||||
}
|
||||
|
||||
// ── Validate config compatibility ─────────────────────────────────────
|
||||
let ref0 = sources[0];
|
||||
for src in sources {
|
||||
if src.state()? != IndexState::Indexed {
|
||||
return Err(OKIError::NotIndexed(src.root_path.clone()));
|
||||
return Err(format!("{}: source index is not fully built", src.dir().display()).into());
|
||||
}
|
||||
if src.kmer_size() != ref0.kmer_size()
|
||||
|| src.minimizer_size() != ref0.minimizer_size()
|
||||
|| src.n_partitions() != ref0.n_partitions()
|
||||
{
|
||||
return Err(OKIError::IncompatibleConfig);
|
||||
return Err("merge sources have incompatible k-mer size / minimizer size / partition count".into());
|
||||
}
|
||||
if mode == MergeMode::Count && !src.meta.config.with_counts {
|
||||
return Err(OKIError::MismatchedMode);
|
||||
if self.mode == MergeMode::Count && !src.meta().config.with_counts {
|
||||
return Err(format!(
|
||||
"{}: source index has no per-genome counts, cannot merge in count mode",
|
||||
src.dir().display()
|
||||
)
|
||||
.into());
|
||||
}
|
||||
}
|
||||
|
||||
// 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<_>>()?;
|
||||
let src_genomes: Vec<Vec<GenomeInfo>> = sources
|
||||
.iter()
|
||||
.map(|s| s.meta().genomes())
|
||||
.collect::<io::Result<_>>()?;
|
||||
|
||||
// ── Log source characteristics and choose base ────────────────────────
|
||||
let mode_str = if mode == MergeMode::Presence {
|
||||
let mode_str = if self.mode == MergeMode::Presence {
|
||||
"presence"
|
||||
} else {
|
||||
"count"
|
||||
};
|
||||
info!(
|
||||
"merge: {} source(s), smer-size={}, mode={}",
|
||||
"merge: {} source(s), kmer-size={}, mode={}",
|
||||
sources.len(),
|
||||
sources[0].kmer_size(),
|
||||
mode_str,
|
||||
@@ -89,7 +151,7 @@ impl KmerIndex {
|
||||
} else {
|
||||
format!("{} genomes", src_genomes[i].len())
|
||||
};
|
||||
let trivial_str = if is_trivial(&src_genomes[i], mode) {
|
||||
let trivial_str = if is_trivial(&src_genomes[i], self.mode) {
|
||||
" [trivial: no data approximation]"
|
||||
} else {
|
||||
""
|
||||
@@ -97,32 +159,28 @@ impl KmerIndex {
|
||||
info!(
|
||||
" [{}] {} — {}, {}, {}{}",
|
||||
i,
|
||||
src.root_path.display(),
|
||||
format_evidence(&src.meta.config.evidence),
|
||||
src.dir().display(),
|
||||
format_evidence(&src.meta().config.evidence),
|
||||
genome_str,
|
||||
mode_str,
|
||||
trivial_str,
|
||||
);
|
||||
}
|
||||
|
||||
let base_idx = choose_base(sources, &src_genomes, mode);
|
||||
let base_idx = choose_base(sources, &src_genomes, self.mode);
|
||||
let needs_approx = sources.iter().enumerate().any(|(i, src)| {
|
||||
!is_trivial(&src_genomes[i], mode)
|
||||
!is_trivial(&src_genomes[i], self.mode)
|
||||
&& matches!(
|
||||
src.meta.config.evidence,
|
||||
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 {
|
||||
""
|
||||
},
|
||||
format_evidence(&sources[base_idx].meta().config.evidence),
|
||||
if needs_approx { "forced approx — " } else { "" },
|
||||
base_idx,
|
||||
sources[base_idx].root_path.display(),
|
||||
sources[base_idx].dir().display(),
|
||||
);
|
||||
|
||||
let mut ordered: Vec<&KmerIndex> = Vec::with_capacity(sources.len());
|
||||
@@ -137,56 +195,56 @@ impl KmerIndex {
|
||||
}
|
||||
let sources: &[&KmerIndex] = &ordered;
|
||||
let src_genomes = ordered_genomes;
|
||||
let evidence = sources[0].meta.config.evidence.clone();
|
||||
let evidence = sources[0].meta().config.evidence.clone();
|
||||
|
||||
// ── Compute final genome labels ────────────────────────────────────────
|
||||
let (source_labels, all_genomes) = compute_labels(&src_genomes, rename_duplicates)?;
|
||||
let (source_labels, all_genomes) = compute_labels(&src_genomes, self.rename_duplicates)?;
|
||||
|
||||
// ── Prepare output directory ──────────────────────────────────────────
|
||||
KmerIndex::clear_output_for_create(output, force)?;
|
||||
KmerIndex::clear_output_for_create(&output, self.force)?;
|
||||
|
||||
// ── Bootstrap: copy first source to output ────────────────────────────
|
||||
// ── Bootstrap: create the destination with its final config/genomes
|
||||
// up front, then copy the base source's partition data into it ─────
|
||||
info!(
|
||||
"bootstrap: copying {} → {} ({} genome(s))",
|
||||
sources[0].root_path.display(),
|
||||
"bootstrap: {} → {} ({} genome(s))",
|
||||
sources[0].dir().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;
|
||||
let mut config = sources[0].meta().config.clone();
|
||||
config.with_counts = self.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")?;
|
||||
let dst = KmerIndex::create(&output, config, None)?;
|
||||
dst.meta()
|
||||
.set_genomes(all_genomes)
|
||||
.map_err(|e| io::Error::other(e.to_string()))?;
|
||||
|
||||
copy_dir_all(
|
||||
&sources[0].dir().join(PARTITIONS_SUBDIR),
|
||||
&output.join(PARTITIONS_SUBDIR),
|
||||
)?;
|
||||
|
||||
if self.mode == MergeMode::Presence {
|
||||
remove_dirs_named(&output, "counts")?;
|
||||
}
|
||||
pb.finish_and_clear();
|
||||
rep.push(t.stop());
|
||||
self.reporter.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)?;
|
||||
copy_spectrums(src.dir(), &output, &old_labels, new_labels)?;
|
||||
}
|
||||
pb.finish_and_clear();
|
||||
rep.push(t.stop());
|
||||
self.reporter.push(t.stop());
|
||||
|
||||
// ── Open destination ──────────────────────────────────────────────────
|
||||
let dst = KmerIndex::open(output)?;
|
||||
let n_partitions = dst.n_partitions();
|
||||
let n_dst_genomes = src_genomes[0].len();
|
||||
|
||||
@@ -199,11 +257,9 @@ impl KmerIndex {
|
||||
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;
|
||||
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..])
|
||||
@@ -224,30 +280,27 @@ impl KmerIndex {
|
||||
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 mode = self.mode;
|
||||
|
||||
let runner = PartitionRunner::new();
|
||||
let mut part_stats: Vec<PartStat> = Vec::with_capacity(n_partitions);
|
||||
let mut done: u64 = 0;
|
||||
let on_progress = &mut self.on_progress;
|
||||
|
||||
runner
|
||||
.run(
|
||||
&order,
|
||||
|i| {
|
||||
dst.merge_partition(
|
||||
i,
|
||||
srcs,
|
||||
mode,
|
||||
n_dst_genomes,
|
||||
block_bits,
|
||||
evidence,
|
||||
)
|
||||
},
|
||||
|i| merge_partition(&dst, i, srcs, mode, n_dst_genomes, block_bits, evidence),
|
||||
|i, g_len, dur| {
|
||||
pb.inc(1);
|
||||
done += 1;
|
||||
if let Some(cb) = on_progress.as_mut() {
|
||||
cb(Progress {
|
||||
position: done,
|
||||
total: Some(n_partitions as u64),
|
||||
});
|
||||
}
|
||||
debug!(
|
||||
"partition {i}: done in {:.1}s — {} new kmers",
|
||||
dur.as_secs_f64(),
|
||||
@@ -260,29 +313,19 @@ impl KmerIndex {
|
||||
});
|
||||
},
|
||||
)
|
||||
.map_err(OKIError::Partition)?;
|
||||
.map_err(|e| io::Error::other(e.to_string()))?;
|
||||
|
||||
pb.finish_and_clear();
|
||||
|
||||
// ── Diagnostic report ─────────────────────────────────────────────
|
||||
print_merge_partition_report(&part_stats, runner.max_workers());
|
||||
|
||||
rep.push(t.stop());
|
||||
self.reporter.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());
|
||||
}
|
||||
// ── Finalize: pack matrices, mark indexed, reopen ──────────────────────
|
||||
let t = Stage::start("finalize");
|
||||
let dst = KmerIndex::finalize_indexed(&output, &mut self.reporter)?;
|
||||
self.reporter.push(t.stop());
|
||||
|
||||
KmerIndex::open(output)
|
||||
Ok(dst)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -339,7 +382,9 @@ fn fmt_bytes(b: u64) -> String {
|
||||
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);
|
||||
let Ok(p) = src.layer_unitigs_path(i, l) else {
|
||||
break;
|
||||
};
|
||||
if !p.exists() {
|
||||
break;
|
||||
}
|
||||
@@ -353,7 +398,7 @@ fn partition_unitig_bytes(src: &KmerIndex, i: usize) -> u64 {
|
||||
fn compute_labels(
|
||||
src_genomes: &[Vec<GenomeInfo>],
|
||||
rename_duplicates: bool,
|
||||
) -> OKIResult<(Vec<Vec<String>>, Vec<GenomeInfo>)> {
|
||||
) -> io::Result<(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();
|
||||
@@ -368,7 +413,10 @@ fn compute_labels(
|
||||
} else if rename_duplicates {
|
||||
format!("{label}.{count}")
|
||||
} else {
|
||||
return Err(OKIError::DuplicateGenomeLabel(label.clone()));
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::AlreadyExists,
|
||||
format!("duplicate genome label: {label}"),
|
||||
));
|
||||
};
|
||||
*count += 1;
|
||||
labels.push(new_label.clone());
|
||||
@@ -437,7 +485,7 @@ fn choose_base(sources: &[&KmerIndex], src_genomes: &[Vec<GenomeInfo>], mode: Me
|
||||
let needs_approx = sources.iter().enumerate().any(|(i, src)| {
|
||||
!is_trivial(&src_genomes[i], mode)
|
||||
&& matches!(
|
||||
src.meta.config.evidence,
|
||||
src.meta().config.evidence,
|
||||
IndexMode::Approx { .. } | IndexMode::Hybrid { .. }
|
||||
)
|
||||
});
|
||||
@@ -448,7 +496,7 @@ fn choose_base(sources: &[&KmerIndex], src_genomes: &[Vec<GenomeInfo>], mode: Me
|
||||
.filter(|(_, src)| {
|
||||
!needs_approx
|
||||
|| matches!(
|
||||
src.meta.config.evidence,
|
||||
src.meta().config.evidence,
|
||||
IndexMode::Approx { .. } | IndexMode::Hybrid { .. }
|
||||
)
|
||||
})
|
||||
|
||||
@@ -0,0 +1,645 @@
|
||||
//! Merging a source partition's new layer into a destination partition's
|
||||
//! de Bruijn graph union (pass 1) then column fill (pass 2). Free function,
|
||||
//! not a `KmerIndex` method — `KmerIndex` stays the stateless
|
||||
//! `Index { Partition { Layer } }` data model; this is one of the
|
||||
//! algorithms that reads/writes through it, living here in `obikmerge`
|
||||
//! instead.
|
||||
|
||||
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 obikindex::layer::{IndexMode, KmerLayer, TypedLayer};
|
||||
use obikindex::{ColBuilder, KmerIndex, OKIError, OKIResult};
|
||||
use obikindexer::{build_graph, materialize_layer};
|
||||
use obikseq::CanonicalKmer;
|
||||
use obiskio::UnitigFileReader;
|
||||
|
||||
use crate::MergeMode;
|
||||
|
||||
// ── 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]);
|
||||
}
|
||||
}
|
||||
|
||||
// ── dst cross-layer lookup ──────────────────────────────────────────────────
|
||||
|
||||
/// Batch membership lookup across `dst_layers` — first layer that carries
|
||||
/// each kmer wins. Groups hits by layer index (as `(kmer batch index,
|
||||
/// slot)` pairs) rather than returning them in query order, so pass 2 can
|
||||
/// follow with one column-major write per `(layer, column)` instead of
|
||||
/// interleaving across layers kmer by kmer. The kmers with no hit in any
|
||||
/// layer come back as a separate index list, for the new-layer MPHF batch
|
||||
/// lookup on a miss.
|
||||
fn find_batch_grouped(
|
||||
layers: &[KmerLayer],
|
||||
kmers: &[CanonicalKmer],
|
||||
) -> (Vec<Vec<(usize, usize)>>, Vec<usize>) {
|
||||
let mut by_layer: Vec<Vec<(usize, usize)>> = vec![Vec::new(); layers.len()];
|
||||
let mut not_found: Vec<usize> = Vec::new();
|
||||
for (idx, &kmer) in kmers.iter().enumerate() {
|
||||
match layers
|
||||
.iter()
|
||||
.enumerate()
|
||||
.find_map(|(li, l)| l.find_slot(kmer).map(|slot| (li, slot)))
|
||||
{
|
||||
Some((li, slot)) => by_layer[li].push((idx, slot)),
|
||||
None => not_found.push(idx),
|
||||
}
|
||||
}
|
||||
(by_layer, not_found)
|
||||
}
|
||||
|
||||
// ── merge_partition ────────────────────────────────────────────────────────
|
||||
|
||||
/// Merge `sources` into destination partition `i`.
|
||||
///
|
||||
/// Each entry in `sources` is `(index, 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 the bootstrap source before this
|
||||
/// call).
|
||||
pub(crate) fn merge_partition(
|
||||
dst: &KmerIndex,
|
||||
i: usize,
|
||||
sources: &[(&KmerIndex, usize)],
|
||||
mode: MergeMode,
|
||||
n_dst_genomes: usize,
|
||||
block_bits: u8,
|
||||
evidence: &IndexMode,
|
||||
) -> OKIResult<usize> {
|
||||
let dst_partition = dst.partition(i)?;
|
||||
let n_dst_layers = dst_partition.n_layers();
|
||||
let n_src_total: usize = sources.iter().map(|(_, n)| *n).sum();
|
||||
|
||||
// First merge in presence mode: materialize implicit single-genome
|
||||
// presence matrices on existing layers (all slots true — every kmer in
|
||||
// those layers belongs to genome_0) so a resumed builder can append
|
||||
// source columns onto them below.
|
||||
if n_dst_genomes == 1 && mode == MergeMode::Presence {
|
||||
for l in 0..n_dst_layers {
|
||||
let layer = dst_partition.layer(l)?.open()?;
|
||||
TypedLayer::<()>::init_presence_matrix(layer.dir(), layer.n())?;
|
||||
}
|
||||
}
|
||||
|
||||
// dst's `0..n_dst_layers` layers and every source's own layers, for
|
||||
// this one partition — all opened the same direct way, not via
|
||||
// `IndexCache`. Every layer read here (dst's pre-existing ones, and
|
||||
// every source's) is just as stable for the duration of this call: dst
|
||||
// only ever gets a brand new layer *appended* below, never one of these
|
||||
// read from, so there's no "dst isn't stable" distinction between the
|
||||
// two loops — `IndexCache` still doesn't fit either one, for a single,
|
||||
// uniform reason: `obikindexer::build_graph`/`obipipeline::WorkerPool`
|
||||
// (pass 1/2 below) are plain-OS-thread pipelines requiring `Send + Sync
|
||||
// + 'static`, and `IndexCache<'_>` always borrows `&KmerIndex` — no
|
||||
// amount of logical stability changes that a borrowed cache can't cross
|
||||
// that boundary, only owned `KmerLayer`s can (each owns its own mmap'd
|
||||
// files, no reference back to the index it came from).
|
||||
let dst_layers: Vec<KmerLayer> = (0..n_dst_layers)
|
||||
.map(|l| dst_partition.layer(l)?.open())
|
||||
.collect::<OKIResult<_>>()?;
|
||||
let dst_layers = Arc::new(dst_layers);
|
||||
|
||||
let src_layers: Vec<Vec<Arc<KmerLayer>>> = sources
|
||||
.iter()
|
||||
.map(|(src, _)| -> OKIResult<Vec<Arc<KmerLayer>>> {
|
||||
let partition = src.partition(i)?;
|
||||
(0..partition.n_layers())
|
||||
.map(|l| partition.layer(l)?.open().map(Arc::new))
|
||||
.collect()
|
||||
})
|
||||
.collect::<OKIResult<_>>()?;
|
||||
|
||||
// ── Pass 1: pipeline — parallel file read + dst 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_layers — Arc<[KmerLayer]>: Sync
|
||||
// Sink : push new kmers into GraphDeBruijn (single thread, no locks needed)
|
||||
|
||||
let mut unitig_paths: Vec<PathBuf> = Vec::new();
|
||||
for layers in &src_layers {
|
||||
for layer in layers {
|
||||
let p = layer.unitigs_path();
|
||||
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_layers);
|
||||
|
||||
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.iter().all(|l| l.find_slot(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 = KmerLayer::new(&dst_partition, new_layer_idx)
|
||||
.dir()
|
||||
.to_path_buf();
|
||||
|
||||
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<obikindex::layer::MphfOnly>> = if any_new {
|
||||
Some(Arc::new(obikindex::layer::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 = dst_layers[l].dir().to_path_buf();
|
||||
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 errors propagate via ?
|
||||
// before any worker thread is spawned.
|
||||
let mut pass2_items: Vec<(usize, usize, PathBuf, Arc<KmerLayer>)> = Vec::new();
|
||||
{
|
||||
let mut col_offset = 0usize;
|
||||
for ((_, src_n), layers) in sources.iter().zip(src_layers.iter()) {
|
||||
for layer in layers {
|
||||
let unitigs_path = layer.unitigs_path();
|
||||
if unitigs_path.exists() {
|
||||
pass2_items.push((col_offset, *src_n, unitigs_path, Arc::clone(layer)));
|
||||
}
|
||||
}
|
||||
col_offset += src_n;
|
||||
}
|
||||
}
|
||||
|
||||
enum Pass2Data {
|
||||
SrcLayer((usize, usize, PathBuf, Arc<KmerLayer>, ThrottleGuard)),
|
||||
RawBatch((usize, usize, Arc<KmerLayer>, Vec<CanonicalKmer>)),
|
||||
/// One entry per `(dst layer, absolute column)` touched by this
|
||||
/// batch, each carrying every `(slot, value)` pair for that column
|
||||
/// together — column-major, so the sink below locks that column's
|
||||
/// `ColBuilder` once per batch and writes the whole run, instead of
|
||||
/// once per `(kmer, genome)` pair.
|
||||
WriteBatch(Vec<(Option<usize>, usize, Vec<(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_layers_t2 = Arc::clone(&dst_layers);
|
||||
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, unitigs_path, src_layer) = t.item;
|
||||
(col_offset, src_n, unitigs_path, src_layer, 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, unitigs_path, src_layer, _guard)) =
|
||||
data
|
||||
{
|
||||
// _guard dropped at end of block, releasing the slot.
|
||||
// `src_layer` (MPHF + matrix) was already opened up
|
||||
// front — only the unitigs.bin stream itself is
|
||||
// opened (and throttled) here.
|
||||
let reader = match UnitigFileReader::open_sequential(&unitigs_path) {
|
||||
Ok(r) => r,
|
||||
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_layer),
|
||||
b,
|
||||
))))
|
||||
.ok();
|
||||
count += 1;
|
||||
}
|
||||
}
|
||||
if !batch.is_empty() {
|
||||
push.send(Ok(Pass2Data::RawBatch((col_offset, src_n, src_layer, batch))))
|
||||
.ok();
|
||||
count += 1;
|
||||
}
|
||||
delta.send(count - 1).ok();
|
||||
}
|
||||
},
|
||||
) as SharedFlatFn<Pass2Data>),
|
||||
make_transform!(
|
||||
Pass2Data,
|
||||
{
|
||||
move |(col_offset, src_n, src_layer, kmers): (
|
||||
usize,
|
||||
usize,
|
||||
Arc<KmerLayer>,
|
||||
Vec<CanonicalKmer>,
|
||||
)|
|
||||
-> Vec<(Option<usize>, usize, Vec<(usize, u32)>)> {
|
||||
// Column-first, one pass for the whole batch: a single
|
||||
// MPHF hash sweep (`hash_batch`) instead of one hash
|
||||
// per kmer, then one read pass per source genome
|
||||
// column (`fill_sub_matrix`) instead of one row read
|
||||
// per kmer — matches the underlying matrix's own
|
||||
// column-major layout.
|
||||
let slots = src_layer.hash_batch(&kmers);
|
||||
let mut cols: Vec<Vec<u32>> =
|
||||
(0..src_n).map(|_| Vec::with_capacity(kmers.len())).collect();
|
||||
src_layer.fill_sub_matrix(&slots, &mut cols);
|
||||
|
||||
// Membership against dst, grouped by dst layer
|
||||
// instead of interleaved kmer by kmer — same
|
||||
// grouping `PartitionCache::find_presence_batch`
|
||||
// does in `obikphylo`, via the shared primitive.
|
||||
let (by_dst_layer, not_found) = find_batch_grouped(&dst_layers_t2, &kmers);
|
||||
|
||||
// Genome-major within each layer group (not kmer-major)
|
||||
// so every `(layer, col)` pair's values land together —
|
||||
// one entry per column, not one per `(kmer, genome)`.
|
||||
let mut ops: Vec<(Option<usize>, usize, Vec<(usize, u32)>)> =
|
||||
Vec::with_capacity(src_n * (by_dst_layer.len() + 1));
|
||||
for (dst_layer, hits) in by_dst_layer.into_iter().enumerate() {
|
||||
if hits.is_empty() {
|
||||
continue;
|
||||
}
|
||||
for g in 0..src_n {
|
||||
let col_vals: Vec<(usize, u32)> = hits
|
||||
.iter()
|
||||
.map(|&(idx, slot)| (slot, cols[g][idx]))
|
||||
.collect();
|
||||
ops.push((Some(dst_layer), col_offset + g, col_vals));
|
||||
}
|
||||
}
|
||||
if let (Some(mphf), false) = (&new_mphf_t2, not_found.is_empty()) {
|
||||
let not_found_kmers: Vec<CanonicalKmer> =
|
||||
not_found.iter().map(|&idx| kmers[idx]).collect();
|
||||
let new_slots = mphf.index_batch(¬_found_kmers);
|
||||
for g in 0..src_n {
|
||||
let col_vals: Vec<(usize, u32)> = not_found
|
||||
.iter()
|
||||
.zip(new_slots.iter())
|
||||
.map(|(&idx, &slot)| (slot, cols[g][idx]))
|
||||
.collect();
|
||||
ops.push((None, col_offset + g, col_vals));
|
||||
}
|
||||
}
|
||||
ops
|
||||
}
|
||||
},
|
||||
RawBatch,
|
||||
WriteBatch
|
||||
),
|
||||
],
|
||||
make_sink!(
|
||||
Pass2Data,
|
||||
{
|
||||
move |ops: Vec<(Option<usize>, usize, Vec<(usize, u32)>)>| {
|
||||
for (layer_opt, col, col_vals) in ops {
|
||||
let builder = match layer_opt {
|
||||
Some(l) => &exist_sink[l][col],
|
||||
None => &new_sink[col],
|
||||
};
|
||||
let mut b = builder.lock().unwrap();
|
||||
for (slot, val) in col_vals {
|
||||
b.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 ──────────────────────────────────────────────────
|
||||
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)?;
|
||||
}
|
||||
|
||||
debug!(
|
||||
"partition {i}: builders closed in {:.3}s",
|
||||
t_close.elapsed().as_secs_f64()
|
||||
);
|
||||
|
||||
Ok(n_new)
|
||||
}
|
||||
Reference in New Issue
Block a user