Files
cubelinux-2/cubesys/src/tenant.rs
T
CUBELinux-2 ec84fbc725 cubesys: coalesce WAL fsync on commit; stop over-reporting durable seq
commit_txn issued an unconditional fsync per COMMIT, so N concurrent
writers serialized behind N disk syncs (p99 hit the 3s socket timeout
under 8 writers). Add Wal::sync_upto: committers queue on an fsync_gate,
the first one flushes the whole accumulated buffer, and waiters that find
committed_seq past their target return with zero I/O. N commits now cost
~1 fsync with the same durability guarantee.

Also fix a durability over-report: flush_pending stamped committed_seq
from the LIVE seq counter, so sequences taken by appenders that had not
yet buffered their bytes were reported durable. Track max_seq alongside
the pending buffer and advance committed_seq only to what was written.

Wire --wal-fsync-ms / --checkpoint-ms in cube-server (previously
hardcoded to defaults, so the documented knob did nothing).

Both regressions are mutation-verified: each test fails when its bug is
reintroduced.
2026-08-11 18:09:29 -04:00

588 lines
23 KiB
Rust

//! Multi-tenant registry — the routing layer for a genuinely concurrent,
//! per-tenant-isolated CUBELinux-2 store (see the plan at
//! `.hermes/plans/2026-08-11_041500-concurrent-multitenant-db.md`).
//!
//! **Task 1** landed the [`TenantId`] type and a [`TenantRegistry`] skeleton.
//! **Task 2** added disk-backed per-tenant stores: each tenant's
//! [`ConcurrentStore`] is opened under its own sanitized directory under
//! `store_dir/<tenant>/`, so tenants are isolated on disk and survive a
//! process restart (reusing the durable WAL + checkpoint machinery, including
//! the `49698af` delta-path fix). Tasks 3-7 wire the daemon, add the
//! read/write gate, transactions, and owner/grant enforcement.
use std::collections::HashMap;
use std::path::PathBuf;
use std::str::FromStr;
use std::sync::{Arc, RwLock};
use crate::store::{ConcurrentStore, DurabilityConfig};
/// Identifies a tenant — the PDF's C-axis environment selector (e.g.
/// `agent-a`, `hermes`, `cloud`). Tenant ids are cheap to copy-compare and
/// serve as the key that isolates one tenant's store from every other.
///
/// Ids are normalized (trimmed) and must be non-empty; they are *not* Copy
/// because the backing string is owned. `Clone/Eq/Hash` is all the registry
/// needs for keying.
#[derive(Clone, Eq, PartialEq, Hash, Debug)]
pub struct TenantId(String);
/// Error type for parsing a [`TenantId`] from a string.
#[derive(Clone, Eq, PartialEq, Debug)]
pub struct TenantIdError(pub String);
impl std::fmt::Display for TenantIdError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "invalid tenant id: {}", self.0)
}
}
impl std::error::Error for TenantIdError {}
impl FromStr for TenantId {
type Err = TenantIdError;
/// Parse a tenant id, trimming surrounding whitespace and rejecting an
/// empty result. Path separators are allowed (the id is a logical name,
/// not a filesystem path — Task 2 derives the on-disk dir from it by
/// sanitization, not by trusting the raw string).
fn from_str(s: &str) -> Result<Self, Self::Err> {
let trimmed = s.trim();
if trimmed.is_empty() {
return Err(TenantIdError("tenant id must not be empty".into()));
}
Ok(TenantId(trimmed.to_string()))
}
}
impl TenantId {
/// The raw id string (trimmed, non-empty).
pub fn as_str(&self) -> &str {
&self.0
}
}
/// One tenant's view of the store. For Task 1 this is an in-memory
/// Configuration for how a tenant's store is materialized.
///
/// `Memory` keeps the Task 1 in-memory behavior (handy for tests and the
/// REPL). `Disk` opens the tenant's [`ConcurrentStore`] at
/// `store_dir/<tenant>/` so the data is isolated on disk and survives a
/// restart — the real multi-tenant layout. `SharedFile` is the legacy
/// single-store mode (one `db`/`wal`/`recovery` triple for every tenant),
/// kept so the pre-existing daemon invocation (`--store PATH`) and `stress.sh`
/// keep working unchanged — see [`TenantRegistry::get_or_provision_shared`].
#[derive(Clone, Debug)]
pub enum TenantConfig {
/// In-memory, non-durable (Task 1 default).
Memory,
/// Disk-backed durable store rooted at `store_dir/`. The tenant subdir is
/// derived from the (sanitized) tenant id, so the on-disk path can never
/// escape `store_dir` regardless of what id the client presents.
Disk {
/// Root directory under which each tenant gets its own subdirectory.
store_dir: PathBuf,
/// Durability tuning forwarded to the store's WAL + checkpoint.
durability: DurabilityConfig,
},
/// Legacy: every tenant (and the no-HELLO default) shares ONE durable
/// store at the given `db`/`wal`/`recovery` paths. Used by the daemon when
/// invoked with `--store PATH` and by `stress.sh`; it is single-tenant by
/// construction but lets Task 3 land without breaking existing callers.
SharedFile {
/// Durable database file path.
db: PathBuf,
/// Write-ahead log path.
wal: PathBuf,
/// Recovery-log path (written when WAL replay is needed at startup).
recovery: PathBuf,
/// Durability tuning forwarded to the store's WAL + checkpoint.
durability: DurabilityConfig,
},
}
impl TenantConfig {
/// Convenience: disk-backed tenant config at `store_dir` with default
/// durability.
pub fn disk(store_dir: impl Into<PathBuf>) -> Self {
TenantConfig::Disk {
store_dir: store_dir.into(),
durability: DurabilityConfig::default(),
}
}
/// Disk-backed tenant config at `store_dir` with explicit durability
/// tuning, so the daemon's `--wal-fsync-ms` / `--checkpoint-ms` flags
/// actually reach each per-tenant store's WAL.
pub fn disk_with(store_dir: impl Into<PathBuf>, durability: DurabilityConfig) -> Self {
TenantConfig::Disk {
store_dir: store_dir.into(),
durability,
}
}
}
/// A client's asserted identity, declared via `HELLO <tenant> <owner_local>`
/// `[<owner_remote>]`. This is the PDF's owner-centric model: the same
/// `owner_local_user`/`owner_remote_user` fields `cubecoords::CubeHeader`
/// already carries on every record. Task 6 enforces it; Task 3 only stamps it
/// so subsequent commands run "as" that owner.
#[derive(Clone, Eq, PartialEq, Debug)]
pub struct TenantIdentity {
/// The tenant axis (C-axis selector).
pub tenant: TenantId,
/// The local owner user the client asserts (maps to `CubeHeader::owner_local_user`).
pub owner_local: String,
/// Optional remote owner user (maps to `CubeHeader::owner_remote_user`).
pub owner_remote: Option<String>,
}
/// One tenant's view of the store. Holds a durable `ConcurrentStore` (memory,
/// disk, or shared-file) plus the tenant's identity. Task 3 adds the client
/// identity stamp; Tasks 4-7 add the read/write gate, transactions, and auth.
///
/// Holding the store in an `Arc` lets many connection threads share one
/// tenant's store cheaply, and lets the registry hand out a cheap clone of
/// the handle on every `get_or_provision`.
pub struct TenantSession {
/// The tenant this session belongs to (useful for diagnostics/telemetry).
pub id: TenantId,
/// The tenant-isolated store. Disk-backed when configured; Task 1 used the
/// in-memory variant.
pub store: Arc<ConcurrentStore>,
/// The client identity stamped by `HELLO` (Task 3). `None` until a HELLO
/// arrives on the connection (or for the no-HELLO default tenant).
identity: RwLock<Option<TenantIdentity>>,
}
impl TenantSession {
/// Stamp (replace) the client identity for this session — called after a
/// successful `HELLO`. Returns the previously-stamped identity (if any).
pub fn set_identity(&self, id: TenantIdentity) -> Option<TenantIdentity> {
let mut g = self.identity.write().unwrap();
g.replace(id)
}
/// The currently-stamped identity, if any.
pub fn identity(&self) -> Option<TenantIdentity> {
self.identity.read().unwrap().clone()
}
/// Provision a tenant session per `cfg`. `Memory` yields an in-memory
/// store; `Disk` opens the tenant's store at `store_dir/<safe-id>/`;
/// `SharedFile` opens the single shared store at `db`/`wal`/`recovery`.
/// `<safe-id>` is the tenant id with every path-separator replaced by `_`
/// so a malicious id can't traverse out of `store_dir`.
pub fn open(id: TenantId, cfg: &TenantConfig) -> std::io::Result<Self> {
let store = match cfg {
TenantConfig::Memory => ConcurrentStore::memory(),
TenantConfig::Disk {
store_dir,
durability,
} => {
let safe = id.as_str().replace(['/', '\\', '\0'], "_");
let tenant_dir = store_dir.join(&safe);
let db = tenant_dir.join("cube-store.json");
let wal = tenant_dir.join("cube-store.wal");
let rec = tenant_dir.join("cube-store.recovery.ndjson");
ConcurrentStore::open(
db.to_str().unwrap(),
wal.to_str().unwrap(),
rec.to_str().unwrap(),
*durability,
)?
}
TenantConfig::SharedFile {
db,
wal,
recovery,
durability,
} => ConcurrentStore::open(
db.to_str().unwrap(),
wal.to_str().unwrap(),
recovery.to_str().unwrap(),
*durability,
)?,
};
Ok(TenantSession {
id,
store: Arc::new(store),
identity: RwLock::new(None),
})
}
}
/// The registry that maps tenant ids to their (isolated) sessions.
///
/// Replacement for the single global `Arc<Mutex<Session>>` in `cube-server`:
/// look up a tenant's `TenantSession` here at connection time. Different
/// tenants run fully in parallel; within a tenant, the `ConcurrentStore`
/// already serializes its own writes.
pub struct TenantRegistry {
tenants: RwLock<HashMap<TenantId, Arc<TenantSession>>>,
/// How each (newly provisioned) tenant's store is materialized. Set at
/// registry creation; existing sessions keep the config they were opened
/// with.
config: TenantConfig,
/// When set, `get_or_provision` returns this single session for ANY tenant
/// id — the legacy single-store mode (`--store PATH`, used by old `cubec`
/// and `stress.sh`). `None` means true per-tenant isolation (Disk/Memory).
shared: RwLock<Option<Arc<TenantSession>>>,
}
impl TenantRegistry {
/// Create an empty registry with the given tenant materialization config.
pub fn with_config(config: TenantConfig) -> Self {
TenantRegistry {
tenants: RwLock::new(HashMap::new()),
config,
shared: RwLock::new(None),
}
}
/// Create an empty registry using in-memory tenant stores (Task 1 default).
pub fn new() -> Self {
Self::with_config(TenantConfig::Memory)
}
/// Fetch the session for `id`. In shared mode (`get_or_provision_shared`
/// was called) every tenant id — and the no-HELLO default — resolves to the
/// one registered store. Otherwise the tenant's own store is provisioned on
/// first use per [`TenantConfig`]; a disk open failure propagates to the
/// caller (the daemon should reject the connection rather than silently
/// fall back to memory).
///
/// Many readers can hold the read lock simultaneously; only a true miss
/// (provisioning) takes the write lock, and it drops the lock immediately
/// after inserting so concurrent provisioning of *different* tenants does
/// not serialize.
pub fn get_or_provision(&self, id: TenantId) -> std::io::Result<Arc<TenantSession>> {
// Shared (legacy single-store) mode: every tenant uses the one store.
if let Some(shared) = self.shared.read().unwrap().as_ref() {
return Ok(shared.clone());
}
// Fast path: already provisioned — take only the read lock.
if let Some(existing) = self.tenants.read().unwrap().get(&id).cloned() {
return Ok(existing);
}
// Miss: take the write lock, but re-check in case another thread
// provisioned the same tenant while we were waiting.
let mut guard = self.tenants.write().unwrap();
if let Some(existing) = guard.get(&id).cloned() {
return Ok(existing);
}
let session = Arc::new(TenantSession::open(id.clone(), &self.config)?);
guard.insert(id, session.clone());
Ok(session)
}
/// The tenant id used for connections that never send `HELLO` (legacy
/// `cubec`/`stress.sh` behavior). In `SharedFile` config every HELLO
/// tenant also resolves here, so the whole daemon shares one store.
pub const DEFAULT_TENANT: &'static str = "default";
/// Enable legacy single-store mode: every tenant id (and the no-HELLO
/// default) resolves to `shared`. The daemon uses this when invoked with
/// `--store PATH` (or `stress.sh`), so it keeps working while still
/// accepting (and ignoring the isolation of) `HELLO` frames. The caller is
/// responsible for having opened `shared` against the desired durable
/// paths. Returns the shared session.
///
/// In `Disk`/`Memory` mode this is not called; tenants are provisioned
/// individually by [`get_or_provision`].
pub fn get_or_provision_shared(&self, shared: Arc<TenantSession>) -> Arc<TenantSession> {
// Register under the default tenant id for callers that look it up by
// DEFAULT_TENANT, and flip the shared fallback so HELLO tenants (any
// id) also resolve here.
self.tenants.write().unwrap().insert(
TenantId::from_str(Self::DEFAULT_TENANT).unwrap(),
shared.clone(),
);
*self.shared.write().unwrap() = Some(shared.clone());
shared
}
/// Parse a `HELLO` frame: `HELLO <tenant> <owner_local> [<owner_remote>]`.
/// Returns the identity on success, or an error string describing what was
/// wrong. The tenant is validated (non-empty after trim) but NOT checked
/// against any allow-list here — admission policy (auto-provision vs
/// deny-unknown) is the daemon's concern, not the parser's.
pub fn parse_hello(line: &str) -> Result<TenantIdentity, String> {
let mut it = line.split_whitespace();
let cmd = it.next().ok_or_else(|| "empty HELLO".to_string())?;
if !cmd.eq_ignore_ascii_case("hello") {
return Err(format!("not a HELLO frame: {cmd}"));
}
let tenant_s = it
.next()
.ok_or_else(|| "HELLO needs <tenant>".to_string())?;
let owner_local = it
.next()
.ok_or_else(|| "HELLO needs <owner_local_user>".to_string())?;
let owner_remote = it.next().map(|s| s.to_string());
let tenant = TenantId::from_str(tenant_s).map_err(|e| format!("HELLO: bad tenant: {e}"))?;
if owner_local.trim().is_empty() {
return Err("HELLO: owner_local_user must not be empty".to_string());
}
Ok(TenantIdentity {
tenant,
owner_local: owner_local.to_string(),
owner_remote,
})
}
}
impl Default for TenantRegistry {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn tenant_id_roundtrips() {
let a = TenantId::from_str("agent-a").unwrap();
let a2 = TenantId::from_str("agent-a").unwrap();
let b = TenantId::from_str("agent-b").unwrap();
assert_eq!(a, a2, "same id parses equal");
assert_ne!(a, b, "different ids are distinct");
assert_eq!(a.as_str(), "agent-a");
assert_eq!(b.as_str(), "agent-b");
}
#[test]
fn tenant_id_rejects_empty() {
assert!(TenantId::from_str("").is_err());
assert!(TenantId::from_str(" ").is_err());
// surrounding whitespace is trimmed, not rejected
assert_eq!(
TenantId::from_str(" agent-a ").unwrap().as_str(),
"agent-a"
);
}
#[test]
fn registry_provisions_and_reuses() {
let reg = TenantRegistry::new();
let a = reg
.get_or_provision(TenantId::from_str("agent-a").unwrap())
.unwrap();
let a2 = reg
.get_or_provision(TenantId::from_str("agent-a").unwrap())
.unwrap();
let b = reg
.get_or_provision(TenantId::from_str("agent-b").unwrap())
.unwrap();
// Same tenant returns the SAME Arc (not a new provisioning).
assert!(Arc::ptr_eq(&a, &a2), "repeat lookup reuses the session");
// Different tenant is a different session.
assert!(!Arc::ptr_eq(&a, &b), "distinct tenants are isolated");
assert_eq!(a.id.as_str(), "agent-a");
assert_eq!(b.id.as_str(), "agent-b");
}
#[test]
fn tenants_are_independent_stores() {
let reg = TenantRegistry::new();
let a = reg
.get_or_provision(TenantId::from_str("agent-a").unwrap())
.unwrap();
let b = reg
.get_or_provision(TenantId::from_str("agent-b").unwrap())
.unwrap();
// Each session carries its own store handle.
assert!(!Arc::ptr_eq(&a.store, &b.store));
}
// ---- Task 2: disk-backed per-tenant isolation ----
/// Two tenants must live in separate on-disk directories, and a write to
/// tenant A must be invisible to tenant B after both are reopened from
/// disk. This is the core "hard per-namespace partition" guarantee.
#[test]
fn tenant_store_is_isolated_on_disk() {
let root = std::env::temp_dir().join(format!(
"cubelinux-tenant-isol-{}-{}",
std::process::id(),
"a1b2"
));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(&root).unwrap();
let reg = TenantRegistry::with_config(TenantConfig::disk(&root));
let a = reg
.get_or_provision(TenantId::from_str("agent-a").unwrap())
.unwrap();
let b = reg
.get_or_provision(TenantId::from_str("agent-b").unwrap())
.unwrap();
// Physically separate dirs.
let dir_a = root.join("agent-a");
let dir_b = root.join("agent-b");
assert!(dir_a.is_dir(), "tenant A has its own dir");
assert!(dir_b.is_dir(), "tenant B has its own dir");
// Write a record to A only.
let coord = cubecoords::Czyx::new(1, 1, 1, 1);
a.store
.put_record(coord, &cubecoords::CubeHeader::new(), b"hello-a");
// B must NOT see it (different store).
assert!(
b.store.get_record(&coord).is_none(),
"tenant B must not see tenant A's record"
);
// Force durability, then drop both sessions so the on-disk stores
// checkpoint and their background threads stop.
a.store.checkpoint();
b.store.checkpoint();
drop(a);
drop(b);
// Reopen from disk under a fresh registry.
let reg2 = TenantRegistry::with_config(TenantConfig::disk(&root));
let a2 = reg2
.get_or_provision(TenantId::from_str("agent-a").unwrap())
.unwrap();
let b2 = reg2
.get_or_provision(TenantId::from_str("agent-b").unwrap())
.unwrap();
// A's record survived the restart; B still has nothing there.
assert_eq!(
a2.store.get_record(&coord).map(|(_, v)| v),
Some(b"hello-a".to_vec()),
"tenant A's record persisted on disk"
);
assert!(
b2.store.get_record(&coord).is_none(),
"tenant B still isolated after reopen"
);
// Cleanup.
let _ = std::fs::remove_dir_all(&root);
}
/// A tenant id containing path separators must be sanitized so its store
/// cannot escape the configured `store_dir`.
#[test]
fn tenant_id_cannot_traverse_store_dir() {
let root = std::env::temp_dir().join(format!(
"cubelinux-tenant-traverse-{}-{}",
std::process::id(),
"c3d4"
));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(&root).unwrap();
let reg = TenantRegistry::with_config(TenantConfig::disk(&root));
let _ = reg
.get_or_provision(TenantId::from_str("../escapee").unwrap())
.unwrap();
let _ = reg
.get_or_provision(TenantId::from_str("normal").unwrap())
.unwrap();
// The malicious id lands under root/<sanitized>, NOT outside root.
assert!(
root.join(".._escapee").is_dir(),
"sanitized id stays under store_dir"
);
// The tenant dir must be a child of (contained by) root — i.e. no
// traversal escape actually happened.
assert!(
root.join(".._escapee").starts_with(&root),
"no traversal escape occurred"
);
// Cleanup.
let _ = std::fs::remove_dir_all(&root);
}
// ---- Task 3: HELLO parsing + legacy shared routing ----
#[test]
fn hello_parser_roundtrips() {
let id = TenantRegistry::parse_hello("HELLO alpha luulu").expect("valid HELLO");
assert_eq!(id.tenant.as_str(), "alpha");
assert_eq!(id.owner_local, "luulu");
assert!(id.owner_remote.is_none());
let id2 = TenantRegistry::parse_hello("hello beta luulu remote")
.expect("valid HELLO with remote");
assert_eq!(id2.tenant.as_str(), "beta");
assert_eq!(id2.owner_local, "luulu");
assert_eq!(id2.owner_remote.as_deref(), Some("remote"));
}
#[test]
fn hello_parser_rejects_garbage() {
assert!(TenantRegistry::parse_hello("").is_err());
assert!(TenantRegistry::parse_hello("PING x y").is_err());
assert!(TenantRegistry::parse_hello("HELLO").is_err());
assert!(TenantRegistry::parse_hello("HELLO alpha").is_err());
// empty owner_local
assert!(TenantRegistry::parse_hello("HELLO alpha ").is_err());
}
/// In legacy SharedFile mode every HELLO tenant — and the no-HELLO default
/// — must resolve to the SAME underlying store (backward-compatible with
/// old cubec/stress.sh, which never send HELLO and expect one global
/// store).
#[test]
fn legacy_shared_mode_shares_one_store() {
let tmp = std::env::temp_dir().join(format!(
"cubelinux-legacy-shared-{}-{}",
std::process::id(),
"e5f6"
));
let _ = std::fs::remove_dir_all(&tmp);
std::fs::create_dir_all(&tmp).unwrap();
let cfg = TenantConfig::SharedFile {
db: tmp.join("cube-store.json"),
wal: tmp.join("cube-store.wal"),
recovery: tmp.join("cube-store.recovery.ndjson"),
durability: DurabilityConfig::default(),
};
let reg = TenantRegistry::with_config(cfg.clone());
let shared = Arc::new(
TenantSession::open(
TenantId::from_str(TenantRegistry::DEFAULT_TENANT).unwrap(),
&cfg,
)
.unwrap(),
);
reg.get_or_provision_shared(shared);
// Default tenant and a HELLO tenant both resolve to the shared store.
let def = reg
.get_or_provision(TenantId::from_str(TenantRegistry::DEFAULT_TENANT).unwrap())
.unwrap();
let hello = reg
.get_or_provision(TenantId::from_str("alpha").unwrap())
.unwrap();
assert!(
Arc::ptr_eq(&def.store, &hello.store),
"legacy mode: all tenants share one store"
);
// A write under HELLO is visible under the default (same store).
let coord = cubecoords::Czyx::new(1, 2, 3, 4);
hello
.store
.put_record(coord, &cubecoords::CubeHeader::new(), b"shared");
assert_eq!(
def.store.get_record(&coord).map(|(_, v)| v),
Some(b"shared".to_vec())
);
hello.store.checkpoint();
let _ = std::fs::remove_dir_all(&tmp);
}
}