cubesys: durability fixes A/B/C verified live on VM
- Fix A (store.rs): checkpoint boundary persists wal.committed_seq (highest fsync'd) instead of wal.seq() (next-to-assign), which skipped all WAL entries since last checkpoint -> silent data loss on reboot. - Fix B (commands.rs open): run decrypted program in isolated read_snapshot() clone instead of put_raw plaintext over sealed envelope (stopped reboot-time EnvelopeTooShort / clobber). - Fix C (commands.rs keyinit): flush OS key cells to WAL via log_put so they fold into base snapshot and survive reboot (keyinit #2 issues 0, not 2); previously re-minted random material each boot -> sealed records unopenable. - Regression guards durable_sealed_record_survives_restart + open_does_not_clobber_sealed_record in cubesys/src/commands.rs. - STARTUP-README: replace stale 'EPHEMERAL across restarts' caveat with the fixed/verified durability note. Verified live: systemctl restart cube-server (VM reboot path) -> sealed record decrypts+executes after reboot; key cell byte-identical; keyinit idempotent.
This commit is contained in:
@@ -20,8 +20,9 @@
|
||||
//! run <path> # load the code cell at <path> and run the VM
|
||||
//! ls <dir> # list a cubefs directory
|
||||
//! stat <path> # getattr via cubefs
|
||||
//! seal <path> <K.C.Z.Y.X> <tf> # encrypt the record at <path> (tf: none|gcm|chacha|xts)
|
||||
//! open <path> <K.C.Z.Y.X> <tf> # decrypt + decode + run the sealed record
|
||||
//! seal <path> <K.Z.Y.X|auto> <tf> # encrypt the record at <path> (tf: none|gcm|chacha|xts)
|
||||
//! open <path> <K.Z.Y.X|auto> <tf> # decrypt + decode + run the sealed record
|
||||
//! keyinit # ensure the OS Null-space keystore exists
|
||||
//!
|
||||
//! Coordinates are written `C.Z.Y.X` (decimal). Key cells live in Null space,
|
||||
//! so they are given directly as coordinates, not as cubefs paths.
|
||||
@@ -33,7 +34,7 @@ use cubesys::commands::Session;
|
||||
/// Command words understood by the shared interpreter. When `cube`'s first
|
||||
/// argument is one of these, it is run as a single command against a fresh
|
||||
/// in-memory store (the same path as `cube repl`), so the OS / a script can
|
||||
/// invoke e.g. `cube open /c001/z001/y001/x001 001.001.001.001 none` directly
|
||||
/// invoke e.g. `cube open /c001/z001/y001/x001 auto none` directly
|
||||
/// — this is the Phase-3 "open by CZYX + flags" surface made a first-class
|
||||
/// CLI command rather than REPL-only.
|
||||
fn is_command_word(w: &str) -> bool {
|
||||
@@ -46,6 +47,7 @@ fn is_command_word(w: &str) -> bool {
|
||||
| "stat"
|
||||
| "seal"
|
||||
| "open"
|
||||
| "keyinit"
|
||||
| "query"
|
||||
| "begin"
|
||||
| "commit"
|
||||
@@ -132,7 +134,8 @@ fn print_help() {
|
||||
run <path> run the code cell at <path>\n \
|
||||
ls <dir> list a cubefs directory\n \
|
||||
stat <path> getattr via cubefs\n \
|
||||
seal <path> <K.Z.Y.X> <tf> encrypt a record (tf: none|gcm|chacha|xts)\n \
|
||||
open <path> <K.Z.Y.X> <tf> decrypt + decode + run a sealed record\n"
|
||||
seal <path> <K.Z.Y.X|auto> <tf> encrypt a record (tf: none|gcm|chacha|xts)\n \
|
||||
open <path> <K.Z.Y.X|auto> <tf> decrypt + decode + run a sealed record\n \
|
||||
keyinit ensure the OS Null-space keystore exists\n"
|
||||
);
|
||||
}
|
||||
|
||||
+207
-8
@@ -787,6 +787,49 @@ impl Session {
|
||||
a.ino, a.kind, a.size, a.mode
|
||||
))
|
||||
}
|
||||
"keyinit" => {
|
||||
// Phase 3 key-management flow: ensure the OS keystore exists in
|
||||
// the (durable) store. Idempotent — only issues key cells when
|
||||
// absent, so previously-sealed records stay openable across
|
||||
// reboots. OS services can then `open`/`seal` by CZYX + flags
|
||||
// unattended, pointing at the Null-space key cells we mint here.
|
||||
if self.txn.is_some() {
|
||||
return Err(
|
||||
"keyinit inside a transaction is not supported; commit or rollback first"
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
let (ks, issued) = self.store.with_mut(cubecrypt::ensure_os_keystore);
|
||||
// CRITICAL: the key cells live in the Null-space keystore and
|
||||
// must survive a daemon restart, otherwise every previously
|
||||
// sealed record becomes unopenable (AuthFailed) after reboot
|
||||
// because `keyinit` re-mints random material on each boot.
|
||||
// `ensure_os_keystore` writes them into the *live* store only;
|
||||
// if we don't also log them to the WAL, a delta-append
|
||||
// checkpoint never folds them into the base snapshot and they
|
||||
// are lost on the next reload. Re-log every key cell after
|
||||
// provisioning so it is durable (idempotent: log_put of an
|
||||
// unchanged value is harmless).
|
||||
for cell in [ks.default_key.cell, ks.xts_key.cell] {
|
||||
if let Some(v) = self.store.get_raw(&cell) {
|
||||
self.store.log_put(cell, v);
|
||||
}
|
||||
}
|
||||
Ok(format!(
|
||||
"ok: OS keystore ready (issued {issued} new key cell(s)); \
|
||||
default key @ {default} ({dt}), xts key @ {xts} ({xt})",
|
||||
issued = issued,
|
||||
default = ks.default_key.cell.pack_u32(),
|
||||
dt = match ks.default_key.transform {
|
||||
cubecrypt::TransformId::Aes256Gcm => "gcm",
|
||||
cubecrypt::TransformId::ChaCha20Poly1305 => "chacha",
|
||||
cubecrypt::TransformId::Aes256Xts => "xts",
|
||||
cubecrypt::TransformId::None => "none",
|
||||
},
|
||||
xts = ks.xts_key.cell.pack_u32(),
|
||||
xt = "xts",
|
||||
))
|
||||
}
|
||||
"seal" | "open" => {
|
||||
if self.txn.is_some() {
|
||||
// R6: denied — log the attempted (unsupported) op.
|
||||
@@ -805,8 +848,23 @@ impl Session {
|
||||
.next()
|
||||
.ok_or_else(|| format!("{cmd} needs <transform>"))?;
|
||||
let coord = crate::path_to_czyx(path).map_err(|e| e.to_string())?;
|
||||
let kc = parse_coord(keyc)
|
||||
.ok_or_else(|| "bad key-cell coord (use C.Z.Y.X)".to_string())?;
|
||||
|
||||
// Key-cell resolution (Phase 3 key-management flow):
|
||||
// * `auto` (or the OS canonical Null coord `0.20.0.1`) uses the
|
||||
// boot-issued OS keystore — provisioned idempotently if
|
||||
// missing — so an OS service can open/seal by CZYX + flags
|
||||
// UNATTENDED without a human passing a key coordinate.
|
||||
// * any other `K.Z.Y.X` is used verbatim (operator-supplied key).
|
||||
let keyc_raw = keyc.trim();
|
||||
let kc = if keyc_raw == "auto"
|
||||
|| keyc_raw == cubecrypt::OS_KEY_DEFAULT.pack_u32().to_string()
|
||||
{
|
||||
let (ks, _issued) = self.store.with_mut(cubecrypt::ensure_os_keystore);
|
||||
ks.default_key.cell
|
||||
} else {
|
||||
parse_coord(keyc_raw)
|
||||
.ok_or_else(|| "bad key-cell coord (use C.Z.Y.X or 'auto')".to_string())?
|
||||
};
|
||||
let transform = parse_transform(tf)
|
||||
.ok_or_else(|| "unknown transform (none|gcm|chacha|xts)".to_string())?;
|
||||
|
||||
@@ -821,9 +879,12 @@ impl Session {
|
||||
// Audit the destructive op that is about to run (seal vs open).
|
||||
self.audit_now(if cmd == "seal" { OP_SEAL } else { OP_OPEN }, coord, true);
|
||||
|
||||
if store.get_record(&kc).is_none() {
|
||||
store.put_raw(kc, b"demo-key-material-32-bytes-long!!".to_vec());
|
||||
}
|
||||
// The key cell at `kc` must already hold key material:
|
||||
// * `auto` -> `keyinit` just provisioned it in Null space.
|
||||
// * explicit -> the operator supplied a real key coordinate.
|
||||
// We no longer fall back to the old hard-coded demo string; a
|
||||
// missing explicit key cell surfaces as a clear KeyCellMissing
|
||||
// error rather than silently sealing under a constant key.
|
||||
let env = CubeEnv::new(
|
||||
vec![KeySlot {
|
||||
key_cell: kc,
|
||||
@@ -857,9 +918,14 @@ impl Session {
|
||||
.map_err(|e| format!("open: {e:?}"))?;
|
||||
let cell = CodeCell::from_record(coord, &CubeHeader::new(), &pt)
|
||||
.ok_or_else(|| "open: decrypted body is not valid bytecode".to_string())?;
|
||||
store.put_raw(coord, pt);
|
||||
let sn = store.read_snapshot();
|
||||
let mut vm = Vm::new(sn);
|
||||
// Execute the decrypted program WITHOUT writing it back over
|
||||
// the sealed record: substitute the plaintext into an
|
||||
// isolated clone of the snapshot so the persisted (encrypted)
|
||||
// record at `coord` stays intact and can be re-opened later
|
||||
// (e.g. after a reboot) without being clobbered by plaintext.
|
||||
let mut snap = store.read_snapshot();
|
||||
snap.put_raw(coord, pt);
|
||||
let mut vm = Vm::new(snap);
|
||||
let res = vm.run(cell.label);
|
||||
Ok(format!(
|
||||
"open+run {path} (key {}) => {res:?}",
|
||||
@@ -1604,4 +1670,137 @@ mod tests {
|
||||
assert!(r.unwrap_err().contains("require a HELLO identity"));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
}
|
||||
|
||||
// REGRESSION (2026-08-13, Fix A): a sealed record written AFTER the last
|
||||
// checkpoint must survive a daemon restart. The old checkpoint boundary
|
||||
// persisted `wal.seq()` (next-to-assign, one past the last entry) as the
|
||||
// `.seq` boundary, so `replay_after` skipped every post-checkpoint WAL
|
||||
// entry on reopen — silent data loss. The fix persists `committed_seq`.
|
||||
#[test]
|
||||
fn durable_sealed_record_survives_restart() {
|
||||
let dir = std::env::temp_dir().join(format!("cube2-seal-restart-{}", std::process::id()));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
let _ = std::fs::create_dir_all(&dir);
|
||||
let (db, wal, rec) = (
|
||||
dir.join("db.cubedb"),
|
||||
dir.join("wal.ndjson"),
|
||||
dir.join("recovery.jsonl"),
|
||||
);
|
||||
let cfg = DurabilityConfig::default();
|
||||
|
||||
// Session 1: provision keystore, seal x090, checkpoint (folds into base),
|
||||
// then seal x091 *after* the checkpoint (WAL-only), then drop (= flush).
|
||||
{
|
||||
let cs = Arc::new(
|
||||
ConcurrentStore::open(
|
||||
db.to_str().unwrap(),
|
||||
wal.to_str().unwrap(),
|
||||
rec.to_str().unwrap(),
|
||||
cfg,
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
exec_on_store(&cs, "keyinit").unwrap();
|
||||
exec_on_store(&cs, "prog /c200/z099/y001/x090 const 7 const 3 add halt").unwrap();
|
||||
let s = exec_on_store(&cs, "seal /c200/z099/y001/x090 auto gcm").unwrap();
|
||||
assert!(
|
||||
s.contains("sealed /c200/z099/y001/x090 under key 1310721"),
|
||||
"seal must report auto key: {s}"
|
||||
);
|
||||
cs.checkpoint(); // fold x090 into base + record boundary
|
||||
// Post-checkpoint write: only in the WAL, NOT in the base.
|
||||
exec_on_store(&cs, "prog /c200/z099/y001/x091 const 7 const 3 add halt").unwrap();
|
||||
let s2 = exec_on_store(&cs, "seal /c200/z099/y001/x091 auto gcm").unwrap();
|
||||
assert!(
|
||||
s2.contains("sealed /c200/z099/y001/x091 under key 1310721"),
|
||||
"second seal must report auto key: {s2}"
|
||||
);
|
||||
// Wait for the group-commit fsync gate so x091 is durably in the
|
||||
// WAL (NOT folded into the base) before we drop the store. This is
|
||||
// exactly the path Fix A protects: a post-checkpoint WAL entry that
|
||||
// must be replayed on reopen. Without the wait the entry would only
|
||||
// live in the in-memory pending buffer and be lost on drop.
|
||||
std::thread::sleep(std::time::Duration::from_millis(100));
|
||||
// drop(cs) flushes the pending WAL so x091 is durable on disk.
|
||||
}
|
||||
|
||||
// Session 2 (simulated reboot): reopen the same store files.
|
||||
{
|
||||
let cs = Arc::new(
|
||||
ConcurrentStore::open(
|
||||
db.to_str().unwrap(),
|
||||
wal.to_str().unwrap(),
|
||||
rec.to_str().unwrap(),
|
||||
cfg,
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
let k = exec_on_store(&cs, "keyinit").unwrap();
|
||||
assert!(
|
||||
k.contains("issued 0 new key cell"),
|
||||
"keystore must persist across restart: {k}"
|
||||
);
|
||||
// The post-checkpoint sealed record must have been replayed.
|
||||
let o = exec_on_store(&cs, "open /c200/z099/y001/x091 auto gcm").unwrap();
|
||||
assert!(
|
||||
o.contains("open+run /c200/z099/y001/x091 (key 1310721)"),
|
||||
"post-checkpoint sealed record must survive restart (WAL replay): {o}"
|
||||
);
|
||||
assert!(
|
||||
!o.contains("no record at"),
|
||||
"record must not be lost after restart: {o}"
|
||||
);
|
||||
assert!(
|
||||
!o.contains("EnvelopeTooShort"),
|
||||
"reopened record must still be a valid envelope: {o}"
|
||||
);
|
||||
// The pre-checkpoint record (folded into the base) also still opens.
|
||||
let o0 = exec_on_store(&cs, "open /c200/z099/y001/x090 auto gcm").unwrap();
|
||||
assert!(
|
||||
o0.contains("open+run /c200/z099/y001/x090 (key 1310721)"),
|
||||
"pre-checkpoint sealed record must also survive: {o0}"
|
||||
);
|
||||
}
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
}
|
||||
|
||||
// REGRESSION (2026-08-13, Fix B): `open` must not overwrite the sealed
|
||||
// record with the decrypted plaintext, or a second open (or any reopen)
|
||||
// finds plaintext where it expects a cubecrypt envelope (EnvelopeTooShort).
|
||||
#[test]
|
||||
fn open_does_not_clobber_sealed_record() {
|
||||
let dir = std::env::temp_dir().join(format!("cube2-open-clobber-{}", std::process::id()));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
let _ = std::fs::create_dir_all(&dir);
|
||||
let cs = Arc::new(
|
||||
ConcurrentStore::open(
|
||||
dir.join("db.cubedb").to_str().unwrap(),
|
||||
dir.join("wal.ndjson").to_str().unwrap(),
|
||||
dir.join("recovery.jsonl").to_str().unwrap(),
|
||||
DurabilityConfig::default(),
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
exec_on_store(&cs, "keyinit").unwrap();
|
||||
exec_on_store(&cs, "prog /c200/z099/y001/x090 const 7 const 3 add halt").unwrap();
|
||||
exec_on_store(&cs, "seal /c200/z099/y001/x090 auto gcm").unwrap();
|
||||
|
||||
let first = exec_on_store(&cs, "open /c200/z099/y001/x090 auto gcm").unwrap();
|
||||
assert!(
|
||||
first.contains("open+run /c200/z099/y001/x090 (key 1310721)"),
|
||||
"first open must decrypt + run: {first}"
|
||||
);
|
||||
// Second open must still succeed identically — the sealed record was
|
||||
// NOT overwritten by plaintext after the first open.
|
||||
let second = exec_on_store(&cs, "open /c200/z099/y001/x090 auto gcm").unwrap();
|
||||
assert!(
|
||||
second.contains("open+run /c200/z099/y001/x090 (key 1310721)"),
|
||||
"second open must still decrypt the sealed record (no clobber): {second}"
|
||||
);
|
||||
assert!(
|
||||
!second.contains("EnvelopeTooShort"),
|
||||
"second open must not find plaintext where the envelope was: {second}"
|
||||
);
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
}
|
||||
}
|
||||
|
||||
+17
-7
@@ -801,9 +801,16 @@ fn checkpoint_store(
|
||||
));
|
||||
}
|
||||
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));
|
||||
// Boundary must be the highest *durable* (fsync'd) WAL seq,
|
||||
// NOT `wal.seq()` (which is the next-to-assign counter and
|
||||
// sits one past the last entry). Persisting the next-to-
|
||||
// assign value made `replay_after` skip every still-valid
|
||||
// WAL entry on restart — i.e. silent data loss of any
|
||||
// record written since the previous checkpoint.
|
||||
let durable = wal.committed_seq.load(Ordering::SeqCst);
|
||||
persist_seq(cp_seq_path, durable);
|
||||
wal.set_cp_seq_wal(durable);
|
||||
wal.set_base_seq(durable);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -857,10 +864,13 @@ fn fold_delta_into_base(
|
||||
}
|
||||
// Delta is now fully represented by the base; truncate it.
|
||||
let _ = fs::write(delta_path, b"");
|
||||
let seq = wal.seq.load(Ordering::SeqCst);
|
||||
persist_seq(cp_seq_path, seq);
|
||||
wal.set_cp_seq_wal(seq);
|
||||
wal.set_base_seq(seq);
|
||||
// Boundary = highest *durable* WAL seq (fsync'd), not `wal.seq()` (the
|
||||
// next-to-assign counter, which sits one past the last entry and would
|
||||
// make `replay_after` skip still-valid entries on restart).
|
||||
let durable = wal.committed_seq.load(Ordering::SeqCst);
|
||||
persist_seq(cp_seq_path, durable);
|
||||
wal.set_cp_seq_wal(durable);
|
||||
wal.set_base_seq(durable);
|
||||
}
|
||||
|
||||
/// Read `db_path` (full base) then apply the delta file; returns the
|
||||
|
||||
Reference in New Issue
Block a user