Add resume API to persistent matrix builders for incremental appending

Introduce a `resume` method for persistent bit and int matrix builders to restore dimensions from persisted metadata. This enables incremental column appending across separate build sessions without requiring manual dimension tracking. Consolidate builder lifecycle management in the merge layer using a unified `MatrixBuilder` enum, simplify closure logic, and add instrumentation. Include tests verifying state preservation and data integrity across multiple resume cycles.
This commit is contained in:
Eric Coissac
2026-08-20 12:36:47 +02:00
parent 308d2b9f92
commit dc3d82f8db
6 changed files with 305 additions and 115 deletions
@@ -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<Self> {
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 }
+12
View File
@@ -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<Self> {
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]
+47
View File
@@ -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();
+47
View File
@@ -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();
+1 -20
View File
@@ -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<Partition
}
}
// ── path helpers ──────────────────────────────────────────────────────────────
pub(crate) fn col_path_bit(dir: &Path, col: usize) -> 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 {
+187 -95
View File
@@ -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<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]);
}
}
// ── 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<ColBuilder> = 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<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)?;
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<ColBuilder> {
let b = PersistentBitVecBuilder::new(
n_new,
&col_path_bit(&data_dir, n_dst_genomes + g),
)?;
Ok(ColBuilder::Bit(b))
})
.collect::<SKResult<_>>()?
}
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<ColBuilder> {
let b = PersistentCompactIntVecBuilder::new(
n_new,
&col_path_int(&data_dir, n_dst_genomes + g),
)?;
Ok(ColBuilder::Int(b))
})
.collect::<SKResult<_>>()?
}
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::<SKResult<Vec<_>>>()?;
(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<Vec<ColBuilder>> = (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<ColBuilder> {
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::<SKResult<_>>()
})
.collect::<SKResult<_>>()?;
// 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_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::<SKResult<Vec<_>>>()?;
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;