diff --git a/cubesys/src/commands.rs b/cubesys/src/commands.rs index aef3e1b..68da0e8 100644 --- a/cubesys/src/commands.rs +++ b/cubesys/src/commands.rs @@ -652,4 +652,43 @@ mod tests { ); let _ = std::fs::remove_dir_all(&dir); } + + #[test] + fn commit_replays_from_wal_without_checkpoint() { + // The daemon never calls checkpoint() between a commit and a crash; it + // relies on WAL replay. This must hold WITHOUT an explicit checkpoint() + // — i.e. the WalOp::Txn entry must round-trip through decode_wal/encode_wal. + let dir = std::env::temp_dir().join(format!("cube2-txn-wal-{}", std::process::id())); + let _ = std::fs::create_dir_all(&dir); + let db = dir.join("db.cubedb"); + let wal = dir.join("wal.ndjson"); + let rec = dir.join("recovery.jsonl"); + let _ = std::fs::remove_file(&db); + let _ = std::fs::remove_file(&wal); + let cs = ConcurrentStore::open( + db.to_str().unwrap(), + wal.to_str().unwrap(), + rec.to_str().unwrap(), + DurabilityConfig::default(), + ) + .unwrap(); + let mut s = Session::with_store(Arc::new(cs)); + s.exec("begin").unwrap(); + s.exec("prog /c005/z001/y001/x001 const 9 halt").unwrap(); + s.exec("commit").unwrap(); + drop(s); // NO checkpoint() — pure WAL replay on reopen + + let cs2 = ConcurrentStore::open( + db.to_str().unwrap(), + wal.to_str().unwrap(), + rec.to_str().unwrap(), + DurabilityConfig::default(), + ) + .unwrap(); + assert!( + cs2.get_raw(&Czyx::new(5, 1, 1, 1)).is_some(), + "committed txn must replay from WAL even without a checkpoint" + ); + let _ = std::fs::remove_dir_all(&dir); + } } diff --git a/cubesys/src/store.rs b/cubesys/src/store.rs index ac8165f..0479811 100644 --- a/cubesys/src/store.rs +++ b/cubesys/src/store.rs @@ -562,6 +562,10 @@ impl ConcurrentStore { /// Apply a transaction batch atomically: under ONE store write lock, apply /// every put/delete, then append a SINGLE `WalOp::Txn` WAL entry so the /// whole batch is durable as one unit and replays idempotently. + /// + /// The WAL is fsync'd before returning so a `COMMIT` is durable the moment + /// the caller gets control back (not merely "eventually" via the group + /// thread). This is what makes `commit` a real transaction boundary. pub fn commit_txn(&self, batch: &[TxnEntry]) { let mut g = self.inner.write().unwrap(); for e in batch { @@ -572,6 +576,7 @@ impl ConcurrentStore { } drop(g); // release the write lock before touching the WAL self.wal.append_txn(batch); + self.wal.flush_pending(); // synchronous fsync: COMMIT == durable } /// Store a code cell at `path` (the path->code bridge), durability-logged. @@ -901,6 +906,21 @@ fn decode_wal(line: &str) -> Option { } let seq = field_u64(line, "seq").ok()?; let op_s = field_str(line, "op"); + // The txn encoder emits `{"seq","op":"txn","batch":...}` with NO c/z/y/x + // (the per-entry coords live inside `batch`). Handle it before the + // c/z/y/x extraction below, which would otherwise fail and drop the entry. + if op_s == "txn" { + let raw = field_str(line, "batch"); + let bytes = from_hex(raw).ok()?; + let batch = unpack_txn(&bytes).ok()?; + return Some(WalEntry { + seq, + op: WalOp::Txn, + coord: Czyx::new(0, 0, 0, 0), // unused; replay reads batch coords + value: Vec::new(), + batch, + }); + } let c = field_u8(line, "c").ok()?; let z = field_u8(line, "z").ok()?; let y = field_u8(line, "y").ok()?; @@ -925,18 +945,6 @@ fn decode_wal(line: &str) -> Option { value: Vec::new(), batch: Vec::new(), }), - "txn" => { - let raw = field_str(line, "batch"); - let bytes = from_hex(raw).ok()?; - let batch = unpack_txn(&bytes).ok()?; - Some(WalEntry { - seq, - op: WalOp::Txn, - coord, - value: Vec::new(), - batch, - }) - } _ => None, } }