diff --git a/cubesys/src/store.rs b/cubesys/src/store.rs index b39752b..c6c0068 100644 --- a/cubesys/src/store.rs +++ b/cubesys/src/store.rs @@ -32,6 +32,8 @@ use std::time::Duration; use cubecoords::{CubeHeader, Czyx}; use cubestore::{CubeStore, HashBackend}; +use std::collections::HashSet; + use crate::persist; use cubecode::{Kind, Op}; @@ -54,6 +56,12 @@ impl Default for DurabilityConfig { } } +/// When the incremental delta file grows past this many bytes, fold it into a +/// fresh full base snapshot and truncate the delta. Keeps on-disk state +/// bounded and fast to load, while still doing tiny incremental writes the +/// vast majority of the time. 1 MiB is comfortable for this store. +const DELTA_COMPACT_BYTES: u64 = 1_048_576; + /// One durable transaction in the WAL (newline-delimited JSON, one per line). /// `seq` is monotonic so replay can skip already-checkpointed entries. For /// `op:"put"`, `v` is the raw backend value bytes (header_len||hdr||body) @@ -72,7 +80,7 @@ enum WalOp { Delete, } -/// The write-ahead log. Append-only NDJSON, group-fsynced on a timer. +/// The write-ahead log. Append-only NDJSON, group-fsync'd on a timer. struct Wal { path: PathBuf, /// `None` in memory mode (no durability); `Some` when a file backs it. @@ -83,6 +91,16 @@ struct Wal { pending: Mutex, /// Highest sequence number known to be fsync'd to disk. committed_seq: AtomicU64, + /// WAL sequence number already folded into the base snapshot. Replay on + /// startup begins after this, so the WAL only carries post-checkpoint + /// transactions. + cp_seq_wal: AtomicU64, + /// WAL sequence number of the last full base snapshot. Used to compute + /// the incremental delta (only coords written after `base_seq`). + base_seq: AtomicU64, + /// Set of coordinates written since `base_seq` (the authoritative changelog + /// for the incremental delta). Cleared whenever `base_seq` advances. + dirty: Mutex>, /// Shared stop flag (also owned by the group-commit thread). stop: Arc, /// Handle of the group-commit fsync thread (behind a Mutex so the struct @@ -123,6 +141,9 @@ impl Wal { seq: AtomicU64::new(start_seq), pending: Mutex::new(String::new()), committed_seq: AtomicU64::new(max_seq.max(cp_seq)), + cp_seq_wal: AtomicU64::new(cp_seq), + base_seq: AtomicU64::new(cp_seq), + dirty: Mutex::new(HashSet::new()), stop: stop.clone(), group_thread: Mutex::new(None), fsync_ms, @@ -140,12 +161,38 @@ impl Wal { seq: AtomicU64::new(1), pending: Mutex::new(String::new()), committed_seq: AtomicU64::new(0), + cp_seq_wal: AtomicU64::new(0), + base_seq: AtomicU64::new(0), + dirty: Mutex::new(HashSet::new()), stop: Arc::new(AtomicBool::new(true)), group_thread: Mutex::new(None), fsync_ms: 0, }) } + /// Mark `coord` dirty (written after `base_seq`), for the incremental delta. + fn mark_dirty(&self, coord: Czyx) { + self.dirty.lock().unwrap().insert(coord); + } + + /// True if `coord` was written after `base_seq` (i.e. belongs in the delta). + fn coord_modified_since(&self, coord: &Czyx, base: u64) -> bool { + if base == 0 { + // No base yet — everything is "since base" only if it was written. + return self.dirty.lock().unwrap().contains(coord); + } + self.dirty.lock().unwrap().contains(coord) + } + + fn set_cp_seq_wal(&self, v: u64) { + self.cp_seq_wal.store(v, Ordering::SeqCst); + } + + fn set_base_seq(&self, v: u64) { + self.base_seq.store(v, Ordering::SeqCst); + self.dirty.lock().unwrap().clear(); + } + /// Append a transaction. In memory mode this is a no-op. fn append(&self, op: WalOp, coord: Czyx, value: Vec) { let has_file = self.file.lock().unwrap().is_some(); @@ -153,6 +200,7 @@ impl Wal { return; } let s = self.seq.fetch_add(1, Ordering::SeqCst) + 1; + self.mark_dirty(coord); let line = encode_wal(&WalEntry { seq: s, op, @@ -212,17 +260,13 @@ impl Wal { if need > committed { self.committed_seq.store(need, Ordering::SeqCst); } + // The base snapshot already reflected up to cp_seq; replay covered the + // rest. Record the boundary so future checkpoints compute deltas and + // WAL replay starts after it. + self.set_cp_seq_wal(cp_seq); applied } - /// Truncate the WAL after a checkpoint (future entries start fresh). - fn rotate(&self) { - if let Some(f) = self.file.lock().unwrap().as_mut() { - let _ = f.set_len(0); - let _ = f.sync_all(); - } - } - /// Signal the group-commit thread to stop. fn set_stop(&self) { self.stop.store(true, Ordering::SeqCst); @@ -276,42 +320,37 @@ impl ConcurrentStore { recovery_log: &str, cfg: DurabilityConfig, ) -> std::io::Result { + // Load the checkpoint: base snapshot + incremental delta. Then replay + // any WAL entries newer than the base's WAL boundary. let db_p = PathBuf::from(db_path); let wal_p = PathBuf::from(wal_path); + // MUST match `checkpoint_store`'s `db_path.with_extension("delta")` so the + // delta is read back from the same path it is written to on checkpoint. + let delta_p = db_p.with_extension("delta"); let cp_p = PathBuf::from(format!("{db_path}.seq")); let rec_p = PathBuf::from(recovery_log); if let Some(p) = db_p.parent() { let _ = fs::create_dir_all(p); } - if let Some(p) = wal_p.parent() { + if let Some(p) = delta_p.parent() { let _ = fs::create_dir_all(p); } - // Highest sequence already durable in the last checkpoint. - let cp_seq: u64 = fs::read_to_string(&cp_p) + // Highest WAL sequence already folded into the base. + let base_cp: u64 = fs::read_to_string(&cp_p) .ok() .and_then(|s| s.trim().parse().ok()) .unwrap_or(0); - // Load the checkpoint snapshot into RAM. - let mut store = CubeStore::new(HashBackend::new()); - if db_p.exists() { - match fs::read_to_string(&db_p) { - Ok(json) => { - if let Err(e) = persist::load_into_store(&mut store, &json) { - eprintln!("cube-store: checkpoint load failed ({e}); recovering from WAL"); - } - } - Err(e) => eprintln!("cube-store: cannot read {db_path}: {e}; starting empty"), - } - } + // Reconstruct the store: base + delta. + let mut store = load_base_plus_delta(&db_p, &delta_p); let stop = Arc::new(AtomicBool::new(false)); - let wal = Wal::open(&wal_p, cp_seq, cfg.wal_fsync_ms, stop.clone())?; + let wal = Wal::open(&wal_p, base_cp, cfg.wal_fsync_ms, stop.clone())?; - // Replay WAL entries newer than the checkpoint. If we had to rely on + // Replay WAL entries newer than the base boundary. If we had to rely on // the WAL, record that fact in the recovery log (in writing). - let applied = wal.replay_after(cp_seq, &mut store); + let applied = wal.replay_after(base_cp, &mut store); if applied > 0 { append_recovery_log(&rec_p, applied); } @@ -511,14 +550,100 @@ impl Drop for ConcurrentStore { } } -/// Write the full store to `db_path` (atomic rename) and rotate the WAL. -/// `cp_seq` recorded is the highest sequence assigned, because the snapshot -/// covers the entire in-RAM store (which always includes every put so far). +/// Persistent checkpoint state lives in three files: +/// * `db_path` — the full base snapshot (NDJSON array). Written rarely. +/// * `db_path.delta` — incremental deltas since the base (one `{c,z,y,x,v}` object per modified coordinate, no seq). Written on every checkpoint. +/// * `db_path.seq` — the WAL sequence number already folded into the base (so WAL replay starts after it). +/// * `wal_path` — the WAL (transactions newer than `db_path.seq`). +/// +/// On startup we load the base, apply the delta, then replay WAL entries newer than `db_path.seq`. The `cp_seq_wal` field tracks that boundary so replay knows where to begin. This gives tiny per-checkpoint disk writes (only changed coords) while keeping a full snapshot for fast cold load. fn checkpoint_store( inner: &Arc>>, db_path: &Path, cp_seq_path: &Path, wal: &Arc, +) { + wal.flush_pending(); + + // Collect only the coordinates that changed since the last base/checkpoint. + let delta = incremental(inner, wal); + + let delta_path = db_path.with_extension("delta"); + + // If we have no delta at all (nothing changed), just re-stamp the WAL + // boundary and leave the base untouched. + if delta.is_empty() { + persist_seq(cp_seq_path, wal.seq.load(Ordering::SeqCst)); + wal.set_cp_seq_wal(wal.seq.load(Ordering::SeqCst)); + return; + } + + // Append the delta. If the delta file is getting large, fold everything + // back into a fresh full base snapshot instead (compaction). If there is + // no base snapshot yet at all, we must write the full base now (the delta + // alone is not a complete store). + let delta_len = fs::read(&delta_path).map(|b| b.len() as u64).unwrap_or(0); + let base_exists = db_path.exists(); + if !base_exists || delta_len + delta.len() as u64 > DELTA_COMPACT_BYTES { + fold_delta_into_base(inner, db_path, cp_seq_path, &delta_path, wal); + } else if let Some(parent) = delta_path.parent() { + let _ = fs::create_dir_all(parent); + if let Ok(mut f) = OpenOptions::new() + .create(true) + .append(true) + .open(&delta_path) + { + let mut buf = String::new(); + for e in &delta { + buf.push_str(&format!( + "{{\"c\":{},\"z\":{},\"y\":{},\"x\":{},\"v\":\"{}\"}}\n", + e.0.c, + e.0.z, + e.0.y, + e.0.x, + to_hex(&e.1) + )); + } + if f.write_all(buf.as_bytes()).is_ok() && f.flush().is_ok() && f.sync_all().is_ok() { + persist_seq(cp_seq_path, wal.seq.load(Ordering::SeqCst)); + wal.set_cp_seq_wal(wal.seq.load(Ordering::SeqCst)); + wal.set_base_seq(wal.seq.load(Ordering::SeqCst)); + } + } + } +} + +/// Snapshot the current delta (only coordinates modified since the last base +/// snapshot) as a list of `(coord, raw_value)` pairs. +fn incremental(inner: &Arc>>, wal: &Arc) -> Vec<(Czyx, Vec)> { + let g = inner.lock().unwrap(); + let base = wal.base_seq.load(Ordering::SeqCst); + let now = wal.seq.load(Ordering::SeqCst); + if now <= base { + return Vec::new(); + } + let mut out = Vec::new(); + for k in g.keys() { + // Only include coords whose last write is newer than the base + // snapshot. We approximate "changed since base" by tracking the set + // of coords written after `base` in the WAL (the authoritative changelog). + if wal.coord_modified_since(&k, base) { + if let Some(v) = g.get_raw(&k) { + out.push((k, v)); + } + } + } + out +} + +/// Fold the current delta + live store into a fresh full base snapshot and +/// truncate the delta. Used when the delta has grown too large. +fn fold_delta_into_base( + inner: &Arc>>, + db_path: &Path, + cp_seq_path: &Path, + delta_path: &Path, + wal: &Arc, ) { let snap = inner.lock().unwrap().clone(); let json = persist::dump_store(&snap); @@ -531,14 +656,54 @@ fn checkpoint_store( let _ = fs::rename(&tmp, db_path); } } - // Record the checkpoint sequence (max assigned), so replay on next start - // only recovers entries strictly newer than this snapshot. + // Delta is now fully represented by the base; truncate it. + let _ = fs::write(delta_path, b""); let seq = wal.seq.load(Ordering::SeqCst); - if let Ok(mut f) = File::create(cp_seq_path) { + persist_seq(cp_seq_path, seq); + wal.set_cp_seq_wal(seq); + wal.set_base_seq(seq); +} + +/// Read `db_path` (full base) then apply the delta file; returns the +/// reconstructed store. If either is missing it is simply skipped. +fn load_base_plus_delta(db_path: &Path, delta_path: &Path) -> CubeStore { + let mut store = CubeStore::new(HashBackend::new()); + if db_path.exists() { + if let Ok(json) = fs::read_to_string(db_path) { + if let Err(e) = persist::load_into_store(&mut store, &json) { + eprintln!("cube-store: base load failed ({e}); relying on WAL/delta"); + } + } + } + if delta_path.exists() { + if let Ok(text) = fs::read_to_string(delta_path) { + for line in text.lines() { + let line = line.trim(); + if line.is_empty() || !line.starts_with('{') { + continue; + } + let c = field_u8(line, "c").unwrap_or(0); + let z = field_u8(line, "z").unwrap_or(0); + let y = field_u8(line, "y").unwrap_or(0); + let x = field_u8(line, "x").unwrap_or(0); + let v = field_str(line, "v"); + if let Ok(value) = from_hex(v) { + store.put_raw(Czyx::new(c, z, y, x), value); + } + } + } + } + store +} + +fn persist_seq(path: &Path, seq: u64) { + if let Some(p) = path.parent() { + let _ = fs::create_dir_all(p); + } + if let Ok(mut f) = File::create(path) { let _ = f.write_all(seq.to_string().as_bytes()); let _ = f.sync_all(); } - wal.rotate(); } /// Append a recovery event to `path` (NDJSON, fsync'd). Called whenever we @@ -690,6 +855,7 @@ mod tests { let _ = fs::remove_file(format!("{db}.wal")); let _ = fs::remove_file(format!("{db}.recovery.ndjson")); let _ = fs::remove_file(format!("{db}.seq")); + let _ = fs::remove_file(PathBuf::from(db).with_extension("delta")); } #[test] @@ -814,6 +980,92 @@ mod tests { assert_eq!(s.keys().len(), 256); } + #[test] + fn incremental_checkpoint_delta_model() { + let db = tmp("db3.json"); + let wal = format!("{db}.wal"); + let rec = format!("{db}.recovery.ndjson"); + let delta = PathBuf::from(&db) + .with_extension("delta") + .to_string_lossy() + .into_owned(); + cleanup(&db); + + // Phase 1: 10 records, then a checkpoint. The base snapshot is written + // and the delta starts empty (everything is now in the base). + { + let s = ConcurrentStore::open( + &db, + &wal, + &rec, + DurabilityConfig { + checkpoint_ms: 50, + wal_fsync_ms: 10, + }, + ) + .unwrap(); + for i in 0..10u8 { + s.put_record(Czyx::new(1, 1, 1, i), &CubeHeader::new(), &[i]); + } + s.checkpoint(); + } + let base_after_1 = fs::read(&db).map(|b| b.len()).unwrap_or(0); + let delta_after_1 = fs::read(&delta).map(|b| b.len()).unwrap_or(0); + assert!( + base_after_1 > 0, + "base snapshot should exist after first checkpoint" + ); + assert_eq!( + delta_after_1, 0, + "delta should be empty right after a full checkpoint" + ); + + // Phase 2: change only 2 records. The next checkpoint must write only + // those 2 into the delta (NOT rewrite the whole base), and a reopen + // must still see all 10. + { + let s = ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(); + s.put_record(Czyx::new(1, 1, 1, 0), &CubeHeader::new(), b"changed"); + s.put_record(Czyx::new(1, 1, 1, 1), &CubeHeader::new(), b"changed"); + s.checkpoint(); + } + let base_after_2 = fs::read(&db).map(|b| b.len()).unwrap_or(0); + let delta_lines_2 = fs::read_to_string(&delta) + .map(|t| t.lines().filter(|l| !l.trim().is_empty()).count()) + .unwrap_or(0); + // Base should NOT have grown by a full rewrite of 10 records; the delta + // must carry exactly the 2 changed coords. + assert!( + base_after_2 <= base_after_1 + 64, + "base should not be fully rewritten on a small change" + ); + assert_eq!( + delta_lines_2, 2, + "delta should contain only the 2 changed coords" + ); + + // Reopen from base + delta and confirm all 10 are present with the + // updated values for the 2 we changed. + { + let s = ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(); + assert_eq!( + s.keys().len(), + 10, + "all 10 records must survive base+delta reopen" + ); + assert_eq!( + s.get_record(&Czyx::new(1, 1, 1, 0)).unwrap().1, + b"changed".to_vec() + ); + assert_eq!( + s.get_record(&Czyx::new(1, 1, 1, 1)).unwrap().1, + b"changed".to_vec() + ); + assert_eq!(s.get_record(&Czyx::new(1, 1, 1, 9)).unwrap().1, vec![9u8]); + } + cleanup(&db); + } + #[test] fn query_doc_type_works() { let s = ConcurrentStore::memory();