diff --git a/src/obicompactvec/src/bitmatrix/builder.rs b/src/obicompactvec/src/bitmatrix/builder.rs index 1f76d225..4cb85f1b 100644 --- a/src/obicompactvec/src/bitmatrix/builder.rs +++ b/src/obicompactvec/src/bitmatrix/builder.rs @@ -23,6 +23,17 @@ impl PersistentBitMatrixBuilder { Ok(Self { dir: dir.to_path_buf(), n, n_cols: 0 }) } + /// Resume appending columns to a Columnar matrix directory already + /// closed by a previous builder session — reads `n`/`n_cols` from + /// `meta.json` instead of starting a fresh matrix at `n_cols == 0`. + /// Lets callers that build columns incrementally over time (e.g. a + /// streamed merge pipeline) reach `add_col`/`close` without knowing the + /// on-disk column naming or `MatrixMeta` themselves. + pub fn resume(dir: &Path) -> io::Result { + let meta = MatrixMeta::load(dir)?; + Ok(Self { dir: dir.to_path_buf(), n: meta.n, n_cols: meta.n_cols }) + } + pub fn n(&self) -> usize { self.n } pub fn n_cols(&self) -> usize { self.n_cols } diff --git a/src/obicompactvec/src/intmatrix.rs b/src/obicompactvec/src/intmatrix.rs index 4680bdb5..4a2015ba 100644 --- a/src/obicompactvec/src/intmatrix.rs +++ b/src/obicompactvec/src/intmatrix.rs @@ -481,6 +481,18 @@ impl PersistentCompactIntMatrixBuilder { fs::create_dir_all(dir)?; Ok(Self { dir: dir.to_path_buf(), n, n_cols: 0 }) } + + /// Resume appending columns to a Columnar matrix directory already + /// closed by a previous builder session — reads `n`/`n_cols` from + /// `meta.json` instead of starting a fresh matrix at `n_cols == 0`. + /// Lets callers that build columns incrementally over time (e.g. a + /// streamed merge pipeline) reach `add_col`/`close` without knowing the + /// on-disk column naming or `MatrixMeta` themselves. + pub fn resume(dir: &Path) -> io::Result { + let meta = MatrixMeta::load(dir)?; + Ok(Self { dir: dir.to_path_buf(), n: meta.n, n_cols: meta.n_cols }) + } + #[inline] pub fn n(&self) -> usize { self.n } #[inline] diff --git a/src/obicompactvec/src/tests/bitmatrix.rs b/src/obicompactvec/src/tests/bitmatrix.rs index 7600ac39..ca7002db 100644 --- a/src/obicompactvec/src/tests/bitmatrix.rs +++ b/src/obicompactvec/src/tests/bitmatrix.rs @@ -62,6 +62,53 @@ fn two_cols_roundtrip() { assert_eq!(&*m.row(2), &[true, false]); } +#[test] +fn resume_continues_n_cols_and_appends_columns() { + let dir = tempdir().unwrap(); + let presence = dir.path().join("presence"); + + let mut b = PersistentBitMatrixBuilder::new(3, &presence).unwrap(); + let mut col0 = b.add_col().unwrap(); + col0.set(0, true); col0.set(1, false); col0.set(2, true); + col0.close().unwrap(); + b.close().unwrap(); + + let mut resumed = PersistentBitMatrixBuilder::resume(&presence).unwrap(); + assert_eq!(resumed.n(), 3); + assert_eq!(resumed.n_cols(), 1); + let mut col1 = resumed.add_col().unwrap(); + col1.set(0, false); col1.set(1, true); col1.set(2, false); + col1.close().unwrap(); + resumed.close().unwrap(); + + let m = PersistentBitMatrix::open(dir.path()).unwrap(); + assert_eq!(m.n_cols(), 2); + assert_eq!(&*m.row(0), &[true, false]); + assert_eq!(&*m.row(1), &[false, true]); + assert_eq!(&*m.row(2), &[true, false]); +} + +#[test] +fn resume_twice_keeps_appending() { + let dir = tempdir().unwrap(); + let presence = dir.path().join("presence"); + + PersistentBitMatrixBuilder::new(2, &presence).unwrap().close().unwrap(); + + for v in [true, false, true] { + let mut b = PersistentBitMatrixBuilder::resume(&presence).unwrap(); + let mut col = b.add_col().unwrap(); + col.set(0, v); col.set(1, !v); + col.close().unwrap(); + b.close().unwrap(); + } + + let m = PersistentBitMatrix::open(dir.path()).unwrap(); + assert_eq!(m.n_cols(), 3); + assert_eq!(&*m.row(0), &[true, false, true]); + assert_eq!(&*m.row(1), &[false, true, false]); +} + #[test] fn col_accessor() { let dir = tempdir().unwrap(); diff --git a/src/obicompactvec/src/tests/intmatrix.rs b/src/obicompactvec/src/tests/intmatrix.rs index 966f6da5..5af4b51e 100644 --- a/src/obicompactvec/src/tests/intmatrix.rs +++ b/src/obicompactvec/src/tests/intmatrix.rs @@ -59,6 +59,53 @@ fn two_cols_roundtrip() { assert_eq!(&*m.row(2), &[3u32, 30]); } +#[test] +fn resume_continues_n_cols_and_appends_columns() { + let dir = tempdir().unwrap(); + let counts = dir.path().join("counts"); + + let mut b = PersistentCompactIntMatrixBuilder::new(3, &counts).unwrap(); + let mut col0 = b.add_col().unwrap(); + col0.set(0, 1); col0.set(1, 2); col0.set(2, 3); + col0.close().unwrap(); + b.close().unwrap(); + + let mut resumed = PersistentCompactIntMatrixBuilder::resume(&counts).unwrap(); + assert_eq!(resumed.n(), 3); + assert_eq!(resumed.n_cols(), 1); + let mut col1 = resumed.add_col().unwrap(); + col1.set(0, 10); col1.set(1, 20); col1.set(2, 30); + col1.close().unwrap(); + resumed.close().unwrap(); + + let m = PersistentCompactIntMatrix::open(dir.path()).unwrap(); + assert_eq!(m.n_cols(), 2); + assert_eq!(&*m.row(0), &[1u32, 10]); + assert_eq!(&*m.row(1), &[2u32, 20]); + assert_eq!(&*m.row(2), &[3u32, 30]); +} + +#[test] +fn resume_twice_keeps_appending() { + let dir = tempdir().unwrap(); + let counts = dir.path().join("counts"); + + PersistentCompactIntMatrixBuilder::new(2, &counts).unwrap().close().unwrap(); + + for v in [1u32, 2, 3] { + let mut b = PersistentCompactIntMatrixBuilder::resume(&counts).unwrap(); + let mut col = b.add_col().unwrap(); + col.set(0, v); col.set(1, v * 10); + col.close().unwrap(); + b.close().unwrap(); + } + + let m = PersistentCompactIntMatrix::open(dir.path()).unwrap(); + assert_eq!(m.n_cols(), 3); + assert_eq!(&*m.row(0), &[1u32, 2, 3]); + assert_eq!(&*m.row(1), &[10u32, 20, 30]); +} + #[test] fn col_accessor() { let dir = tempdir().unwrap(); diff --git a/src/obikpartitionner/src/common.rs b/src/obikpartitionner/src/common.rs index 99e345ee..2939e0a8 100644 --- a/src/obikpartitionner/src/common.rs +++ b/src/obikpartitionner/src/common.rs @@ -1,6 +1,4 @@ -use std::fs; -use std::io; -use std::path::{Path, PathBuf}; +use std::path::Path; use obicompactvec::{PersistentBitVecBuilder, PersistentCompactIntVecBuilder}; use obilayeredmap::meta::PartitionMeta; @@ -43,23 +41,6 @@ pub(crate) fn load_meta(dir: &Path, context: &'static str) -> SKResult PathBuf { - dir.join(format!("col_{col:06}.pbiv")) -} - -pub(crate) fn col_path_int(dir: &Path, col: usize) -> PathBuf { - dir.join(format!("col_{col:06}.pciv")) -} - -pub(crate) fn write_matrix_meta(dir: &Path, n: usize, n_cols: usize) -> io::Result<()> { - fs::write( - dir.join("meta.json"), - format!("{{\"n\":{n},\"n_cols\":{n_cols}}}\n"), - ) -} - // ── ColBuilder ──────────────────────────────────────────────────────────────── pub(crate) enum ColBuilder { diff --git a/src/obikpartitionner/src/merge_layer/mod.rs b/src/obikpartitionner/src/merge_layer/mod.rs index a822c87c..c02011b5 100644 --- a/src/obikpartitionner/src/merge_layer/mod.rs +++ b/src/obikpartitionner/src/merge_layer/mod.rs @@ -8,7 +8,8 @@ //! builder setup, and pass 2), not a set of independently callable steps. use std::fs; -use std::path::PathBuf; +use std::io; +use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex}; use tracing::debug; @@ -18,16 +19,13 @@ use obipipeline::{ make_sink, make_source, make_transform, }; -use obicompactvec::{ - PersistentBitMatrix, PersistentBitMatrixBuilder, PersistentBitVecBuilder, - PersistentCompactIntMatrix, PersistentCompactIntMatrixBuilder, PersistentCompactIntVecBuilder, -}; +use obicompactvec::{PersistentBitMatrixBuilder, PersistentCompactIntMatrixBuilder}; use obikseq::CanonicalKmer; use obilayeredmap::meta::PartitionMeta; use obilayeredmap::{IndexMode, Layer, LayeredMap, MphfOnly}; use obiskio::{SKError, SKResult, UnitigFileReader}; -use crate::common::{ColBuilder, col_path_bit, col_path_int, load_meta, olm_to_sk, write_matrix_meta}; +use crate::common::{ColBuilder, load_meta, olm_to_sk}; use crate::graph_pipeline::{build_graph, materialize_layer}; use crate::partition::KmerPartition; @@ -43,6 +41,152 @@ pub enum MergeMode { 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 { + 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 { + 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 { + 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]); + } +} + // ── helpers ─────────────────────────────────────────────────────────────────── const INDEX_SUBDIR: &str = "index"; @@ -170,89 +314,48 @@ impl KmerPartition { 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) are written via append_column (all-zero/false). - // Source-genome columns are created as mutable builders for pass 2. - let new_src_builders: Vec = if any_new { + // 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, Option) = 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)?; - match mode { - MergeMode::Presence => { - PersistentBitMatrixBuilder::new(n_new, &data_dir) - .map_err(SKError::Io)? - .close() - .map_err(SKError::Io)?; - for _ in 0..n_dst_genomes { - PersistentBitMatrix::append_column(&data_dir, |_| false) - .map_err(SKError::Io)?; - } - (0..n_src_total) - .map(|g| -> SKResult { - let b = PersistentBitVecBuilder::new( - n_new, - &col_path_bit(&data_dir, n_dst_genomes + g), - )?; - Ok(ColBuilder::Bit(b)) - }) - .collect::>()? - } - MergeMode::Count => { - PersistentCompactIntMatrixBuilder::new(n_new, &data_dir) - .map_err(SKError::Io)? - .close() - .map_err(SKError::Io)?; - for _ in 0..n_dst_genomes { - PersistentCompactIntMatrix::append_column(&data_dir, |_| 0) - .map_err(SKError::Io)?; - } - (0..n_src_total) - .map(|g| -> SKResult { - let b = PersistentCompactIntVecBuilder::new( - n_new, - &col_path_int(&data_dir, n_dst_genomes + g), - )?; - Ok(ColBuilder::Int(b)) - }) - .collect::>()? - } + let mut mb = MatrixBuilder::new(mode, n_new, &data_dir).map_err(SKError::Io)?; + for _ in 0..n_dst_genomes { + mb.add_absent_col().map_err(SKError::Io)?; } + let cols = (0..n_src_total) + .map(|_| mb.add_col().map_err(SKError::Io)) + .collect::>>()?; + (cols, Some(mb)) } else { - vec![] + (vec![], None) }; let t_builders = std::time::Instant::now(); - // Builders for existing layers: n_src_total per layer. - // Columns n_dst_genomes .. n_dst_genomes + n_src_total - 1. - let exist_builders: Vec> = (0..n_dst_layers) - .map(|l| { - let layer_dir = dst_index_dir.join(format!("layer_{l}")); - let n = dst_map.layer(l).n(); - (0..n_src_total) - .map(|src_g| -> SKResult { - match mode { - MergeMode::Presence => { - let data_dir = layer_dir.join("presence"); - let b = PersistentBitVecBuilder::new( - n, - &col_path_bit(&data_dir, n_dst_genomes + src_g), - )?; - Ok(ColBuilder::Bit(b)) - } - MergeMode::Count => { - let data_dir = layer_dir.join("counts"); - let b = PersistentCompactIntVecBuilder::new( - n, - &col_path_int(&data_dir, n_dst_genomes + src_g), - )?; - Ok(ColBuilder::Int(b)) - } - } - }) - .collect::>() - }) - .collect::>()?; + // 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 = Vec::with_capacity(n_dst_layers); + let mut exist_builders: Vec> = Vec::with_capacity(n_dst_layers); + for l in 0..n_dst_layers { + let layer_dir = dst_index_dir.join(format!("layer_{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(SKError::Io)?; + let cols = (0..n_src_total) + .map(|_| mb.add_col().map_err(SKError::Io)) + .collect::>>()?; + exist_mbs.push(mb); + exist_builders.push(cols); + } debug!("partition {i}: builders ready in {:.3}s", t_builders.elapsed().as_secs_f64()); @@ -407,21 +510,15 @@ impl KmerPartition { let t_close = std::time::Instant::now(); // ── Close builders and update metadata ──────────────────────────────── - for (l, builders) in exist_locked.into_iter().enumerate() { - let layer_dir = dst_index_dir.join(format!("layer_{l}")); + 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[{l}] not uniquely owned")) + .unwrap_or_else(|_| panic!("pass2: exist_builder not uniquely owned")) .into_inner() .unwrap_or_else(|e| e.into_inner()) .close()?; } - let n = dst_map.layer(l).n(); - let data_dir = match mode { - MergeMode::Presence => layer_dir.join("presence"), - MergeMode::Count => layer_dir.join("counts"), - }; - write_matrix_meta(&data_dir, n, n_dst_genomes + n_src_total).map_err(SKError::Io)?; + mb.close().map_err(SKError::Io)?; } for b in new_locked { @@ -431,13 +528,8 @@ impl KmerPartition { .unwrap_or_else(|e| e.into_inner()) .close()?; } - if any_new { - let data_dir = match mode { - MergeMode::Presence => new_layer_dir.join("presence"), - MergeMode::Count => new_layer_dir.join("counts"), - }; - write_matrix_meta(&data_dir, n_new, n_dst_genomes + n_src_total) - .map_err(SKError::Io)?; + if let Some(mb) = new_mb { + mb.close().map_err(SKError::Io)?; let mut part_meta = PartitionMeta::load(&dst_index_dir).map_err(|e| olm_to_sk(e, "merge"))?; part_meta.n_layers = new_layer_idx + 1;