fix(cubesys): make transaction commits durable + replayable
Two correctness bugs found via ad-hoc daemon verification (T5 was
compile-verified only before):
1. decode_wal dropped WalOp::Txn entries on replay: the txn encoder emits
{"seq","op":"txn","batch"} with NO c/z/y/x fields, but decode_wal
read c/z/y/x unconditionally -> field_u8("c") returned Err -> the
whole entry was skipped. Committed transactions silently vanished on
restart. Fix: branch on op=='txn' before the c/z/y/x extraction.
2. commit was not synchronously durable: append_txn only buffered to the
WAL pending buffer; fsync happened on the 25ms group thread. A
clean stop within that window lost the commit. Fix: commit_txn now
calls wal.flush_pending() (fsync) before returning, so COMMIT is
durable on return -- a real transaction boundary.
Adds unit test commit_replays_from_wal_without_checkpoint (would have
failed before fix 1). Ad-hoc verifier exercises all 3 changed paths on
the live cube-server socket.
This commit is contained in:
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
+20
-12
@@ -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<WalEntry> {
|
||||
}
|
||||
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<WalEntry> {
|
||||
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,
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user