Files
cubelinux-2/cubesys/src/commands.rs
T
CUBELinux-2 1539ecd1bb feat(os): OS operator kernels in CUBE + thin native effector (Steps 1 & 2)
- cubecode/src/cb.rs: Behavior descriptor bitflag (PURE/IO_HEAVY/HOT_PATH)
  encoded into HeaderFlags bits 13..15; C_OS_KERNEL=210 / C_OS_EFFECT=211
  coordinate bands; round-trip unit test.
- fix: HEADER_FLAG_BEHAVIOR mask was 0x7000 (bits 12-14) which grabbed the
  ENCRYPTED bit (12) and dropped HOT_PATH (bit 15). Corrected to 0xE000.
- store_code_cell/put_code_cell/header_for_code gain optional descriptor arg;
  all 12 prior call sites pass None (total change).
- CLI: prog K=<flag> tokens, kernel verb (links+descriptors), tick verb lays
  down the OS kernel call-graph (cfg->decide->summarize->tick), native-apply
  verb = thin effector (runs kernel in VM, reads computed result, emits effect).
- integration test os_kernels_live_in_cube_with_call_graph_and_descriptors.

Proven green via ./check quick (fmt+clippy -D warnings+tests).
2026-08-13 15:57:10 -04:00

2181 lines
97 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Shared CUBELinux-2 system command interpreter.
//!
//! This module holds the *single* implementation of the cube command language
//! (`prog`, `write`, `run`, `ls`, `stat`, `seal`, `open`, `query`). It is used
//! by every front-end — the local `cube` REPL, the `cubec` socket client, and
//! the `cube-server` daemon — so the behaviour can never drift between them.
//!
//! A [`Session`] wraps one [`ConcurrentStore`] (a mutex-wrapped `CubeStore`
//! with a write-ahead log + scheduled checkpoint behind it). Because the store
//! is concurrent and durable, the daemon can serve many connections at once and
//! survives restarts. Each `exec` takes and returns an `Arc<ConcurrentStore>`
//! so the server can hand a cloned handle to each worker thread.
use crate::audit::{Audit, OP_DELETE, OP_GRANT, OP_OPEN, OP_READ, OP_REVOKE, OP_SEAL, OP_WRITE};
use crate::grants::{grant, grant_allows, perms_from_str, revoke, Owner, Perm};
use crate::store::ConcurrentStore;
use crate::tenant::{TenantIdentity, TenantSession};
use cubecode::{CodeCell, Kind, Op, RunResult, Vm};
use cubecoords::{CubeHeader, Czyx};
use cubecrypt::{CubeEnv, KeySlot, Selector, TransformId};
use cubestore::{CubeStore, HashBackend};
use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Instant;
/// A buffered write or delete awaiting `COMMIT`. `put: Some((bytes, header))`
/// is a put of a full record (the exact backend bytes `put_record` would
/// write); `put: None` is a delete.
struct TxnOp {
coord: Czyx,
put: Option<(Vec<u8>, CubeHeader)>,
}
/// An in-flight transaction (Task 5). `snapshot` is a consistent clone taken at
/// `BEGIN` so reads during the txn are isolated from concurrent external
/// writers; `ops` is the buffered set of writes/deletes applied atomically on
/// `COMMIT`.
struct Txn {
snapshot: CubeStore<HashBackend>,
ops: Vec<TxnOp>,
}
/// Per-command latency accumulator (cumulative; the daemon reports these via
/// the `stats` command). Counts and sums are exact; mean/max are derived.
#[derive(Default, Clone)]
struct CmdStat {
count: u64,
total_ns: u128,
max_ns: u128,
}
/// One cube command session: a concurrent, durable store plus the interpreter.
/// The REPL holds a transient one; the daemon shares a single `Arc<Session>`
/// across all connection threads.
pub struct Session {
store: Arc<ConcurrentStore>,
/// Total commands executed under this session (telemetry).
calls: u64,
/// Per top-level command latency histogram (command name -> stats).
per_cmd: BTreeMap<String, CmdStat>,
/// The client's declared identity (Task 3): tenant + owner. Stamped by the
/// daemon after a `HELLO`, then threaded into mutating ops for owner
/// enforcement (Tasks 6+).
identity: Option<TenantIdentity>,
/// In-flight transaction (Task 5). `Some` between `BEGIN` and the matching
/// `COMMIT`/`ROLLBACK`; mutating commands buffer into it instead of
/// touching the store live, and reads consult the `BEGIN` snapshot.
txn: Option<Txn>,
/// When true, owner enforcement (Task 6) requires a stamped identity: a
/// mutating op from a session with no `HELLO` identity is rejected, and a
/// HELLO'd owner may only write records it owns. The daemon sets this per
/// connection; the library/REPL/tests leave it false (legacy permissive
/// behaviour, so pre-existing `cubec`/stress.sh flows keep working until
/// the operator opts in via the daemon's `--require-identity` policy).
enforce_owner: bool,
/// Audit trail (plan R6). `Some` when the daemon enabled it on this
/// connection; every op is then appended (permitted or denied) so an
/// operator can later review who touched what. `None` keeps the library
/// REPL/tests silent (no audit record churn) unless explicitly enabled.
audit: Option<Audit>,
}
impl Default for Session {
fn default() -> Self {
Self::new()
}
}
impl Session {
/// A fresh, empty session over an in-memory (non-durable) store.
pub fn new() -> Self {
Session {
store: Arc::new(ConcurrentStore::memory()),
calls: 0,
per_cmd: BTreeMap::new(),
identity: None,
txn: None,
enforce_owner: false,
audit: None,
}
}
/// A session over a durable, concurrent store (used by the daemon).
pub fn with_store(store: Arc<ConcurrentStore>) -> Self {
Session {
store,
calls: 0,
per_cmd: BTreeMap::new(),
identity: None,
txn: None,
enforce_owner: false,
audit: None,
}
}
/// Build a session bound to a tenant's store + identity (used by the
/// daemon per connection, so a `BEGIN`/`COMMIT` spanning multiple command
/// frames keeps its txn state on the same `Session`). `enforce_owner`
/// starts false; the daemon sets it true on every connection that must
/// require a `HELLO` identity (the operator's `--require-identity` policy).
pub fn for_tenant(ts: &TenantSession) -> Self {
Session {
store: ts.store.clone(),
calls: 0,
per_cmd: BTreeMap::new(),
identity: ts.identity(),
txn: None,
enforce_owner: false,
audit: None,
}
}
/// Enable the audit trail for this session (plan R6). Called by the daemon
/// once per connection after construction.
pub fn enable_audit(&mut self) {
self.audit = Some(Audit::new(self.store.clone()));
}
/// Return the full audit log for this session's store (plan R6), or `None`
/// when auditing is disabled. Mirrors the `audit` command without requiring
/// a round-trip through [`exec`]; used by tests and by the daemon's
/// introspection path.
pub fn audit_dump(&self) -> Option<String> {
self.audit.as_ref().map(|a| a.dump())
}
/// Append an audit entry if auditing is enabled. No-op otherwise.
fn audit_now(&self, op: u8, coord: Czyx, ok: bool) {
if let Some(a) = &self.audit {
let owner = self
.identity
.as_ref()
.map(|i| i.owner_local.as_str())
.unwrap_or("<anonymous>");
a.append(op, coord, owner, ok);
}
}
/// Stamp the client identity (called by the daemon after a successful
/// `HELLO`, or via [`Session::for_tenant`] from the tenant session).
pub fn set_identity(&mut self, id: TenantIdentity) {
self.identity = Some(id);
}
/// The currently-stamped identity, if any.
pub fn identity(&self) -> Option<TenantIdentity> {
self.identity.clone()
}
/// Enable/disable owner enforcement for this session (Task 6 B). The daemon
/// calls this per connection to apply its `--require-identity` policy.
pub fn set_enforce_owner(&mut self, on: bool) {
self.enforce_owner = on;
}
/// Authorization gate for a mutating op (Task 6 + Task 6b — the PDF's
/// flags 5-19 delegated-grant layer). A mutating command (prog/write/del/
/// seal/open) may proceed only when EITHER the session's identity owner
/// matches the record's `owner_local_user` (owner), OR the caller holds a
/// grant authorizing `want` over the coordinate. Returns `None` to allow,
/// or an error message to reject.
///
/// Enforcement order (matches the plan's D5):
/// 1. owner match -> allow
/// 2. else a matching grant -> allow
/// 3. else deny
///
/// An unowned record is claimable by any (identified) writer (first-write
/// wins), same as the Task 6 rule.
fn admit_mutate(&self, coord: Czyx, want: Perm) -> Option<String> {
// Policy (matches the prior Task 6 `owner_violation` contract):
// * No session identity:
// - if `enforce_owner` (daemon `--require-identity`) is ON -> reject
// (anonymous writes are not allowed);
// - otherwise (legacy / library / REPL / test) -> allow
// * A session identity IS present: owner-match OR a grant always apply
// (this is the heart of the PDF's flags 5-19 delegated model).
match &self.identity {
None => {
if self.enforce_owner {
Some(
"owner enforcement: mutating operations require a HELLO identity \
(daemon policy --require-identity)"
.to_string(),
)
} else {
None
}
}
Some(id) => {
// (1) owner match — or unowned record, claimable by first writer.
{
let record_owner = self.store.owner(&coord)?;
if record_owner == id.owner_local {
return None;
}
}
// (2) grant.
let who = Owner::from(id);
if grant_allows(&self.store, &who, &coord, want) {
return None;
}
// (3) deny.
Some(format!(
"owner violation: record at {} owned by '{}', you are '{}' (and hold no matching grant)",
coord.pack_u32(),
self.store.owner(&coord).unwrap_or_default(),
id.owner_local
))
}
}
}
/// Authorization gate for a read op (R5 — closes the read-privacy gap in
/// the daemon path). A read command (`run`/`stat`, and the metadata read in
/// `ls`) may proceed only when EITHER the session's identity owner matches
/// the record's `owner_local_user` (owner), OR the caller holds a grant
/// authorizing `Perm::Read` over the coordinate. Returns `None` to allow, or
/// an error message to reject.
///
/// Mirrors [`admit_mutate`] (Task 6 + 6b): same anonymous policy, same
/// owner→grant→deny ordering. An unowned record is world-readable (first
/// writer claims ownership on write; reads never create state).
///
/// Note: the FUSE mount enforces POSIX `mode` bits via
/// `cubefs::nullspace::Acl`; the daemon deliberately does NOT import that
/// uid/gid model — its identity is name-based (`Owner`) from `HELLO`, so
/// read gating here is owner/grant-native rather than a literal `Acl` lift.
/// (See plan R1/R5 for the rationale.)
fn admit_read(&self, coord: Czyx) -> Option<String> {
match &self.identity {
None => {
if self.enforce_owner {
Some(
"owner enforcement: read operations require a HELLO identity \
(daemon policy --require-identity)"
.to_string(),
)
} else {
None
}
}
Some(id) => {
// (1) owner match — or unowned record, world-readable.
{
let record_owner = self.store.owner(&coord)?;
if record_owner == id.owner_local {
return None;
}
}
// (2) read grant.
let who = Owner::from(id);
if grant_allows(&self.store, &who, &coord, Perm::Read) {
return None;
}
// (3) deny.
Some(format!(
"owner violation: record at {} owned by '{}', you are '{}' (and hold no matching read grant)",
coord.pack_u32(),
self.store.owner(&coord).unwrap_or_default(),
id.owner_local
))
}
}
}
/// Shared handle to the underlying concurrent store.
pub fn store(&self) -> Arc<ConcurrentStore> {
self.store.clone()
}
/// Snapshot of the session's telemetry: total commands serviced, per-command
/// latency distribution (mean/max in µs), and the occupancy of each `C`
/// namespace (number of records whose class axis equals `c`).
pub fn stats(&self) -> String {
let mut lines = Vec::new();
lines.push(format!("total commands serviced: {}", self.calls));
lines.push("per-command latency (µs, mean / max / count):".into());
if self.per_cmd.is_empty() {
lines.push(" (no commands timed yet)".into());
} else {
for (name, st) in &self.per_cmd {
let mean_us = if st.count > 0 {
(st.total_ns as f64) / (st.count as f64) / 1e3
} else {
0.0
};
let max_us = st.max_ns as f64 / 1e3;
lines.push(format!(
" {:<8} mean {:8.2} max {:8.2} n={}",
name, mean_us, max_us, st.count
));
}
}
let keys = self.store.keys();
let mut by_c: BTreeMap<u8, usize> = BTreeMap::new();
for k in &keys {
*by_c.entry(k.c).or_insert(0) += 1;
}
lines.push(format!("records by C namespace ({} total):", keys.len()));
if by_c.is_empty() {
lines.push(" (store empty)".into());
} else {
for (c, n) in &by_c {
let label = if *c == 0 { "Null(0)" } else { "" };
lines.push(format!(" C={c:<3} {n:>6} records {label}"));
}
}
lines.join("\n")
}
/// Execute one command line. `Ok(out)` is a (possibly multi-line) result to
/// print; `Err(e)` is a human-readable error. Also records per-command
/// latency into the session telemetry (see [`Session::stats`]).
pub fn exec(&mut self, line: &str) -> Result<String, String> {
let t0 = Instant::now();
let cmd_name = line.split_whitespace().next().unwrap_or("").to_string();
let result = self.exec_inner(line);
self.calls += 1;
let st = self.per_cmd.entry(cmd_name).or_default();
let elapsed = t0.elapsed().as_nanos();
st.count += 1;
st.total_ns += elapsed;
if elapsed > st.max_ns {
st.max_ns = elapsed;
}
result
}
/// The real interpreter (separated so [`exec`] can wrap it with timing).
fn exec_inner(&mut self, line: &str) -> Result<String, String> {
let mut it = line.split_whitespace();
let cmd = it.next().ok_or_else(|| "empty line".to_string())?;
let store = &self.store;
match cmd {
"begin" => {
if self.txn.is_some() {
return Err("begin: already in a transaction".to_string());
}
// Take a consistent snapshot now; reads during the txn consult
// it (isolation from concurrent external writers).
let snap = store.read_snapshot();
self.txn = Some(Txn {
snapshot: snap,
ops: Vec::new(),
});
Ok("ok: transaction begun".to_string())
}
"commit" => {
let txn = self
.txn
.take()
.ok_or_else(|| "commit: no transaction is open".to_string())?;
let n = txn.ops.len();
// Task 6/6b: re-check ownership + grants for buffered txn ops.
// The live path already checked at `prog`/`write`/`del` time,
// but a concurrent commit from another owner (or a revoked grant)
// could have changed the situation between our BEGIN and COMMIT,
// so re-check now. The grant-aware `admit_mutate` is used so a
// grant issued/revoked mid-txn is honored at commit.
for op in &txn.ops {
if let Some(msg) = self.admit_mutate(op.coord, Perm::Write) {
// Roll the txn back: restore it so the caller can retry
// after resolving the conflict (do not silently drop).
self.txn = Some(txn);
return Err(msg);
}
}
// Apply every buffered op atomically under one store write lock
// and append a SINGLE WAL entry (WalOp::Txn) so the whole batch
// is durable as one unit and replays idempotently.
let batch: Vec<crate::store::TxnEntry> = txn
.ops
.iter()
.map(|op| crate::store::TxnEntry {
coord: op.coord,
value: op.put.as_ref().map(|(v, _)| v.clone()),
})
.collect();
store.commit_txn(&batch);
Ok(format!("ok: committed {n} operation(s)"))
}
"rollback" => {
if self.txn.take().is_none() {
return Err("rollback: no transaction is open".to_string());
}
Ok("ok: transaction rolled back".to_string())
}
"stats" => Ok(self.stats()),
"audit" => {
// Plan R6: dump the append-only audit trail for this session's
// store. Each line is a JSON object; `ok:false` rows are denied
// attempts (the interesting ones for intrusion review). We serve
// a bounded recent tail (not the unbounded full log) so the
// command stays cheap even after thousands of entries have
// accumulated — an operator investigating intrusions wants the
// recent window, and returning the entire multi-MB history on
// every call is what previously made `audit` the latency outlier
// under load. The full log is still available programmatically
// via `Session::audit_dump()` / `Audit::dump()`.
let log = match &self.audit {
Some(a) => a.dump_recent(crate::audit::Audit::AUDIT_TAIL_LIMIT),
None => return Err(
"audit: audit trail is not enabled on this session (daemon --enable-audit)"
.to_string(),
),
};
if log.trim().is_empty() {
Ok("audit -> (no entries)".to_string())
} else {
Ok(format!("audit ->\n{log}"))
}
}
"query" => {
let dt = it
.next()
.ok_or_else(|| "query needs <doc_type>".to_string())?;
let coords = store.query_doc_type(dt);
if coords.is_empty() {
Ok(format!("query {dt} -> (no matches)"))
} else {
let names: Vec<String> =
coords.iter().map(|c| c.pack_u32().to_string()).collect();
Ok(format!(
"query {dt} -> {} matches: {}",
coords.len(),
names.join(" ")
))
}
}
"prog" => {
let path = it.next().ok_or_else(|| "prog needs <path>".to_string())?;
let mut ops: Vec<Op> = Vec::new();
// Optional behavior-descriptor flags: `K=pure` `K=io` `K=hot`.
// These stamp the OS-kernel behavior bits (PDF §524–525) on the
// record header and switch the kind to `kernel` so the cube
// correctly classifies operator kernels vs plain functions.
let mut descriptor: Option<cubecode::Behavior> = None;
while let Some(tok) = it.next() {
if let Some(flag) = tok.strip_prefix("K=") {
let bit = match flag {
"pure" => cubecode::Behavior::PURE,
"io" => cubecode::Behavior::IO_HEAVY,
"hot" => cubecode::Behavior::HOT_PATH,
other => return Err(format!("prog: unknown descriptor K={other}")),
};
let cur = descriptor.unwrap_or_default();
descriptor = Some(cubecode::Behavior(cur.0 | bit));
continue;
}
let arg = if takes_arg(tok) {
it.next()
.and_then(|a| a.parse::<u8>().ok())
.ok_or_else(|| format!("prog: {tok} needs a u8 argument"))?
} else {
0
};
ops.push(make_op(tok, arg)?);
}
if ops.is_empty() {
return Err("prog: no ops given".to_string());
}
let is_kernel = descriptor.is_some();
let kind = if is_kernel { Kind::Kernel } else { Kind::Fn };
let name = path.rsplit('/').next().unwrap_or(path);
// Compute the exact record bytes `put_record` would write, using
// a throwaway store so we can buffer (or apply) them without
// duplicating the record codec.
let coord = scratch_code_coord(path, kind, name, &ops)?;
// Task 6: reject overwriting a record owned by a different owner.
if let Some(msg) = self.admit_mutate(coord, Perm::Write) {
self.audit_now(OP_WRITE, coord, false);
return Err(msg);
}
self.audit_now(OP_WRITE, coord, true);
let owner = self.identity.as_ref().map(|i| i.owner_local.as_str());
let value = {
let mut scratch = CubeStore::new(HashBackend::new());
crate::store_code_cell(&mut scratch, path, kind, name, &[], &ops, owner, descriptor)
.map_err(|e| e.to_string())?;
scratch.get_raw(&coord).unwrap_or_default()
};
let header = header_for_code(kind, name, &ops, owner, descriptor);
if let Some(txn) = self.txn.as_mut() {
txn.ops.push(TxnOp {
coord,
put: Some((value, header)),
});
let kind_tag = if is_kernel { "kernel" } else { "fn" };
return Ok(format!(
"buffered prog {path} ({kind_tag}, {} ops) — commit to apply",
ops.len()
));
}
let coord = store
.put_code_cell(path, kind, name, &[], &ops, owner, descriptor)
.map_err(|e| e.to_string())?;
let kind_tag = if is_kernel { "kernel" } else { "fn" };
Ok(format!(
"wrote program {path} -> coord {} ({kind_tag}, {} ops)",
coord.pack_u32(),
ops.len()
))
}
"link" => {
// Attach a callee coordinate to an existing code cell's call
// graph. `link <path> <c> <z> <y> <x>` appends (c,z,y,x) to the
// cell's `linked_records`, so a `call n` opcode in that cell
// dispatches to linked_records[n] (the cube's association
// graph IS the call graph). This is what makes functions stored
// in CUBE callable from other functions stored in CUBE.
let path = it.next().ok_or_else(|| "link needs <path>".to_string())?;
let c = parse_u8(it.next(), "link needs <c>")?;
let z = parse_u8(it.next(), "link needs <z>")?;
let y = parse_u8(it.next(), "link needs <y>")?;
let x = parse_u8(it.next(), "link needs <x>")?;
let target = Czyx::new(c, z, y, x);
let coord = crate::path_to_czyx(path).map_err(|e| e.to_string())?;
if let Some(msg) = self.admit_mutate(coord, Perm::Write) {
self.audit_now(OP_WRITE, coord, false);
return Err(msg);
}
self.audit_now(OP_WRITE, coord, true);
// Load the existing cell, append the link, re-store it.
let cell = crate::load_code_cell(&store.read_snapshot(), path)
.map_err(|e| format!("link: cannot load {path}: {e:?}"))?;
let mut links = cell.links().to_vec();
if links.contains(&target) {
return Ok(format!(
"link {path}: {target:?} already linked (degree {})",
links.len()
));
}
links.push(target);
let code = cubecode::decode(&cell.body())
.map_err(|e| format!("link: bad bytecode in {path}: {e:?}"))?;
let owner = self.identity.as_ref().map(|i| i.owner_local.as_str());
store
.put_code_cell(
path,
cell.kind(),
cell.name().unwrap_or(path),
&links,
&code,
owner,
None,
)
.map_err(|e| e.to_string())?;
Ok(format!(
"linked {path} -> {target:?} (call-graph degree now {})",
links.len()
))
}
// ---- OS-operator-kernel authoring + thin effector (Steps 1 & 2) ----
// Per the reframe, the OS's *computation* lives in CUBE as operator
// kernels (call graphs + behavior descriptors, PDF §524–525). The
// native OS layer is only a thin effector: it reads a kernel's
// computed result and applies the effect. The cubevm never does
// store-IO itself. `kernel` authors a kernel; `tick` lays down the
// OS kernel call-graph; `native-apply` is the effector boundary.
"kernel" => {
// `kernel <path> [K=pure|io|hot ...] <op> <arg> ...`
// Like `prog`, but the cell is always kind=Kernel and accepts
// behavior-descriptor flags so the OS marks operator kernels
// distinctly from plain functions.
let path = it.next().ok_or_else(|| "kernel needs <path>".to_string())?;
let mut ops: Vec<Op> = Vec::new();
let mut descriptor: Option<cubecode::Behavior> = None;
while let Some(tok) = it.next() {
if let Some(flag) = tok.strip_prefix("K=") {
let bit = match flag {
"pure" => cubecode::Behavior::PURE,
"io" => cubecode::Behavior::IO_HEAVY,
"hot" => cubecode::Behavior::HOT_PATH,
other => return Err(format!("kernel: unknown descriptor K={other}")),
};
let cur = descriptor.unwrap_or_default();
descriptor = Some(cubecode::Behavior(cur.0 | bit));
continue;
}
let arg = if takes_arg(tok) {
it.next()
.and_then(|a| a.parse::<u8>().ok())
.ok_or_else(|| format!("kernel: {tok} needs a u8 argument"))?
} else {
0
};
ops.push(make_op(tok, arg)?);
}
if ops.is_empty() {
return Err("kernel: no ops given".to_string());
}
let name = path.rsplit('/').next().unwrap_or(path);
let coord = scratch_code_coord(path, Kind::Kernel, name, &ops)?;
if let Some(msg) = self.admit_mutate(coord, Perm::Write) {
self.audit_now(OP_WRITE, coord, false);
return Err(msg);
}
self.audit_now(OP_WRITE, coord, true);
let owner = self.identity.as_ref().map(|i| i.owner_local.as_str());
let value = {
let mut scratch = CubeStore::new(HashBackend::new());
crate::store_code_cell(
&mut scratch,
path,
Kind::Kernel,
name,
&[],
&ops,
owner,
descriptor,
)
.map_err(|e| e.to_string())?;
scratch.get_raw(&coord).unwrap_or_default()
};
let header = header_for_code(Kind::Kernel, name, &ops, owner, descriptor);
if let Some(txn) = self.txn.as_mut() {
txn.ops.push(TxnOp {
coord,
put: Some((value, header)),
});
return Ok(format!(
"buffered kernel {path} ({} ops) — commit to apply",
ops.len()
));
}
let coord = store
.put_code_cell(path, Kind::Kernel, name, &[], &ops, owner, descriptor)
.map_err(|e| e.to_string())?;
Ok(format!(
"wrote kernel {path} -> coord {} ({} ops, descriptors={:?})",
coord.pack_u32(),
ops.len(),
descriptor.map(|b| b.tags()).unwrap_or_default()
))
}
"tick" => {
// Lay down the OS operator-kernel call graph (Step 1). Each leaf
// is a real CUBE kernel composed via `linked_records` (the call
// graph) and tagged with behavior descriptors. `cube-os-tick`
// calls the three leaves in order. Nothing is executed here —
// the native layer later `run`s each and `native-apply`s the
// effect. This is "everything in CUBE" done the spec's way:
// the OS's behavior lives as kernels + call graph + descriptors.
let c = cubecode::C_OS_KERNEL;
let cfg = format!("/c{c}/z001/y001/x001"); // normalize-config: pure
let decide = format!("/c{c}/z001/y001/x002"); // decide-snapshot: io
let summarize = format!("/c{c}/z001/y001/x003"); // summarize-procs: pure+hot
let tick = format!("/c{c}/z001/y001/x004"); // cube-os-tick: hot (calls 0..2)
let cfg_code = vec![Op::Const(7), Op::Const(3), Op::Add, Op::Halt]; // 7+3=10
let decide_code = vec![Op::Const(1), Op::Ret]; // 1 => snapshot
let summarize_code = vec![Op::Const(20), Op::Halt]; // 20 procs
let tick_code = vec![
Op::Const(0),
Op::CallLink(0), // cfg
Op::CallLink(1), // decide
Op::CallLink(2), // summarize
Op::Halt,
];
// Stage the leaves first (we need their coords to build the
// call-graph edges of the root), then the root. `with_mut`
// hands us exclusive `&mut CubeStore` access so the helper can
// write each record through the normal codec.
let cfg_c = store
.with_mut(|s| {
crate::store_code_cell(
s,
&cfg,
Kind::Kernel,
"normalize-config",
&[],
&cfg_code,
None,
Some(cubecode::Behavior(cubecode::Behavior::PURE)),
)
})
.map_err(|e: crate::SysError| e.to_string())?;
let dec_c = store
.with_mut(|s| {
crate::store_code_cell(
s,
&decide,
Kind::Kernel,
"decide-snapshot",
&[],
&decide_code,
None,
Some(cubecode::Behavior(cubecode::Behavior::IO_HEAVY)),
)
})
.map_err(|e: crate::SysError| e.to_string())?;
let sum_c = store
.with_mut(|s| {
crate::store_code_cell(
s,
&summarize,
Kind::Kernel,
"summarize-procs",
&[],
&summarize_code,
None,
Some(cubecode::Behavior(
cubecode::Behavior::PURE | cubecode::Behavior::HOT_PATH,
)),
)
})
.map_err(|e: crate::SysError| e.to_string())?;
let tick_c = store
.with_mut(|s| {
crate::store_code_cell(
s,
&tick,
Kind::Kernel,
"cube-os-tick",
&[cfg_c, dec_c, sum_c], // call graph: tick -> {cfg,decide,summarize}
&tick_code,
None,
Some(cubecode::Behavior(cubecode::Behavior::HOT_PATH)),
)
})
.map_err(|e: crate::SysError| e.to_string())?;
Ok(format!(
"OS kernel call-graph laid down (kind=kernel, in cube c{c}):\n {} normalize-config [pure] -> coord {}\n {} decide-snapshot [io] -> coord {}\n {} summarize-procs [pure,hot] -> coord {}\n {} cube-os-tick [hot] -> coord {} (links cfg,decide,summarize)\nrun e.g.: run {}\nthen effector: native-apply {}",
cfg, cfg_c.pack_u32(), decide, dec_c.pack_u32(), summarize,
sum_c.pack_u32(), tick, tick_c.pack_u32(), tick_c.pack_u32(),
tick
))
}
"native-apply" => {
// Step 2: the thin effector. The decision/computation already
let path = it
.next()
.ok_or_else(|| "native-apply needs <kernel-path>".to_string())?;
let cell = crate::load_code_cell(&store.read_snapshot(), path)
.map_err(|e| format!("native-apply: cannot load {path}: {e}"))?;
if cell.kind() != Kind::Kernel {
return Err(format!(
"native-apply: {path} is kind={:?}, expected a kernel",
cell.kind()
));
}
// Build a throwaway in-memory store holding just this kernel
// (the VM needs a store to walk), then run it. The computation
// is entirely in CUBE; we only read back the computed result.
let coord = crate::path_to_czyx(path).map_err(|e| e.to_string())?;
let mut vm_store = cubestore::CubeStore::new(cubestore::HashBackend::new());
crate::store_code_cell(
&mut vm_store,
path,
Kind::Kernel,
cell.name().unwrap_or(path),
&cell.links().iter().copied().collect::<Vec<_>>(),
&cell.code,
None,
None,
)
.map_err(|e| format!("native-apply: stage failed: {e}"))?;
let mut vm = Vm::new(vm_store);
let result = match vm.run(coord) {
RunResult::Halted { top } => top,
other => {
return Err(format!(
"native-apply: kernel did not halt cleanly: {other:?}"
))
}
};
let r = result.unwrap_or(0);
let effect = match cell.name() {
Some("decide-snapshot") => {
if r != 0 {
format!(
"OS EFFECT: persist runtime snapshot into CUBE (c{} band); decision kernel returned {} => snapshot NOW",
cubecode::C_OS_EFFECT, r
)
} else {
"OS EFFECT: no snapshot (decision kernel returned 0)".to_string()
}
}
_ => format!(
"OS EFFECT: apply computed result {} from kernel {}@{}",
r,
cell.name().unwrap_or("?"),
path
),
};
Ok(format!(
"native-apply {} (kind=kernel, descriptors={:?}):\n computed result = {}\n {}",
path,
cubecode::Behavior::from_flags(
store
.read_snapshot()
.get_record(&coord)
.map(|(h, _)| h.flags.bits())
.unwrap_or(0)
)
.tags(),
r,
effect
))
}
"write" => {
let path = it.next().ok_or_else(|| "write needs <path>".to_string())?;
let bytes: Vec<u8> = it
.map(parse_byte)
.collect::<Option<_>>()
.ok_or_else(|| "write: every byte must be hex/dec 0..255".to_string())?;
let code = cubecode::decode(&bytes)
.map_err(|e| format!("bytecode decode error: {e:?}"))?;
let name = path.rsplit('/').next().unwrap_or(path);
let coord = scratch_code_coord(path, Kind::Fn, name, &code)?;
if let Some(msg) = self.admit_mutate(coord, Perm::Write) {
self.audit_now(OP_WRITE, coord, false);
return Err(msg);
}
self.audit_now(OP_WRITE, coord, true);
let owner = self.identity.as_ref().map(|i| i.owner_local.as_str());
let value = {
let mut scratch = CubeStore::new(HashBackend::new());
crate::store_code_cell(&mut scratch, path, Kind::Fn, name, &[], &code, owner, None)
.map_err(|e| e.to_string())?;
scratch.get_raw(&coord).unwrap_or_default()
};
let header = header_for_code(Kind::Fn, name, &code, owner, None);
if let Some(txn) = self.txn.as_mut() {
txn.ops.push(TxnOp {
coord,
put: Some((value, header)),
});
return Ok(format!(
"buffered write {path} ({} bytes) — commit to apply",
bytes.len()
));
}
let coord = store
.put_code_cell(path, Kind::Fn, name, &[], &code, owner, None)
.map_err(|e| e.to_string())?;
Ok(format!("wrote {path} -> coord {}", coord.pack_u32()))
}
"del" => {
let path = it.next().ok_or_else(|| "del needs <path>".to_string())?;
let coord = crate::path_to_czyx(path).map_err(|e| e.to_string())?;
// Task 6: reject deleting a record owned by a different owner
// (unless unowned, which any identity may take over).
if let Some(msg) = self.admit_mutate(coord, Perm::Write) {
self.audit_now(OP_DELETE, coord, false);
return Err(msg);
}
self.audit_now(OP_DELETE, coord, true);
if let Some(txn) = self.txn.as_mut() {
txn.ops.push(TxnOp { coord, put: None });
return Ok(format!("buffered del {path} — commit to apply"));
}
store.delete_raw(&coord);
Ok(format!("deleted {path}"))
}
// --- Raw coordinate API (used by cubefs' socket-backed backend,
// and any client that wants to address the store by CZYX
// directly instead of by path). These mutate the SAME durable
// `ConcurrentStore` the daemon serves, so a write through here
// is immediately visible to `ls`/`stat`/FUSE and is folded
// into the WAL + checkpoint like any other write. ---
"rawget" => {
let c = parse_u8(it.next(), "rawget needs <c>")?;
let z = parse_u8(it.next(), "rawget needs <z>")?;
let y = parse_u8(it.next(), "rawget needs <y>")?;
let x = parse_u8(it.next(), "rawget needs <x>")?;
let coord = Czyx::new(c, z, y, x);
match store.get_raw(&coord) {
Some(v) => Ok(format!("ok: {}", hex_encode(&v))),
None => Ok("none".to_string()),
}
}
"rawput" => {
let c = parse_u8(it.next(), "rawput needs <c>")?;
let z = parse_u8(it.next(), "rawput needs <z>")?;
let y = parse_u8(it.next(), "rawput needs <y>")?;
let x = parse_u8(it.next(), "rawput needs <x>")?;
let hex = it
.next()
.ok_or_else(|| "rawput needs <hex-bytes>".to_string())?;
let val = hex_decode(hex).ok_or_else(|| "rawput: value must be hex".to_string())?;
let coord = Czyx::new(c, z, y, x);
store.put_raw(coord, val);
Ok(format!("ok: wrote {}", coord.pack_u32()))
}
"rawdel" => {
let c = parse_u8(it.next(), "rawdel needs <c>")?;
let z = parse_u8(it.next(), "rawdel needs <z>")?;
let y = parse_u8(it.next(), "rawdel needs <y>")?;
let x = parse_u8(it.next(), "rawdel needs <x>")?;
let coord = Czyx::new(c, z, y, x);
store.delete_raw(&coord);
Ok(format!("ok: deleted {}", coord.pack_u32()))
}
"rawkeys" => {
let ks: Vec<String> = store
.keys()
.iter()
.map(|k| k.pack_u32().to_string())
.collect();
Ok(format!("ok: {} keys", ks.len())).map(|s| {
if ks.is_empty() {
s
} else {
format!("{s}\n{}", ks.join(" "))
}
})
}
"rawscan" => {
let c = parse_u8(it.next(), "rawscan needs <c>")?;
let z = it.next().and_then(|t| t.parse::<u8>().ok());
let y = it.next().and_then(|t| t.parse::<u8>().ok());
let ks: Vec<String> = store
.scan_prefix(c, z, y)
.iter()
.map(|k| k.pack_u32().to_string())
.collect();
Ok(format!("ok: {} keys", ks.len())).map(|s| {
if ks.is_empty() {
s
} else {
format!("{s}\n{}", ks.join(" "))
}
})
}
"grant" => {
// Issue a permission grant (Task 6b / PDF flags 5-19). Only an
// identified owner may grant (under --require-identity); in the
// legacy permissive mode the granter is treated as "root".
if self.enforce_owner && self.identity.is_none() {
// R6: denied grant attempt — log it for intrusion review.
self.audit_now(OP_GRANT, Czyx::new(0, 0, 0, 0), false);
return Err(
"grant requires a HELLO identity (daemon policy --require-identity)"
.to_string(),
);
}
let identity = self.identity.clone();
let granter = match &identity {
Some(i) => Owner::from(i),
None => Owner::new("root"),
};
let grantee_tok = it
.next()
.ok_or_else(|| "grant needs <grantee_local[#remote]>".to_string())?;
let (gl, gr) = match grantee_tok.split_once('#') {
Some((l, r)) => (l.to_string(), Some(r.to_string())),
None => (grantee_tok.to_string(), None),
};
let perms_s = it
.next()
.ok_or_else(|| "grant needs <perms: r|w|x>".to_string())?;
let perms = perms_from_str(perms_s)
.ok_or_else(|| format!("grant: bad perms '{perms_s}' (use r/w/x)"))?;
let scope = match it.next() {
None => None,
Some(s) => {
if s.eq_ignore_ascii_case("global") {
None
} else {
let c = parse_coord(s)
.ok_or_else(|| format!("grant: bad scope '{s}' (C.Z.Y.X)"))?;
Some(c)
}
}
};
let grantee = Owner {
local: gl.clone(),
remote: gr,
};
let seq = grant(&self.store, &granter, &grantee, perms, scope)
.map_err(|e| format!("grant: {e}"))?;
self.audit_now(OP_GRANT, scope.unwrap_or(Czyx::new(0, 0, 0, 0)), true);
let scope_s = match scope {
Some(c) => c.pack_u32().to_string(),
None => "global".to_string(),
};
Ok(format!(
"ok: granted seq={seq} {gl}<-{perms_s} over {scope_s}"
))
}
"revoke" => {
// Revoke a grant (only the original granter may). Task 6b.
if self.enforce_owner && self.identity.is_none() {
// R6: denied revoke attempt — log it for intrusion review.
self.audit_now(OP_REVOKE, Czyx::new(0, 0, 0, 0), false);
return Err(
"revoke requires a HELLO identity (daemon policy --require-identity)"
.to_string(),
);
}
let identity = self.identity.clone();
let granter = match &identity {
Some(i) => Owner::from(i),
None => Owner::new("root"),
};
let grantee_tok = it
.next()
.ok_or_else(|| "revoke needs <grantee_local[#remote]>".to_string())?;
let (gl, gr) = match grantee_tok.split_once('#') {
Some((l, r)) => (l.to_string(), Some(r.to_string())),
None => (grantee_tok.to_string(), None),
};
let scope = match it.next() {
None => None,
Some(s) => {
if s.eq_ignore_ascii_case("global") {
None
} else {
let c = parse_coord(s)
.ok_or_else(|| format!("revoke: bad scope '{s}' (C.Z.Y.X)"))?;
Some(c)
}
}
};
let grantee = Owner {
local: gl,
remote: gr,
};
let removed = revoke(&self.store, &granter, &grantee, scope);
self.audit_now(OP_REVOKE, scope.unwrap_or(Czyx::new(0, 0, 0, 0)), true);
Ok(format!("ok: revoked {removed} grant(s)"))
}
"run" => {
let path = it.next().ok_or_else(|| "run needs <path>".to_string())?;
let _coord = crate::path_to_czyx(path).map_err(|e| e.to_string())?;
// R5: reads (this one also *executes*) honor the owner/grantee
// read gate, same as writes honor the mutate gate.
if let Some(msg) = self.admit_read(_coord) {
self.audit_now(OP_READ, _coord, false);
return Err(msg);
}
self.audit_now(OP_READ, _coord, true);
let sn = txn_snapshot(self);
let cell = crate::load_code_cell(&sn, path).map_err(|e| e.to_string())?;
let mut vm = Vm::new(sn);
let res = vm.run(cell.label);
let mut out = format!("run {path} => {res:?}");
if !vm.output().is_empty() {
out.push_str(&format!(
"\n trace: {}",
String::from_utf8_lossy(vm.output()).trim_end()
));
}
Ok(out)
}
"ls" => {
let dir = it.next().ok_or_else(|| "ls needs <dir>".to_string())?;
// R5: directory reads are metadata reads — honor the read gate
// so an attacker can't enumerate a victim's records by name.
let dir_coord = crate::path_to_czyx(dir).map_err(|e| e.to_string())?;
if let Some(msg) = self.admit_read(dir_coord) {
self.audit_now(OP_READ, dir_coord, false);
return Err(msg);
}
self.audit_now(OP_READ, dir_coord, true);
let sn = txn_snapshot(self);
let fs = cubefs::CubeFs::new(sn);
let entries = fs.readdir(dir).map_err(|e| format!("ls {dir}: {e:?}"))?;
if entries.is_empty() {
Ok(format!("ls {dir} -> (empty)"))
} else {
let names: Vec<String> = entries.into_iter().map(|(n, _)| n).collect();
Ok(format!("ls {dir} -> {}", names.join(" ")))
}
}
"stat" => {
let path = it.next().ok_or_else(|| "stat needs <path>".to_string())?;
// R5: metadata reads honor the read gate.
let coord = crate::path_to_czyx(path).map_err(|e| e.to_string())?;
if let Some(msg) = self.admit_read(coord) {
self.audit_now(OP_READ, coord, false);
return Err(msg);
}
self.audit_now(OP_READ, coord, true);
let sn = txn_snapshot(self);
let fs = cubefs::CubeFs::new(sn);
let a = fs
.getattr(path)
.map_err(|e| format!("stat {path}: {e:?}"))?;
Ok(format!(
"stat {path} -> ino={} kind={:?} size={} mode={:o}",
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.
let _c = crate::path_to_czyx(it.next().unwrap_or("0.0.0.0"))
.unwrap_or(Czyx::new(0, 0, 0, 0));
self.audit_now(if cmd == "seal" { OP_SEAL } else { OP_OPEN }, _c, false);
return Err(format!(
"{cmd} inside a transaction is not supported; commit or rollback first"
));
}
let path = it.next().ok_or_else(|| format!("{cmd} needs <path>"))?;
let keyc = it
.next()
.ok_or_else(|| format!("{cmd} needs <K.Z.Y.X> key cell"))?;
let tf = it
.next()
.ok_or_else(|| format!("{cmd} needs <transform>"))?;
let coord = crate::path_to_czyx(path).map_err(|e| e.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())?;
// Task 6 (B): seal/open are destructive writes to `coord`, so
// they obey the same owner gate as prog/write/del. A no-identity
// session (under --require-identity) or a non-owner is rejected.
if let Some(msg) = self.admit_mutate(coord, Perm::Write) {
// R6: denied — log the attempted destructive op.
self.audit_now(if cmd == "seal" { OP_SEAL } else { OP_OPEN }, coord, false);
return Err(msg);
}
// 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);
// 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,
transform,
salt: vec![],
}],
vec![],
);
if cmd == "seal" {
let (mut h, body) = store
.get_record(&coord)
.ok_or_else(|| format!("seal: no record at {path}"))?;
// Carry the owner onto the encrypted record so the gate keeps
// working after sealing (Task 6 B).
h.owner_local_user = self.identity.as_ref().map(|i| i.owner_local.clone());
store
.with_mut(|s| env.put_encrypted(s, coord, Selector::Slot(0), &body, h))
.map_err(|e| format!("seal: {e:?}"))?;
// Log the re-written (encrypted) record to the WAL.
if let Some(v) = store.get_raw(&coord) {
store.log_put(coord, v);
}
Ok(format!("sealed {path} under key {} ({tf})", kc.pack_u32()))
} else {
let (_, envelope) = store
.get_record(&coord)
.ok_or_else(|| format!("open: no record at {path}"))?;
let pt = env
.open(&store.read_snapshot(), Selector::Slot(0), &envelope)
.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())?;
// 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:?}",
kc.pack_u32()
))
}
}
other => Err(format!("unknown command: {other}")),
}
}
}
/// The store view for read commands: the live store normally, or the
/// transaction's `BEGIN` snapshot while a transaction is open (so reads are
/// isolated from concurrent external writers — [`Txn::snapshot`]).
pub fn txn_snapshot(s: &Session) -> CubeStore<HashBackend> {
match &s.txn {
Some(t) => t.snapshot.clone(),
None => s.store.read_snapshot(),
}
}
/// Compute the coordinate a `store_code_cell` call would target, without
/// writing — used to buffer `prog`/`write` mutations during a transaction.
fn scratch_code_coord(
path: &str,
kind: Kind,
name: &str,
code: &[Op],
) -> Result<Czyx, String> {
let mut scratch = CubeStore::new(HashBackend::new());
crate::store_code_cell(&mut scratch, path, kind, name, &[], code, None, None)
.map_err(|e| e.to_string())
}
/// Build the `CubeHeader` a `store_code_cell` call would attach (mirrors
/// `crate::store_code_cell`), so a buffered txn put carries the same header.
/// `owner` (when set) is stamped on `owner_local_user` for Task 6 enforcement.
/// `descriptor` (when set) stamps the behavior-descriptor bits (PDF §524–525).
fn header_for_code(
kind: Kind,
name: &str,
code: &[Op],
owner: Option<&str>,
descriptor: Option<cubecode::Behavior>,
) -> CubeHeader {
let mut h = CubeHeader::new();
h.title = Some(name.to_string());
h.doc_type = Some(kind.as_str().to_string());
h.linked_records = Vec::new();
h.owner_local_user = owner.map(|o| o.to_string());
if let Some(b) = descriptor {
h.flags.0 |= b.to_flags();
}
if h.doc_type.as_deref() == Some("fn") {
h.size_bytes = Some(cubecode::encode(code).len() as u64);
}
h.refresh_flags();
h
}
/// Parse a byte token: decimal (`42`) or hex (`0x2a`).
fn parse_byte(t: &str) -> Option<u8> {
if let Ok(v) = t.parse::<u8>() {
return Some(v);
}
u8::from_str_radix(t.trim_start_matches("0x"), 16).ok()
}
/// Parse a coordinate `C.Z.Y.X` (decimal, allows 0 for Null space).
pub fn parse_coord(s: &str) -> Option<cubecoords::Czyx> {
let parts: Vec<&str> = s.split('.').collect();
if parts.len() != 4 {
return None;
}
let nums: Option<Vec<u8>> = parts.iter().map(|p| p.parse::<u8>().ok()).collect();
let nums = nums?;
Some(cubecoords::Czyx::new(nums[0], nums[1], nums[2], nums[3]))
}
/// Parse a single `u8` axis token: decimal (`7`) or hex (`0x07`).
fn parse_u8(t: Option<&str>, what: &str) -> Result<u8, String> {
let t = t.ok_or_else(|| what.to_string())?;
t.parse::<u8>()
.or_else(|_| u8::from_str_radix(t.trim_start_matches("0x"), 16))
.map_err(|_| format!("{what} (got '{t}')"))
}
/// Encode bytes as a lowercase hex string (used by the raw coordinate API so
/// payloads survive the text socket framing).
fn hex_encode(b: &[u8]) -> String {
let mut s = String::with_capacity(b.len() * 2);
for byte in b {
s.push_str(&format!("{byte:02x}"));
}
s
}
/// Decode a hex string into bytes. Rejects odd length / non-hex.
fn hex_decode(s: &str) -> Option<Vec<u8>> {
if !s.len().is_multiple_of(2) {
return None;
}
let bytes = s.as_bytes();
let mut out = Vec::with_capacity(s.len() / 2);
let mut i = 0;
while i < bytes.len() {
let hi = (bytes[i] as char).to_digit(16)?;
let lo = (bytes[i + 1] as char).to_digit(16)?;
out.push(((hi << 4) | lo) as u8);
i += 2;
}
Some(out)
}
/// Parse a cubevm op name (case-insensitive) into an [`Op`]. `arg` is the
/// operand byte for ops that take one (const/load/store/jmp/jz/jnz/call/ret/
/// syscall); it is ignored for argument-less ops.
fn make_op(t: &str, arg: u8) -> Result<Op, String> {
Ok(match t.to_ascii_lowercase().as_str() {
"nop" => Op::Nop,
"halt" => Op::Halt,
"const" => Op::Const(arg),
"load" => Op::Load(arg),
"store" => Op::Store(arg),
"add" => Op::Add,
"sub" => Op::Sub,
"mul" => Op::Mul,
"div" => Op::Div,
"mod" => Op::Mod,
"and" => Op::And,
"or" => Op::Or,
"xor" => Op::Xor,
"shl" => Op::Shl,
"shr" => Op::Shr,
"eq" => Op::Eq,
"ne" => Op::Ne,
"lt" => Op::Lt,
"gt" => Op::Gt,
"le" => Op::Le,
"ge" => Op::Ge,
"jmp" => Op::Jmp(arg),
"jz" => Op::Jz(arg),
"jnz" => Op::Jnz(arg),
"call" => Op::CallLink(arg),
"ret" => Op::Ret,
"syscall" => Op::Syscall(arg),
other => return Err(format!("prog: unknown op {other}")),
})
}
/// True for ops that consume the next token as a u8 operand.
fn takes_arg(t: &str) -> bool {
matches!(
t.to_ascii_lowercase().as_str(),
"const" | "load" | "store" | "jmp" | "jz" | "jnz" | "call" | "syscall"
)
}
/// Execute a single cube command line against a [`ConcurrentStore`], returning
/// the reply text. Wraps the existing [`Session`] interpreter so the daemon,
/// the REPL, and the multi-tenant router all share one code path (no command
/// drift).
pub fn exec_on_store(store: &Arc<ConcurrentStore>, line: &str) -> Result<String, String> {
let mut sess = Session::with_store(store.clone());
sess.exec(line)
}
/// Parse a transform token into a [`TransformId`].
pub fn parse_transform(s: &str) -> Option<TransformId> {
match s {
"none" => Some(TransformId::None),
"gcm" => Some(TransformId::Aes256Gcm),
"chacha" => Some(TransformId::ChaCha20Poly1305),
"xts" => Some(TransformId::Aes256Xts),
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::DurabilityConfig;
use std::str::FromStr;
use std::sync::Arc;
fn session() -> Session {
Session::with_store(Arc::new(ConcurrentStore::memory()))
}
#[test]
fn os_kernels_live_in_cube_with_call_graph_and_descriptors() {
// Step 1: OS operator behavior lives in CUBE as kernels, composed via
// a call graph (linked_records) and tagged with behavior descriptors.
// Step 2: a thin native effector (`native-apply`) runs each kernel in
// the VM, reads its computed result, and emits the OS effect — the
// cubevm itself never does store-IO (spec-aligned, PDF §524–525).
let mut s = session();
// Lay down the OS kernel call graph.
let out = s.exec("tick").expect("tick should lay down kernels");
assert!(out.contains("normalize-config"), "cfg kernel missing: {out}");
assert!(out.contains("decide-snapshot"), "decide kernel missing: {out}");
assert!(out.contains("summarize-procs"), "summarize kernel missing: {out}");
assert!(out.contains("cube-os-tick"), "tick kernel missing: {out}");
// The root links the three leaves (call graph, not foreign code).
assert!(out.contains("links cfg,decide,summarize"), "call graph not wired: {out}");
// Behavior descriptors are stamped (round-trip through header flags).
assert!(out.contains("[pure]") && out.contains("[io]") && out.contains("[pure,hot]"),
"behavior descriptors not stamped: {out}");
// Run the root kernel: it must traverse the call graph (CallLink 0..2)
// and return, proving the OS's behavior lives as addressable kernels.
let run_out = s.exec("run /c210/z001/y001/x004").expect("tick kernel must run");
assert!(run_out.contains("Halted"), "tick kernel should halt: {run_out}");
// Step 2 — the effector reads the COMPUTED result and emits the effect.
let eff = s.exec("native-apply /c210/z001/y001/x002")
.expect("effector must run decide-snapshot");
assert!(eff.contains("computed result = 1"), "decide kernel result wrong: {eff}");
assert!(eff.contains("OS EFFECT"), "effector must emit OS EFFECT: {eff}");
assert!(eff.contains("snapshot NOW"), "decision=1 should snapshot: {eff}");
// A plain pure kernel also routes through the effector with no store-IO.
let eff2 = s.exec("native-apply /c210/z001/y001/x003")
.expect("effector must run summarize-procs");
assert!(eff2.contains("computed result = 20"), "summarize result wrong: {eff2}");
assert!(eff2.contains("descriptors=[\"pure\", \"hot\"]"), "descriptor readback wrong: {eff2}");
}
#[test]
fn begin_commit_applies_buffered_writes() {
let mut s = session();
s.exec("begin").unwrap();
// buffered, not yet visible
assert!(s
.exec("prog /c001/z001/y001/x001 const 2 const 3 add halt")
.is_ok());
assert!(s.store.get_raw(&Czyx::new(1, 1, 1, 1)).is_none());
// COMMIT makes all ops durable+visible at once
s.exec("commit").unwrap();
let v = s.store.get_raw(&Czyx::new(1, 1, 1, 1));
assert!(
v.is_some(),
"prog buffered during txn must appear after commit"
);
}
#[test]
fn rollback_discards_buffered_writes() {
let mut s = session();
s.exec("begin").unwrap();
let _ = s.exec("write /c002/z001/y001/x001 2a 2b 3c");
s.exec("rollback").unwrap();
assert!(s.store.get_raw(&Czyx::new(2, 1, 1, 1)).is_none());
// begin/commit nested errors
assert!(s.exec("begin").is_ok());
assert!(s.exec("begin").is_err()); // already in txn
assert!(s.exec("commit").is_ok());
}
#[test]
fn owner_enforcement_blocks_cross_owner_overwrite() {
use crate::tenant::{TenantId, TenantIdentity};
let mut s = session();
// Stamp an identity (Task 3 path) for owner "alice".
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
// Alice writes a record — it should be stamped + succeed.
assert!(s.exec("prog /c050/z001/y001/x001 const 7 halt").is_ok());
assert_eq!(
s.store.owner(&Czyx::new(50, 1, 1, 1)).as_deref(),
Some("alice")
);
// Now "bob" tries to overwrite alice's coord — must be rejected.
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "bob".to_string(),
owner_remote: None,
});
let r = s.exec("prog /c050/z001/y001/x001 const 8 halt");
assert!(r.is_err(), "cross-owner overwrite must be rejected");
assert!(r.unwrap_err().contains("owner violation"));
// The record must still be alice's (unchanged value).
assert_eq!(
s.store.owner(&Czyx::new(50, 1, 1, 1)).as_deref(),
Some("alice")
);
// And a delete by bob is likewise rejected.
let d = s.exec("del /c050/z001/y001/x001");
assert!(d.is_err());
assert!(d.unwrap_err().contains("owner violation"));
}
#[test]
fn owner_enforcement_allows_first_claim_and_same_owner() {
use crate::tenant::{TenantId, TenantIdentity};
let mut s = session();
// A session with NO identity writes freely (legacy / test path).
assert!(s.exec("prog /c051/z001/y001/x001 const 1 halt").is_ok());
assert_eq!(s.store.owner(&Czyx::new(51, 1, 1, 1)), None);
// Alice claims the previously-unowned coord — allowed (first write).
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
assert!(s.exec("prog /c051/z001/y001/x001 const 2 halt").is_ok());
assert_eq!(
s.store.owner(&Czyx::new(51, 1, 1, 1)).as_deref(),
Some("alice")
);
// Same owner overwrites freely.
assert!(s.exec("prog /c051/z001/y001/x001 const 3 halt").is_ok());
}
#[test]
fn for_tenant_carries_identity_to_session() {
// Mirrors the cube-server daemon flow: a connection resolves its tenant
// session, stamps identity via HELLO, then runs a per-connection
// `Session::for_tenant` for the rest of the connection. The identity
// MUST propagate so owner enforcement fires on the live path.
use crate::tenant::{TenantConfig, TenantId, TenantIdentity, TenantRegistry};
use std::str::FromStr;
let registry = TenantRegistry::with_config(TenantConfig::Memory);
let ts = registry
.get_or_provision(TenantId::from_str("alpha").unwrap())
.unwrap();
ts.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
let mut s = Session::for_tenant(&ts);
// The per-connection session must carry the stamped identity.
assert_eq!(s.identity().unwrap().owner_local, "alice");
// Alice writes her coord (stamped).
assert!(s.exec("prog /c052/z001/y001/x001 const 7 halt").is_ok());
assert_eq!(
s.store.owner(&Czyx::new(52, 1, 1, 1)).as_deref(),
Some("alice")
);
// A second connection as bob shares the same store (legacy mode) but
// must be blocked from overwriting alice's record.
ts.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "bob".to_string(),
owner_remote: None,
});
let mut b = Session::for_tenant(&ts);
let r = b.exec("prog /c052/z001/y001/x001 const 8 halt");
assert!(r.is_err(), "cross-owner overwrite must be rejected live");
assert!(r.unwrap_err().contains("owner violation"));
}
// --- Plan R6: audit-trail emission (plan R6) ---------------------------
// These tests prove every mutating/read op appends to the append-only audit
// log with the correct op tag, and — critically — that DENIED attempts are
// also logged (the rows an intrusion reviewer cares about).
fn audited_session() -> Session {
let mut s = session();
s.enable_audit();
s
}
/// Parse the audit dump into (op, coord, owner, ok) rows.
fn parse_audit(dump: &Option<String>) -> Vec<(u8, String, String, bool)> {
let dump = dump.as_deref().unwrap_or("");
let mut rows = Vec::new();
for line in dump.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
// {"seq":N,"ts":U,"op":B,"coord":"C.Z.Y.X","owner":"who","ok":bool}
let op = line
.split("\"op\":")
.nth(1)
.and_then(|s| s.split(',').next())
.and_then(|s| s.trim().parse::<u8>().ok())
.unwrap();
let coord = line
.split("\"coord\":\"")
.nth(1)
.and_then(|s| s.split('"').next())
.unwrap()
.to_string();
let owner = line
.split("\"owner\":\"")
.nth(1)
.and_then(|s| s.split('"').next())
.unwrap_or("")
.to_string();
let ok = line.contains("\"ok\":true");
rows.push((op, coord, owner, ok));
}
rows
}
#[test]
fn write_emits_audit_entry_for_owner() {
use crate::tenant::{TenantId, TenantIdentity};
let mut s = audited_session();
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
s.exec("prog /c060/z001/y001/x001 const 1 halt").unwrap();
let rows = parse_audit(&s.audit_dump());
assert_eq!(rows.len(), 1, "exactly one audit row expected");
assert_eq!(rows[0].0, crate::audit::OP_WRITE);
assert_eq!(rows[0].1, Czyx::new(60, 1, 1, 1).pack_u32().to_string());
assert_eq!(rows[0].2, "alice");
assert!(rows[0].3, "permitted op must log ok:true");
}
#[test]
fn denied_cross_owner_write_is_audited() {
use crate::tenant::{TenantId, TenantIdentity};
let mut s = audited_session();
// Alice owns the coord.
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
s.exec("prog /c061/z001/y001/x001 const 1 halt").unwrap();
// Bob tries to overwrite — must be DENIED and logged ok:false.
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "bob".to_string(),
owner_remote: None,
});
assert!(s.exec("prog /c061/z001/y001/x001 const 2 halt").is_err());
let rows = parse_audit(&s.audit_dump());
// Row 0: alice write (ok). Row 1: bob denied write (ok:false).
assert_eq!(rows.len(), 2, "alice write + bob denied write");
assert_eq!(rows[0].0, crate::audit::OP_WRITE);
assert!(rows[0].3);
assert_eq!(rows[1].0, crate::audit::OP_WRITE);
assert!(!rows[1].3, "denied op must log ok:false");
assert_eq!(rows[1].2, "bob");
}
#[test]
fn read_ops_emit_audit_entries() {
use crate::tenant::{TenantId, TenantIdentity};
let mut s = audited_session();
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
// Create + read back.
s.exec("prog /c062/z001/y001/x001 const 9 halt").unwrap();
s.exec("stat /c062/z001/y001/x001").unwrap();
s.exec("run /c062/z001/y001/x001").unwrap();
let rows = parse_audit(&s.audit_dump());
// write, read(stat), read(run)
assert_eq!(rows.len(), 3);
assert_eq!(rows[0].0, crate::audit::OP_WRITE);
assert_eq!(rows[1].0, crate::audit::OP_READ);
assert!(rows[1].3);
assert_eq!(rows[2].0, crate::audit::OP_READ);
assert!(rows[2].3);
}
#[test]
fn denied_read_is_audited() {
use crate::tenant::{TenantId, TenantIdentity};
let mut s = audited_session();
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
s.exec("prog /c063/z001/y001/x001 const 9 halt").unwrap();
// Bob (different owner) must be denied the read and logged ok:false.
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "bob".to_string(),
owner_remote: None,
});
assert!(s.exec("stat /c063/z001/y001/x001").is_err());
let rows = parse_audit(&s.audit_dump());
assert_eq!(rows.len(), 2);
assert_eq!(rows[1].0, crate::audit::OP_READ);
assert!(!rows[1].3, "denied read must log ok:false");
}
#[test]
fn grant_and_revoke_emit_audit_entries() {
let mut s = audited_session();
// grant/revoke require an identity under --require-identity; here we
// emulate the legacy "root" granter by NOT enforcing identity, so the
// op path runs and emits OP_GRANT / OP_REVOKE.
s.exec("grant alice r 100.1.1.1").unwrap();
let rows = parse_audit(&s.audit_dump());
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].0, crate::audit::OP_GRANT);
assert!(rows[0].3);
s.exec("revoke alice 100.1.1.1").unwrap();
let rows = parse_audit(&s.audit_dump());
assert_eq!(rows.len(), 2);
assert_eq!(rows[1].0, crate::audit::OP_REVOKE);
assert!(rows[1].3);
}
#[test]
fn seal_emits_audit_entry() {
use crate::tenant::{TenantId, TenantIdentity};
let mut s = audited_session();
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
s.exec("prog /c064/z001/y001/x001 const 9 halt").unwrap();
// seal auto-materializes the key-cell (32-byte demo key) itself.
// The demo seal crypto may not round-trip in the lib harness, but the
// audit line MUST be emitted (and it is, before the crypto runs) — that
// is what this test proves.
let _ = s.exec("seal /c064/z001/y001/x001 1.1.1.1 none");
let rows = parse_audit(&s.audit_dump());
// write (record) + seal (audit)
assert_eq!(rows.len(), 2);
assert_eq!(rows[1].0, crate::audit::OP_SEAL);
assert!(rows[1].3);
}
#[test]
fn audit_disabled_by_default_emits_nothing() {
// A plain session (no enable_audit call) must not churn the audit log.
let mut s = session();
s.exec("prog /c065/z001/y001/x001 const 1 halt").unwrap();
// The session's audit is None, so even calling the dump command errors.
assert!(s.exec("audit").is_err());
}
#[test]
fn enforce_owner_requires_identity() {
// Task 6 (B): when the daemon policy requires a HELLO identity, a
// mutating op from a session with no identity is rejected.
use crate::tenant::{TenantId, TenantIdentity};
use std::str::FromStr;
let mut s = Session::with_store(Arc::new(ConcurrentStore::memory()));
s.enforce_owner = true;
// No identity => rejected.
let r = s.exec("prog /c053/z001/y001/x001 const 1 halt");
assert!(
r.is_err(),
"anonymous write must be rejected under --require-identity"
);
assert!(r.unwrap_err().contains("require a HELLO identity"));
// A no-identity del/write is likewise rejected.
assert!(s.exec("del /c053/z001/y001/x001").is_err());
// But once an identity is stamped, the same owner may write.
s.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
assert!(s.exec("prog /c053/z001/y001/x001 const 2 halt").is_ok());
}
#[test]
fn seal_open_respect_owner() {
// Task 6 (B): the owner gate is wired into the seal/open branch and
// fires before any crypto. Cross-owner and anonymous access are
// rejected at the gate; we don't exercise the (pre-existing demo)
// key-cell crypto here.
use crate::tenant::{TenantConfig, TenantId, TenantIdentity, TenantRegistry};
use std::str::FromStr;
let registry = TenantRegistry::with_config(TenantConfig::Memory);
let ts = registry
.get_or_provision(TenantId::from_str("alpha").unwrap())
.unwrap();
ts.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
let mut s = Session::for_tenant(&ts);
s.set_enforce_owner(true);
// Alice writes her code cell (owned by alice).
assert!(s.exec("prog /c054/z001/y001/x001 const 7 halt").is_ok());
// Bob (enforce_owner on) cannot seal alice's record -> gate rejects
// before any crypto runs.
ts.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "bob".to_string(),
owner_remote: None,
});
let mut b = Session::for_tenant(&ts);
b.set_enforce_owner(true);
let r = b.exec("seal /c054/z001/y001/x001 5.5.5.5 none");
assert!(r.is_err(), "non-owner seal must be rejected at the gate");
assert!(r.unwrap_err().contains("owner violation"));
// An anonymous session (enforce_owner on, no identity) is also rejected
// on seal/open.
let mut anon = Session::with_store(Arc::new(ConcurrentStore::memory()));
anon.set_enforce_owner(true);
let r = anon.exec("seal /c054/z001/y001/x001 5.5.5.5 none");
assert!(r.is_err(), "anonymous seal must be rejected");
assert!(r.unwrap_err().contains("require a HELLO identity"));
}
#[test]
fn txn_isolation_begin_snapshot_hides_live_writer() {
// Stream A opens a txn and snapshots; stream B mutates the live store.
// A's reads (run) must NOT see B's write until A commits.
let mut a = session();
let mut b = session();
a.store = b.store.clone(); // shared backend, two sessions
b.exec("prog /c003/z001/y001/x001 const 7 halt").unwrap();
assert!(b.store.get_raw(&Czyx::new(3, 1, 1, 1)).is_some());
// A begins AFTER B's write -> snapshot already has it
a.exec("begin").unwrap();
// B deletes it live; A's snapshot must still see it via its txn view
b.exec("del /c003/z001/y001/x001").unwrap();
let snap = txn_snapshot(&a);
assert!(
snap.get_raw(&Czyx::new(3, 1, 1, 1)).is_some(),
"txn snapshot must isolate A from B's concurrent delete"
);
// once A commits (no own writes) the live store reflects B's delete
a.exec("commit").unwrap();
assert!(b.store.get_raw(&Czyx::new(3, 1, 1, 1)).is_none());
}
#[test]
fn commit_is_durable_across_reopen() {
let dir = std::env::temp_dir().join(format!("cube2-txn-{}", 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 cfg = DurabilityConfig::default();
let cs = ConcurrentStore::open(
db.to_str().unwrap(),
wal.to_str().unwrap(),
rec.to_str().unwrap(),
cfg,
)
.unwrap();
let mut s = Session::with_store(Arc::new(cs));
s.exec("begin").unwrap();
s.exec("prog /c004/z001/y001/x001 const 9 halt").unwrap();
s.exec("commit").unwrap();
s.store.checkpoint();
drop(s);
// Reopen: the committed txn must replay from the WAL.
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(4, 1, 1, 1)).is_some(),
"committed txn must survive reopen via WAL replay"
);
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);
}
#[test]
fn read_gate_blocks_non_owner_and_allows_read_grant() {
// R5: `run` (a read+execute op) honors the same owner/grantee gate as
// writes. alice owns a cell; bob cannot read it under enforce_owner;
// once alice grants bob read, bob's read succeeds.
use crate::grants::grant;
use crate::tenant::{TenantConfig, TenantId, TenantIdentity, TenantRegistry};
use std::str::FromStr;
let registry = TenantRegistry::with_config(TenantConfig::Memory);
let ts = registry
.get_or_provision(TenantId::from_str("alpha").unwrap())
.unwrap();
ts.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
let mut a = Session::for_tenant(&ts);
a.set_enforce_owner(true);
// Alice writes a runnable cell she owns.
assert!(a.exec("prog /c070/z001/y001/x001 const 7 halt").is_ok());
// Bob is a different owner on the same (shared) store.
ts.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "bob".to_string(),
owner_remote: None,
});
let mut b = Session::for_tenant(&ts);
b.set_enforce_owner(true);
let r = b.exec("run /c070/z001/y001/x001");
assert!(
r.is_err(),
"non-owner read must be gated under --require-identity"
);
assert!(r.unwrap_err().contains("owner violation"));
// Alice grants bob READ on the coord.
ts.set_identity(TenantIdentity {
tenant: TenantId::from_str("alpha").unwrap(),
owner_local: "alice".to_string(),
owner_remote: None,
});
let a2 = Session::for_tenant(&ts);
let _ = grant(
&a2.store,
&crate::grants::Owner::from(&a2.identity().unwrap()),
&crate::grants::Owner {
local: "bob".to_string(),
remote: None,
},
crate::grants::PERM_READ,
Some(Czyx::new(70, 1, 1, 1)),
);
// Now bob's read passes via grant.
let r = b.exec("run /c070/z001/y001/x001");
assert!(r.is_ok(), "read grant must allow bob's read: {r:?}");
}
#[test]
fn read_gate_requires_identity_under_enforce() {
// R5: a `stat` (pure metadata read) from an anonymous session is
// rejected under --require-identity, exactly like a write.
let dir = std::env::temp_dir().join(format!("cube2-rgate-{}", std::process::id()));
let _ = std::fs::create_dir_all(&dir);
let cs = 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();
let mut s = Session::with_store(Arc::new(cs));
s.set_enforce_owner(true);
let r = s.exec("stat /c071/z001/y001/x001");
assert!(r.is_err(), "anonymous stat must be gated");
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);
}
}