//! Concurrent, durable cube store — the "real database" backend. //! //! Wraps the in-memory [`CubeStore`] behind a `Mutex` so many //! threads can read and write at once, and adds two durability layers: //! //! 1. A **write-ahead log (WAL)** — every transaction is appended as one //! newline-delimited JSON object to a WAL file and *group-fsynced* on a //! short interval (default 25 ms). The WAL is the "JSON backup of all //! transactions": it covers the gap between a command returning and the //! slower scheduled database checkpoint landing on disk, so the daemon //! keeps its speed while still being crash-safe. //! //! 2. A **scheduled checkpoint** — on a configurable interval (default //! 2000 ms) the whole store is serialized to the durable database file //! (atomic rename) and the WAL is rotated. This is the "already scheduled //! writes into the database." //! //! On startup we load the last checkpoint, then replay any WAL entries newer //! than it. If we actually had to rely on the WAL (i.e. there were //! not-yet-checkpointed transactions at crash time), we write a recovery //! event to a recovery log — satisfying the requirement to *log, in writing, //! any event where we had to fall back to the JSON backup*. use std::fs::{self, File, OpenOptions}; use std::io::Write; use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex, RwLock}; use std::thread::{self, JoinHandle}; use std::time::Duration; use cubecoords::{CubeHeader, Czyx}; use cubestore::{CubeStore, HashBackend}; use std::collections::HashSet; use crate::persist; use cubecode::{Kind, Op}; /// Durability tuning. #[derive(Clone, Copy, Debug)] pub struct DurabilityConfig { /// Checkpoint (full store snapshot) interval in milliseconds. pub checkpoint_ms: u64, /// WAL group-commit fsync interval in milliseconds. Smaller = less data /// loss on crash, more fsync overhead. pub wal_fsync_ms: u64, } impl Default for DurabilityConfig { fn default() -> Self { DurabilityConfig { checkpoint_ms: 2000, wal_fsync_ms: 25, } } } /// When the incremental delta file grows past this many bytes, fold it into a /// fresh full base snapshot and truncate the delta. Keeps on-disk state /// bounded and fast to load, while still doing tiny incremental writes the /// vast majority of the time. 1 MiB is comfortable for this store. const DELTA_COMPACT_BYTES: u64 = 1_048_576; /// One durable transaction in the WAL (newline-delimited JSON, one per line). /// `seq` is monotonic so replay can skip already-checkpointed entries. For /// `op:"put"`, `v` is the raw backend value bytes (header_len||hdr||body) /// hex-encoded; replay does `put_raw(coord, v)` — identical to a normal write, /// so replays are idempotent. For `op:"txn"`, `batch` carries the committed /// batch (puts/deletes) as one durable unit. struct WalEntry { seq: u64, op: WalOp, coord: Czyx, value: Vec, /// Only set for `op:"txn"`; empty otherwise. batch: Vec, } #[derive(Clone, Copy, PartialEq, Eq)] enum WalOp { Put, Delete, /// A transaction commit: a batch of puts/deletes applied atomically. /// `Some(value)` is a put (value is the raw backend bytes), `None` is a /// delete. Replayed idempotently (each entry is a put_raw/delete_raw). Txn, } /// One entry in a [`WalOp::Txn`] batch: `value: Some(bytes)` is a put, /// `value: None` is a delete. Carries exactly what the store needs to apply /// or replay the op — no header parsing required. #[derive(Clone)] pub struct TxnEntry { /// Coordinate to write or delete. pub coord: Czyx, /// `Some(raw_backend_bytes)` => put; `None` => delete. pub value: Option>, } /// Buffered-but-not-yet-fsync'd WAL bytes, plus the highest sequence number /// contained in them. /// /// Tracking `max_seq` alongside the buffer is what makes `committed_seq` /// honest: a flush may only advance `committed_seq` to the highest seq that /// was actually in the batch it wrote, NOT to whatever the global `seq` /// counter happens to read at fsync time. Concurrent appenders bump `seq` /// between the buffer swap and the fsync, and claiming those as durable would /// over-report durability (telemetry would say data survived a crash that in /// fact never reached the platter). #[derive(Default)] struct Pending { buf: String, max_seq: u64, } /// The write-ahead log. Append-only NDJSON, group-fsync'd on a timer. struct Wal { path: PathBuf, /// `None` in memory mode (no durability); `Some` when a file backs it. file: Mutex>, /// Next sequence number to assign. seq: AtomicU64, /// Buffered NDJSON lines not yet fsync'd (group commit). pending: Mutex, /// Serializes fsync attempts so concurrent committers COALESCE into one /// fsync instead of each paying for their own. A waiter re-checks /// `committed_seq` after acquiring this gate: if the thread ahead of it /// already flushed past its target seq, it returns without any I/O. This /// is the group-commit win — N concurrent COMMITs cost ~1 fsync. fsync_gate: Mutex<()>, /// Highest sequence number known to be fsync'd to disk. committed_seq: AtomicU64, /// WAL sequence number already folded into the base snapshot. Replay on /// startup begins after this, so the WAL only carries post-checkpoint /// transactions. cp_seq_wal: AtomicU64, /// WAL sequence number of the last full base snapshot. Used to compute /// the incremental delta (only coords written after `base_seq`). base_seq: AtomicU64, /// Set of coordinates written since `base_seq` (the authoritative changelog /// for the incremental delta). Cleared whenever `base_seq` advances. dirty: Mutex>, /// Shared stop flag (also owned by the group-commit thread). stop: Arc, /// Handle of the group-commit fsync thread (behind a Mutex so the struct /// stays `Sync`; `JoinHandle` is `Send` but not `Sync`). group_thread: Mutex>>, fsync_ms: u64, } impl Wal { /// Open (or create) a durable WAL at `path`, replaying nothing here — the /// caller drives replay via [`Wal::replay_after`]. `cp_seq` is the highest /// sequence already durable in the checkpoint; new entries continue after /// `max(existing_file_seq, cp_seq)`. fn open( path: &Path, cp_seq: u64, fsync_ms: u64, stop: Arc, ) -> std::io::Result> { let existing = fs::read_to_string(path).unwrap_or_default(); let mut max_seq = 0u64; for line in existing.lines() { if let Some(e) = decode_wal(line) { if e.seq > max_seq { max_seq = e.seq; } } } let start_seq = max_seq.max(cp_seq) + 1; let file = OpenOptions::new() .create(true) .read(true) .append(true) .open(path)?; let wal = Arc::new(Wal { path: path.to_path_buf(), file: Mutex::new(Some(file)), seq: AtomicU64::new(start_seq), pending: Mutex::new(Pending::default()), fsync_gate: Mutex::new(()), committed_seq: AtomicU64::new(max_seq.max(cp_seq)), cp_seq_wal: AtomicU64::new(cp_seq), base_seq: AtomicU64::new(cp_seq), dirty: Mutex::new(HashSet::new()), stop: stop.clone(), group_thread: Mutex::new(None), fsync_ms, }); let t = spawn_group(wal.clone()); *wal.group_thread.lock().unwrap() = Some(t); Ok(wal) } /// In-memory WAL: appends and fsyncs are no-ops, no background thread. fn memory() -> Arc { Arc::new(Wal { path: PathBuf::new(), file: Mutex::new(None), seq: AtomicU64::new(1), pending: Mutex::new(Pending::default()), fsync_gate: Mutex::new(()), committed_seq: AtomicU64::new(0), cp_seq_wal: AtomicU64::new(0), base_seq: AtomicU64::new(0), dirty: Mutex::new(HashSet::new()), stop: Arc::new(AtomicBool::new(true)), group_thread: Mutex::new(None), fsync_ms: 0, }) } /// Mark `coord` dirty (written after `base_seq`), for the incremental delta. fn mark_dirty(&self, coord: Czyx) { self.dirty.lock().unwrap().insert(coord); } /// True if `coord` was written after `base_seq` (i.e. belongs in the delta). fn coord_modified_since(&self, coord: &Czyx, base: u64) -> bool { if base == 0 { // No base yet — everything is "since base" only if it was written. return self.dirty.lock().unwrap().contains(coord); } self.dirty.lock().unwrap().contains(coord) } fn set_cp_seq_wal(&self, v: u64) { self.cp_seq_wal.store(v, Ordering::SeqCst); } fn set_base_seq(&self, v: u64) { self.base_seq.store(v, Ordering::SeqCst); self.dirty.lock().unwrap().clear(); } /// Append a transaction. In memory mode this is a no-op. Returns the /// sequence number assigned (0 in memory mode) so a caller that needs /// synchronous durability can wait for exactly that seq. fn append(&self, op: WalOp, coord: Czyx, value: Vec) -> u64 { let has_file = self.file.lock().unwrap().is_some(); if !has_file { return 0; } let s = self.seq.fetch_add(1, Ordering::SeqCst) + 1; self.mark_dirty(coord); let line = encode_wal(&WalEntry { seq: s, op, coord, value, batch: Vec::new(), }); let mut p = self.pending.lock().unwrap(); p.buf.push_str(&line); p.buf.push('\n'); if s > p.max_seq { p.max_seq = s; } s } /// Append a transaction commit. In memory mode this is a no-op. The whole /// batch lands as ONE WAL line (`op:"txn"`), so a crash mid-txn replays as /// a single all-or-nothing unit. Returns the assigned sequence number. fn append_txn(&self, batch: &[TxnEntry]) -> u64 { let has_file = self.file.lock().unwrap().is_some(); if !has_file { return 0; } let s = self.seq.fetch_add(1, Ordering::SeqCst) + 1; for e in batch { self.mark_dirty(e.coord); } let line = encode_wal(&WalEntry { seq: s, op: WalOp::Txn, coord: Czyx::new(0, 0, 0, 0), value: Vec::new(), batch: batch.to_vec(), }); let mut p = self.pending.lock().unwrap(); p.buf.push_str(&line); p.buf.push('\n'); if s > p.max_seq { p.max_seq = s; } s } /// Flush buffered lines to disk and fsync them (group commit). /// /// `committed_seq` advances ONLY to the highest seq actually present in the /// bytes this call wrote — never to the live `seq` counter, which /// concurrent appenders may have already pushed past our batch. fn flush_pending(&self) { let _gate = self.fsync_gate.lock().unwrap(); self.flush_locked(); } /// Flush the current buffer. Caller must hold `fsync_gate`. fn flush_locked(&self) { let (bytes, batch_max) = { let mut p = self.pending.lock().unwrap(); if p.buf.is_empty() { return; } let taken = std::mem::take(&mut *p); (taken.buf, taken.max_seq) }; let mut fopt = self.file.lock().unwrap(); if let Some(f) = fopt.as_mut() { if f.write_all(bytes.as_bytes()).is_ok() && f.flush().is_ok() && f.sync_all().is_ok() { self.committed_seq.fetch_max(batch_max, Ordering::SeqCst); } } } /// Make sequence `target` durable, COALESCING with any concurrent flush. /// /// This is the batch-fsync core. Instead of every COMMIT issuing its own /// `sync_all` (which serialized N writers behind N fsyncs and produced the /// multi-second p99 stalls under 8 concurrent writers), committers queue on /// `fsync_gate`. Whoever gets in first flushes the WHOLE accumulated /// buffer — including the lines appended by everyone still waiting. Those /// waiters then observe `committed_seq >= target` and return with zero /// I/O. N concurrent commits therefore cost ~1 fsync, not N, while every /// caller still gets a true durability guarantee before returning. fn sync_upto(&self, target: u64) { if target == 0 || self.committed_seq.load(Ordering::SeqCst) >= target { return; } let _gate = self.fsync_gate.lock().unwrap(); // Re-check: the holder ahead of us may already have flushed past us. if self.committed_seq.load(Ordering::SeqCst) >= target { return; } self.flush_locked(); } /// Replay entries with `seq > cp_seq` into `store`. Returns the number of /// entries applied (0 means the checkpoint alone was sufficient). fn replay_after(&self, cp_seq: u64, store: &mut CubeStore) -> u64 { let text = fs::read_to_string(&self.path).unwrap_or_default(); let mut applied = 0u64; let mut max = 0u64; for line in text.lines() { if let Some(e) = decode_wal(line) { if e.seq > max { max = e.seq; } if e.seq > cp_seq { match e.op { WalOp::Put => store.put_raw(e.coord, e.value), WalOp::Delete => store.delete_raw(&e.coord), WalOp::Txn => { // A committed batch replays as its puts/deletes. // Idempotent (each is put_raw/delete_raw), so a // crash that left the batch half-applied before the // checkpoint still converges on replay. for te in &e.batch { match &te.value { Some(v) => store.put_raw(te.coord, v.clone()), None => store.delete_raw(&te.coord), } } } } applied += 1; } } } // Continue sequence numbering past anything we saw. let need = max.max(cp_seq) + 1; let cur = self.seq.load(Ordering::SeqCst); if need > cur { self.seq.store(need, Ordering::SeqCst); } let committed = self.committed_seq.load(Ordering::SeqCst); if need > committed { self.committed_seq.store(need, Ordering::SeqCst); } // The base snapshot already reflected up to cp_seq; replay covered the // rest. Record the boundary so future checkpoints compute deltas and // WAL replay starts after it. self.set_cp_seq_wal(cp_seq); applied } /// Signal the group-commit thread to stop. fn set_stop(&self) { self.stop.store(true, Ordering::SeqCst); } } /// Spawn the background thread that group-fsyncs buffered WAL lines. fn spawn_group(wal: Arc) -> JoinHandle<()> { thread::spawn(move || loop { thread::sleep(Duration::from_millis(wal.fsync_ms)); if wal.stop.load(Ordering::SeqCst) { break; } wal.flush_pending(); }) } /// The concurrent, durable store handle shared across threads. /// /// Internally the `CubeStore` is guarded by an `RwLock` (Task 4): many reader /// threads run in parallel, while writes and the scheduled checkpoint take the /// write side. This converts the old "one global mutex serializes every /// command" into "reads scale across cores, writes stay serialized per store" /// — real reader/writer sharding without touching the WAL or the coordinate /// encoding. pub struct ConcurrentStore { inner: Arc>>, wal: Arc, db_path: PathBuf, cp_seq_path: PathBuf, stop: Arc, flush_thread: Arc>>>, cfg: DurabilityConfig, /// In-memory mirror of the rendered audit log (one newline-joined string), /// kept per-store so every `Audit` bound to this store shares it. `append` /// extends it by one line and `dump` returns a clone — both O(1). The /// durable per-entry records remain the source of truth; this is just a /// fast read mirror, rebuilt from records on first `dump` if cold (e.g. /// after a restart that loaded a durable store). audit_tail: Mutex, } impl ConcurrentStore { /// In-memory, non-durable store (used by the REPL, scripts, and tests). /// No WAL, no checkpoint thread. pub fn memory() -> Self { ConcurrentStore { inner: Arc::new(RwLock::new(CubeStore::new(HashBackend::new()))), wal: Wal::memory(), db_path: PathBuf::new(), cp_seq_path: PathBuf::new(), stop: Arc::new(AtomicBool::new(true)), flush_thread: Arc::new(Mutex::new(None)), cfg: DurabilityConfig::default(), audit_tail: Mutex::new(String::new()), } } /// Shared, per-store in-memory mirror of the rendered audit log. All /// `Audit` instances bound to this store use this one buffer so concurrent /// connections interleave their entries exactly as the durable per-entry /// records do. `audit.rs` holds the serialization contract. pub(crate) fn audit_tail(&self) -> &Mutex { &self.audit_tail } /// Open (or create) a durable store at `db_path`, with the WAL at /// `wal_path` and recovery events logged to `recovery_log`. Loads the last /// checkpoint, replays any newer WAL entries, and starts the background /// checkpoint thread. pub fn open( db_path: &str, wal_path: &str, recovery_log: &str, cfg: DurabilityConfig, ) -> std::io::Result { // Load the checkpoint: base snapshot + incremental delta. Then replay // any WAL entries newer than the base's WAL boundary. let db_p = PathBuf::from(db_path); let wal_p = PathBuf::from(wal_path); // MUST match `checkpoint_store`'s `db_path.with_extension("delta")` so the // delta is read back from the same path it is written to on checkpoint. let delta_p = db_p.with_extension("delta"); let cp_p = PathBuf::from(format!("{db_path}.seq")); let rec_p = PathBuf::from(recovery_log); if let Some(p) = db_p.parent() { let _ = fs::create_dir_all(p); } if let Some(p) = delta_p.parent() { let _ = fs::create_dir_all(p); } // Highest WAL sequence already folded into the base. let base_cp: u64 = fs::read_to_string(&cp_p) .ok() .and_then(|s| s.trim().parse().ok()) .unwrap_or(0); // Reconstruct the store: base + delta. let mut store = load_base_plus_delta(&db_p, &delta_p); let stop = Arc::new(AtomicBool::new(false)); let wal = Wal::open(&wal_p, base_cp, cfg.wal_fsync_ms, stop.clone())?; // Replay WAL entries newer than the base boundary. If we had to rely on // the WAL, record that fact in the recovery log (in writing). let applied = wal.replay_after(base_cp, &mut store); if applied > 0 { append_recovery_log(&rec_p, applied); } let inner = Arc::new(RwLock::new(store)); let cs = ConcurrentStore { inner, wal, db_path: db_p, cp_seq_path: cp_p, stop, flush_thread: Arc::new(Mutex::new(None)), cfg, audit_tail: Mutex::new(String::new()), }; // Background checkpoint thread. let flush = { let inner = cs.inner.clone(); let db = cs.db_path.clone(); let cp = cs.cp_seq_path.clone(); let wal = cs.wal.clone(); let stop = cs.stop.clone(); let ms = cs.cfg.checkpoint_ms; thread::spawn(move || loop { thread::sleep(Duration::from_millis(ms)); if stop.load(Ordering::SeqCst) { break; } checkpoint_store(&inner, &db, &cp, &wal); }) }; *cs.flush_thread.lock().unwrap() = Some(flush); Ok(cs) } // ---- Reads (read-lock the inner store; many readers run in parallel) ---- /// Raw backend get. pub fn get_raw(&self, key: &Czyx) -> Option> { self.inner.read().unwrap().get_raw(key) } /// Fetch and split a record into `(header, body)`. pub fn get_record(&self, key: &Czyx) -> Option<(CubeHeader, Vec)> { self.inner.read().unwrap().get_record(key) } /// Every coordinate present. pub fn keys(&self) -> Vec { self.inner.read().unwrap().keys() } /// Coordinates under a `C`/`Z`/`Y` prefix. pub fn scan_prefix(&self, c: u8, z: Option, y: Option) -> Vec { self.inner.read().unwrap().scan_prefix(c, z, y) } /// Coordinates whose header lists `target` in `linked_records`. pub fn linked_to(&self, target: &Czyx) -> Vec { self.inner.read().unwrap().linked_to(target) } /// Query by document type (the `doc_type` header field). A real predicate /// over the store, not just a prefix scan. pub fn query_doc_type(&self, dt: &str) -> Vec { let g = self.inner.read().unwrap(); let mut out = Vec::new(); for k in g.keys() { if let Some((h, _)) = g.get_record(&k) { if h.doc_type.as_deref() == Some(dt) { out.push(k); } } } out.sort(); out } /// Reconstruct a directory listing from `path` metatags: every record /// whose `path` header field is under `dir` (the PDF's "operate on CZYX /// records and Null-space flags rather than paths and inodes"). Delegates /// to [`CubeStore::scan_by_path_prefix`]. pub fn query_by_path_prefix(&self, dir: &str) -> Vec { self.inner.read().unwrap().scan_by_path_prefix(dir) } /// Exact-path metatag lookup. Delegates to [`CubeStore::scan_by_path`]. pub fn query_by_path(&self, path: &str) -> Vec { self.inner.read().unwrap().scan_by_path(path) } /// The `owner_local_user` stamped on the record at `key`, if it has one. /// Used by owner enforcement (Task 6): a mutating command may only /// overwrite a record whose owner matches the session's identity owner. pub fn owner(&self, key: &Czyx) -> Option { self.inner .read() .unwrap() .get_record(key) .and_then(|(h, _)| h.owner_local_user.clone()) } /// A consistent point-in-time snapshot of the whole store. Used by the VM /// and cubefs, which take a `CubeStore` by value. A read-lock clone, so it /// does not block concurrent readers (it only waits for an in-flight /// writer to release). pub fn read_snapshot(&self) -> CubeStore { self.inner.read().unwrap().clone() } /// Test/bench only: take the write side of the inner `RwLock` directly. /// Gated behind the `bench` feature so it never ships in the daemon path. #[cfg(feature = "bench")] pub fn inner_write(&self) -> std::sync::RwLockWriteGuard<'_, CubeStore> { self.inner.write().unwrap() } // ---- Writes (write-lock the inner store + log to WAL) ---- /// Raw backend put, durability-logged. pub fn put_raw(&self, key: Czyx, value: Vec) { let v = { let mut g = self.inner.write().unwrap(); g.put_raw(key, value); g.get_raw(&key).unwrap_or_default() }; self.wal.append(WalOp::Put, key, v); } /// Raw backend delete, durability-logged. pub fn delete_raw(&self, key: &Czyx) { self.inner.write().unwrap().delete_raw(key); self.wal.append(WalOp::Delete, *key, Vec::new()); } /// Store `header` + `body` at `label`, durability-logged. pub fn put_record(&self, key: Czyx, header: &CubeHeader, body: &[u8]) { let v = { let mut g = self.inner.write().unwrap(); g.put_record(key, header, body); g.get_raw(&key).unwrap_or_default() }; self.wal.append(WalOp::Put, key, v); } /// Associate `src -> dst` (PDF Package 2 link), durability-logged. pub fn associate(&self, src: Czyx, dst: Czyx) -> bool { let ok = self.inner.write().unwrap().associate(src, dst); if ok { if let Some(v) = self.get_raw(&src) { self.wal.append(WalOp::Put, src, v); } } ok } /// Run a closure with exclusive (write) access to the inner store. Used by /// callers that mutate through the `CubeStore` API directly (e.g. /// `store_code_cell`, `CubeEnv::put_encrypted`). After such a mutation, /// call [`log_put`] with the resulting value to record it in the WAL. pub fn with_mut(&self, f: impl FnOnce(&mut CubeStore) -> R) -> R { let mut g = self.inner.write().unwrap(); f(&mut g) } /// Durability-log a put whose bytes were written via [`with_mut`]. The /// append is buffered; the group-commit thread fsyncs it within /// `wal_fsync_ms`. Callers needing an immediate guarantee use `commit_txn`. pub fn log_put(&self, key: Czyx, value: Vec) { let _seq = self.wal.append(WalOp::Put, key, value); } /// Apply a transaction batch atomically: under ONE store write lock, apply /// every put/delete, then append a SINGLE `WalOp::Txn` WAL entry so the /// whole batch is durable as one unit and replays idempotently. /// /// The WAL is fsync'd before returning so a `COMMIT` is durable the moment /// the caller gets control back (not merely "eventually" via the group /// thread). This is what makes `commit` a real transaction boundary. /// /// The fsync is COALESCED via [`Wal::sync_upto`]: concurrent committers /// share one fsync rather than each paying for their own, so the durability /// guarantee is unchanged while N writers no longer serialize behind N /// disk syncs. pub fn commit_txn(&self, batch: &[TxnEntry]) { let mut g = self.inner.write().unwrap(); for e in batch { match &e.value { Some(v) => g.put_raw(e.coord, v.clone()), None => g.delete_raw(&e.coord), } } drop(g); // release the write lock before touching the WAL let seq = self.wal.append_txn(batch); self.wal.sync_upto(seq); // coalesced fsync: COMMIT == durable } /// Store a code cell at `path` (the path->code bridge), durability-logged. /// `owner` (when set) is stamped on the record's `owner_local_user` field /// so owner enforcement (Task 6) can later reject cross-owner overwrites. pub fn put_code_cell( &self, path: &str, kind: Kind, name: &str, links: &[Czyx], code: &[Op], owner: Option<&str>, descriptor: Option, ) -> Result { let coord = self.with_mut(|store| { crate::store_code_cell(store, path, kind, name, links, code, owner, descriptor) })?; if let Some(v) = self.get_raw(&coord) { self.log_put(coord, v); } Ok(coord) } /// Force a durable checkpoint now (also flushes pending WAL). pub fn checkpoint(&self) { self.wal.flush_pending(); if self.db_path.as_os_str().is_empty() { return; } checkpoint_store(&self.inner, &self.db_path, &self.cp_seq_path, &self.wal); } /// Highest WAL sequence number known to be fsync'd to disk. Useful for /// telemetry — "how much is actually durable right now". pub fn wal_durable_seq(&self) -> u64 { self.wal.committed_seq.load(Ordering::SeqCst) } /// Stop background threads and perform a final checkpoint. Safe to call /// more than once. pub fn shutdown(&self) { self.stop.store(true, Ordering::SeqCst); self.wal.set_stop(); self.checkpoint(); eprintln!( "cube-store: shutdown; durable WAL seq up to {}", self.wal_durable_seq() ); if let Some(t) = self.flush_thread.lock().unwrap().take() { let _ = t.join(); } // Group thread is best-effort joined at Drop if still running. } } impl Drop for ConcurrentStore { fn drop(&mut self) { self.stop.store(true, Ordering::SeqCst); self.wal.set_stop(); if let Some(t) = self.flush_thread.lock().unwrap().take() { let _ = t.join(); } if let Some(t) = self.wal.group_thread.lock().unwrap().take() { let _ = t.join(); } } } /// Persistent checkpoint state lives in three files: /// * `db_path` — the full base snapshot (NDJSON array). Written rarely. /// * `db_path.delta` — incremental deltas since the base (one `{c,z,y,x,v}` object per modified coordinate, no seq). Written on every checkpoint. /// * `db_path.seq` — the WAL sequence number already folded into the base (so WAL replay starts after it). /// * `wal_path` — the WAL (transactions newer than `db_path.seq`). /// /// On startup we load the base, apply the delta, then replay WAL entries newer than `db_path.seq`. The `cp_seq_wal` field tracks that boundary so replay knows where to begin. This gives tiny per-checkpoint disk writes (only changed coords) while keeping a full snapshot for fast cold load. fn checkpoint_store( inner: &Arc>>, db_path: &Path, cp_seq_path: &Path, wal: &Arc, ) { wal.flush_pending(); // Collect only the coordinates that changed since the last base/checkpoint. let delta = incremental(inner, wal); let delta_path = db_path.with_extension("delta"); // If we have no delta at all (nothing changed), just re-stamp the WAL // boundary and leave the base untouched. if delta.is_empty() { persist_seq(cp_seq_path, wal.seq.load(Ordering::SeqCst)); wal.set_cp_seq_wal(wal.seq.load(Ordering::SeqCst)); return; } // Append the delta. If the delta file is getting large, fold everything // back into a fresh full base snapshot instead (compaction). If there is // no base snapshot yet at all, we must write the full base now (the delta // alone is not a complete store). let delta_len = fs::read(&delta_path).map(|b| b.len() as u64).unwrap_or(0); let base_exists = db_path.exists(); if !base_exists || delta_len + delta.len() as u64 > DELTA_COMPACT_BYTES { fold_delta_into_base(inner, db_path, cp_seq_path, &delta_path, wal); } else if let Some(parent) = delta_path.parent() { let _ = fs::create_dir_all(parent); if let Ok(mut f) = OpenOptions::new() .create(true) .append(true) .open(&delta_path) { let mut buf = String::new(); for e in &delta { buf.push_str(&format!( "{{\"c\":{},\"z\":{},\"y\":{},\"x\":{},\"v\":\"{}\"}}\n", e.0.c, e.0.z, e.0.y, e.0.x, to_hex(&e.1) )); } if f.write_all(buf.as_bytes()).is_ok() && f.flush().is_ok() && f.sync_all().is_ok() { // Boundary must be the highest *durable* (fsync'd) WAL seq, // NOT `wal.seq()` (which is the next-to-assign counter and // sits one past the last entry). Persisting the next-to- // assign value made `replay_after` skip every still-valid // WAL entry on restart — i.e. silent data loss of any // record written since the previous checkpoint. let durable = wal.committed_seq.load(Ordering::SeqCst); persist_seq(cp_seq_path, durable); wal.set_cp_seq_wal(durable); wal.set_base_seq(durable); } } } } /// Snapshot the current delta (only coordinates modified since the last base /// snapshot) as a list of `(coord, raw_value)` pairs. fn incremental( inner: &Arc>>, wal: &Arc, ) -> Vec<(Czyx, Vec)> { let g = inner.read().unwrap(); let base = wal.base_seq.load(Ordering::SeqCst); let now = wal.seq.load(Ordering::SeqCst); if now <= base { return Vec::new(); } let mut out = Vec::new(); for k in g.keys() { // Only include coords whose last write is newer than the base // snapshot. We approximate "changed since base" by tracking the set // of coords written after `base` in the WAL (the authoritative changelog). if wal.coord_modified_since(&k, base) { if let Some(v) = g.get_raw(&k) { out.push((k, v)); } } } out } /// Fold the current delta + live store into a fresh full base snapshot and /// truncate the delta. Used when the delta has grown too large. fn fold_delta_into_base( inner: &Arc>>, db_path: &Path, cp_seq_path: &Path, delta_path: &Path, wal: &Arc, ) { let snap = inner.read().unwrap().clone(); let json = persist::dump_store(&snap); if let Some(parent) = db_path.parent() { let _ = fs::create_dir_all(parent); } let tmp = db_path.with_extension("tmp"); if let Ok(mut f) = File::create(&tmp) { if f.write_all(json.as_bytes()).is_ok() && f.flush().is_ok() && f.sync_all().is_ok() { let _ = fs::rename(&tmp, db_path); } } // Delta is now fully represented by the base; truncate it. let _ = fs::write(delta_path, b""); // Boundary = highest *durable* WAL seq (fsync'd), not `wal.seq()` (the // next-to-assign counter, which sits one past the last entry and would // make `replay_after` skip still-valid entries on restart). let durable = wal.committed_seq.load(Ordering::SeqCst); persist_seq(cp_seq_path, durable); wal.set_cp_seq_wal(durable); wal.set_base_seq(durable); } /// Read `db_path` (full base) then apply the delta file; returns the /// reconstructed store. If either is missing it is simply skipped. fn load_base_plus_delta(db_path: &Path, delta_path: &Path) -> CubeStore { let mut store = CubeStore::new(HashBackend::new()); if db_path.exists() { if let Ok(json) = fs::read_to_string(db_path) { if let Err(e) = persist::load_into_store(&mut store, &json) { eprintln!("cube-store: base load failed ({e}); relying on WAL/delta"); } } } if delta_path.exists() { if let Ok(text) = fs::read_to_string(delta_path) { for line in text.lines() { let line = line.trim(); if line.is_empty() || !line.starts_with('{') { continue; } let c = field_u8(line, "c").unwrap_or(0); let z = field_u8(line, "z").unwrap_or(0); let y = field_u8(line, "y").unwrap_or(0); let x = field_u8(line, "x").unwrap_or(0); let v = field_str(line, "v"); if let Ok(value) = from_hex(v) { store.put_raw(Czyx::new(c, z, y, x), value); } } } } store } fn persist_seq(path: &Path, seq: u64) { if let Some(p) = path.parent() { let _ = fs::create_dir_all(p); } if let Ok(mut f) = File::create(path) { let _ = f.write_all(seq.to_string().as_bytes()); let _ = f.sync_all(); } } /// Append a recovery event to `path` (NDJSON, fsync'd). Called whenever we /// had to fall back to the WAL backup at startup. fn append_recovery_log(path: &Path, applied: u64) { if let Some(p) = path.parent() { let _ = fs::create_dir_all(p); } if let Ok(mut f) = OpenOptions::new().create(true).append(true).open(path) { let ts = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map(|d| d.as_secs()) .unwrap_or(0); let line = format!( "{{\"ts\":{ts},\"event\":\"wal_recovery\",\"applied\":{applied},\"note\":\"recovered {applied} transactions from the JSON WAL backup because they had not yet been checkpointed into the database\"}}\n" ); let _ = f.write_all(line.as_bytes()); let _ = f.sync_all(); } } // ---- WAL serialization (dependency-free) ---- fn encode_wal(e: &WalEntry) -> String { match e.op { WalOp::Put => format!( "{{\"seq\":{},\"op\":\"put\",\"c\":{},\"z\":{},\"y\":{},\"x\":{},\"v\":\"{}\"}}", e.seq, e.coord.c, e.coord.z, e.coord.y, e.coord.x, to_hex(&e.value) ), WalOp::Delete => format!( "{{\"seq\":{},\"op\":\"del\",\"c\":{},\"z\":{},\"y\":{},\"x\":{}}}", e.seq, e.coord.c, e.coord.z, e.coord.y, e.coord.x ), WalOp::Txn => format!( "{{\"seq\":{},\"op\":\"txn\",\"batch\":\"{}\"}}", e.seq, to_hex(&pack_txn(&e.batch)) ), } } /// Serialize a txn batch into a flat byte buffer (then hex'd for the WAL): /// for each entry: 1 flag byte (1=put, 0=del) + c,z,y,x (4 bytes) + /// if put: 4-byte BE length + value bytes. fn pack_txn(batch: &[TxnEntry]) -> Vec { let mut out = Vec::new(); for e in batch { match &e.value { Some(v) => { out.push(1); out.extend_from_slice(&[e.coord.c, e.coord.z, e.coord.y, e.coord.x]); out.extend_from_slice(&(v.len() as u32).to_be_bytes()); out.extend_from_slice(v); } None => { out.push(0); out.extend_from_slice(&[e.coord.c, e.coord.z, e.coord.y, e.coord.x]); } } } out } /// Reverse of [`pack_txn`]. fn unpack_txn(mut buf: &[u8]) -> Result, String> { let mut out = Vec::new(); while !buf.is_empty() { let flag = buf[0]; buf = &buf[1..]; if buf.len() < 4 { return Err("txn entry truncated (coord)".into()); } let coord = Czyx::new(buf[0], buf[1], buf[2], buf[3]); buf = &buf[4..]; if flag == 1 { if buf.len() < 4 { return Err("txn entry truncated (len)".into()); } let len = u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]) as usize; buf = &buf[4..]; if buf.len() < len { return Err("txn entry truncated (value)".into()); } let value = buf[..len].to_vec(); buf = &buf[len..]; out.push(TxnEntry { coord, value: Some(value), }); } else { out.push(TxnEntry { coord, value: None }); } } Ok(out) } fn decode_wal(line: &str) -> Option { let line = line.trim(); if line.is_empty() || !line.starts_with('{') { return None; } let seq = field_u64(line, "seq").ok()?; let op_s = field_str(line, "op"); // The txn encoder emits `{"seq","op":"txn","batch":...}` with NO c/z/y/x // (the per-entry coords live inside `batch`). Handle it before the // c/z/y/x extraction below, which would otherwise fail and drop the entry. if op_s == "txn" { let raw = field_str(line, "batch"); let bytes = from_hex(raw).ok()?; let batch = unpack_txn(&bytes).ok()?; return Some(WalEntry { seq, op: WalOp::Txn, coord: Czyx::new(0, 0, 0, 0), // unused; replay reads batch coords value: Vec::new(), batch, }); } let c = field_u8(line, "c").ok()?; let z = field_u8(line, "z").ok()?; let y = field_u8(line, "y").ok()?; let x = field_u8(line, "x").ok()?; let coord = Czyx::new(c, z, y, x); match op_s { "put" => { let v = field_str(line, "v"); let value = from_hex(v).ok()?; Some(WalEntry { seq, op: WalOp::Put, coord, value, batch: Vec::new(), }) } "del" => Some(WalEntry { seq, op: WalOp::Delete, coord, value: Vec::new(), batch: Vec::new(), }), _ => None, } } fn field_u64(obj: &str, key: &str) -> Result { let pat = format!("\"{key}\":"); let pos = obj.find(&pat).ok_or_else(|| format!("missing {key}"))?; let after = &obj[pos + pat.len()..]; let end = after.find([',', '}', ' ']).unwrap_or(after.len()); after[..end] .trim() .parse::() .map_err(|e| format!("bad {key}: {e}")) } fn field_u8(obj: &str, key: &str) -> Result { field_u64(obj, key).and_then(|v| u8::try_from(v).map_err(|_| format!("{key} out of u8 range"))) } fn field_str<'a>(obj: &'a str, key: &str) -> &'a str { let pat = format!("\"{key}\":"); let pos = match obj.find(&pat) { Some(p) => p, None => return "", }; let after = obj[pos + pat.len()..].trim_start(); if let Some(rest) = after.strip_prefix('"') { if let Some(end) = rest.find('"') { return &rest[..end]; } } "" } fn to_hex(bytes: &[u8]) -> String { let mut s = String::with_capacity(bytes.len() * 2); for b in bytes { s.push_str(&format!("{b:02x}")); } s } fn from_hex(s: &str) -> Result, String> { if !s.len().is_multiple_of(2) { return Err("odd-length hex".to_string()); } 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) .ok_or_else(|| "bad hex digit".to_string())?; let lo = (bytes[i + 1] as char) .to_digit(16) .ok_or_else(|| "bad hex digit".to_string())?; out.push((hi * 16 + lo) as u8); i += 2; } Ok(out) } #[cfg(test)] mod tests { use super::*; use std::sync::atomic::AtomicUsize; fn tmp(name: &str) -> String { format!( "{}/cubetest-{}-{}", std::env::temp_dir().display(), std::process::id(), name ) } fn cleanup(db: &str) { let _ = fs::remove_file(db); let _ = fs::remove_file(format!("{db}.wal")); let _ = fs::remove_file(format!("{db}.recovery.ndjson")); let _ = fs::remove_file(format!("{db}.seq")); let _ = fs::remove_file(PathBuf::from(db).with_extension("delta")); } #[test] fn memory_store_roundtrip() { let s = ConcurrentStore::memory(); s.put_raw(Czyx::new(1, 2, 3, 4), vec![9]); assert_eq!(s.get_raw(&Czyx::new(1, 2, 3, 4)), Some(vec![9])); let sn = s.read_snapshot(); assert_eq!(sn.get_raw(&Czyx::new(1, 2, 3, 4)), Some(vec![9])); } #[test] fn durable_checkpoint_and_replay() { let db = tmp("db.json"); let wal = format!("{db}.wal"); let rec = format!("{db}.recovery.ndjson"); cleanup(&db); { let s = ConcurrentStore::open( &db, &wal, &rec, DurabilityConfig { checkpoint_ms: 50, wal_fsync_ms: 10, }, ) .unwrap(); s.put_record(Czyx::new(1, 1, 1, 1), &CubeHeader::new(), b"1"); s.put_record(Czyx::new(2, 2, 2, 2), &CubeHeader::new(), b"2"); s.checkpoint(); } { let s = ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(); assert!( s.keys().contains(&Czyx::new(1, 1, 1, 1)), "coord 1 lost across checkpoint" ); assert!( s.keys().contains(&Czyx::new(2, 2, 2, 2)), "coord 2 lost across checkpoint" ); assert_eq!( s.get_record(&Czyx::new(1, 1, 1, 1)).unwrap().1, b"1".to_vec() ); } cleanup(&db); } #[test] fn wal_recovery_after_crash() { let db = tmp("db2.json"); let wal = format!("{db}.wal"); let rec = format!("{db}.recovery.ndjson"); cleanup(&db); { // Long checkpoint interval so the post-checkpoint write is NOT // checkpointed before we "crash". let s = ConcurrentStore::open( &db, &wal, &rec, DurabilityConfig { checkpoint_ms: 10_000, wal_fsync_ms: 10, }, ) .unwrap(); s.put_record(Czyx::new(1, 1, 1, 1), &CubeHeader::new(), b"1"); s.put_record(Czyx::new(2, 2, 2, 2), &CubeHeader::new(), b"2"); s.checkpoint(); // Simulate a write that is fsync'd to the WAL but never // checkpointed, then a crash (drop without another checkpoint). s.put_record(Czyx::new(3, 3, 3, 3), &CubeHeader::new(), b"3"); s.wal.flush_pending(); } // Reopen: 1,2 from the checkpoint; 3 recovered from the WAL. { let s = ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(); assert!(s.keys().contains(&Czyx::new(1, 1, 1, 1)), "coord 1 lost"); assert!(s.keys().contains(&Czyx::new(2, 2, 2, 2)), "coord 2 lost"); assert!( s.keys().contains(&Czyx::new(3, 3, 3, 3)), "coord 3 not recovered from WAL" ); } let rec_text = fs::read_to_string(&rec).unwrap_or_default(); assert!( rec_text.contains("wal_recovery"), "expected a recovery event to be logged" ); cleanup(&db); } /// Regression: `committed_seq` must never exceed what was actually fsync'd. /// /// The old `flush_pending` stamped `committed_seq = seq.load()` — the LIVE /// counter. `Wal::append` bumps `seq` (fetch_add) and pushes its bytes into /// `pending` under a *separate* lock, so there is a real window where `seq` /// has advanced but the bytes are not yet buffered. A flush landing in that /// window would mark the absent seq durable. /// /// This test drives that window deterministically rather than relying on /// thread timing: it advances `seq` exactly as `append`'s `fetch_add` does, /// WITHOUT buffering bytes, then flushes. `wal_durable_seq()` must not /// claim the phantom sequence. #[test] fn durable_seq_never_over_reports() { let db = tmp("durseq"); let wal_p = format!("{db}.wal"); let rec = format!("{db}.rec"); cleanup(&db); let s = ConcurrentStore::open(&db, &wal_p, &rec, DurabilityConfig::default()).unwrap(); // One real, fully-buffered write, then flush it. s.put_raw(Czyx::new(9, 0, 0, 1), vec![1]); s.wal.flush_pending(); let after_real = s.wal_durable_seq(); // Buffer a real write (takes seq N and pushes its bytes)... s.put_raw(Czyx::new(9, 0, 0, 2), vec![2]); // ...then reproduce the window: a concurrent appender takes seq N+1 but // has not yet pushed its bytes. It must NOT be counted as durable. let phantom = s.wal.seq.fetch_add(1, Ordering::SeqCst) + 1; s.wal.flush_pending(); let claimed = s.wal_durable_seq(); let on_disk = fs::read_to_string(&wal_p) .unwrap_or_default() .lines() .filter_map(decode_wal) .map(|e| e.seq) .max() .unwrap_or(0); assert!( claimed <= on_disk, "durability over-reported: claimed durable seq {claimed} but only {on_disk} is on disk \ (phantom seq {phantom} was never written)" ); assert!( claimed > after_real, "flush failed to advance durability at all: {claimed} vs {after_real}" ); s.shutdown(); cleanup(&db); } /// Concurrent commits must COALESCE into shared fsyncs while EVERY commit /// still returns genuinely durable. This crashes the store without any /// checkpoint, so recovery is pure WAL replay: if the coalescing logic ever /// lets a committer return before its own bytes are fsync'd, a coord goes /// missing here. #[test] fn concurrent_commits_are_durable_and_coalesced() { let db = tmp("coalesce"); let wal = format!("{db}.wal"); let rec = format!("{db}.rec"); cleanup(&db); { let s = Arc::new( ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(), ); let mut hs = Vec::new(); for t in 0..8u32 { let s = s.clone(); hs.push(thread::spawn(move || { for i in 0..25u32 { let x = (t * 25 + i) as u8; s.commit_txn(&[TxnEntry { coord: Czyx::new(8, 0, 0, x), value: Some(vec![x]), }]); } })); } for h in hs { h.join().unwrap(); } // Crash WITHOUT checkpoint: every commit must still be recoverable // purely from the fsync'd WAL. } let s = ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(); let keys = s.keys(); for x in 0..200u8 { assert!( keys.contains(&Czyx::new(8, 0, 0, x)), "commit at x={x} was lost — coalesced fsync broke the durability guarantee" ); } s.shutdown(); cleanup(&db); } #[test] fn concurrent_writes() { let s = Arc::new(ConcurrentStore::memory()); let counter = Arc::new(AtomicUsize::new(0)); let mut hs = Vec::new(); for t in 0..8u32 { let s = s.clone(); let c = counter.clone(); hs.push(thread::spawn(move || { for i in 0..500u32 { let x = (t * 500 + i) as u8; s.put_raw(Czyx::new(1, 0, 0, x), vec![x]); c.fetch_add(1, Ordering::SeqCst); } })); } for h in hs { h.join().unwrap(); } assert_eq!(counter.load(Ordering::SeqCst), 4000); // All 4000 puts landed (last-writer-wins per coordinate); because the X // axis is a u8, the 4000 distinct input x-values wrap to 256 unique // coordinates, so the store holds 256 records — every put was observed // by the mutex, none lost to a race. assert_eq!(s.keys().len(), 256); } #[test] fn incremental_checkpoint_delta_model() { let db = tmp("db3.json"); let wal = format!("{db}.wal"); let rec = format!("{db}.recovery.ndjson"); let delta = PathBuf::from(&db) .with_extension("delta") .to_string_lossy() .into_owned(); cleanup(&db); // Phase 1: 10 records, then a checkpoint. The base snapshot is written // and the delta starts empty (everything is now in the base). { let s = ConcurrentStore::open( &db, &wal, &rec, DurabilityConfig { checkpoint_ms: 50, wal_fsync_ms: 10, }, ) .unwrap(); for i in 0..10u8 { s.put_record(Czyx::new(1, 1, 1, i), &CubeHeader::new(), &[i]); } s.checkpoint(); } let base_after_1 = fs::read(&db).map(|b| b.len()).unwrap_or(0); let delta_after_1 = fs::read(&delta).map(|b| b.len()).unwrap_or(0); assert!( base_after_1 > 0, "base snapshot should exist after first checkpoint" ); assert_eq!( delta_after_1, 0, "delta should be empty right after a full checkpoint" ); // Phase 2: change only 2 records. The next checkpoint must write only // those 2 into the delta (NOT rewrite the whole base), and a reopen // must still see all 10. { let s = ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(); s.put_record(Czyx::new(1, 1, 1, 0), &CubeHeader::new(), b"changed"); s.put_record(Czyx::new(1, 1, 1, 1), &CubeHeader::new(), b"changed"); s.checkpoint(); } let base_after_2 = fs::read(&db).map(|b| b.len()).unwrap_or(0); let delta_lines_2 = fs::read_to_string(&delta) .map(|t| t.lines().filter(|l| !l.trim().is_empty()).count()) .unwrap_or(0); // Base should NOT have grown by a full rewrite of 10 records; the delta // must carry exactly the 2 changed coords. assert!( base_after_2 <= base_after_1 + 64, "base should not be fully rewritten on a small change" ); assert_eq!( delta_lines_2, 2, "delta should contain only the 2 changed coords" ); // Reopen from base + delta and confirm all 10 are present with the // updated values for the 2 we changed. { let s = ConcurrentStore::open(&db, &wal, &rec, DurabilityConfig::default()).unwrap(); assert_eq!( s.keys().len(), 10, "all 10 records must survive base+delta reopen" ); assert_eq!( s.get_record(&Czyx::new(1, 1, 1, 0)).unwrap().1, b"changed".to_vec() ); assert_eq!( s.get_record(&Czyx::new(1, 1, 1, 1)).unwrap().1, b"changed".to_vec() ); assert_eq!(s.get_record(&Czyx::new(1, 1, 1, 9)).unwrap().1, vec![9u8]); } cleanup(&db); } #[test] fn query_doc_type_works() { let s = ConcurrentStore::memory(); let mut h = CubeHeader::new(); h.doc_type = Some("fn".into()); h.refresh_flags(); s.put_record(Czyx::new(1, 1, 1, 1), &h, b"a"); s.put_record(Czyx::new(1, 1, 1, 2), &CubeHeader::new(), b"b"); let fns = s.query_doc_type("fn"); assert_eq!(fns, vec![Czyx::new(1, 1, 1, 1)]); } // Task 4: reader/writer sharding. Many concurrent readers must make // progress WHILE a long writer holds the write side — i.e. reads take the // RwLock read side and are not serialized behind a single global mutex. // We assert that reader threads complete their reads during the writer's // hold window rather than after it. #[test] fn concurrent_reads_dont_block_on_writer() { let s = Arc::new(ConcurrentStore::memory()); let coord = Czyx::new(5, 1, 1, 1); s.put_raw(coord, vec![42]); // A writer that holds the write lock for a while. let writer = { let s = s.clone(); thread::spawn(move || { // Take the write side and hold it long enough that any // serialized reader model would force all readers to wait. let mut g = s.inner.write().unwrap(); thread::sleep(Duration::from_millis(200)); g.put_raw(coord, vec![7]); }) }; // Spawn readers; they should acquire the read side concurrently // (the RwLock admits many readers at once) and finish well before // the writer releases. let mut readers = Vec::new(); let start = std::time::Instant::now(); for _ in 0..8 { let s = s.clone(); readers.push(thread::spawn(move || { for _ in 0..200 { // read lock — must not block on other readers let _ = s.get_raw(&coord); } })); } for r in readers { r.join().unwrap(); } writer.join().unwrap(); // The readers did 8*200=1600 reads. Under a single global mutex behind // a 200ms writer hold they could not all finish before the writer, // because the writer would serialize every other access. They did // finish during the window (the read side is shared), so total time is // well under what a fully-serialized path would cost. assert!( start.elapsed() < Duration::from_secs(2), "reads appear serialized behind the writer (RwLock sharding broken?)" ); assert_eq!(s.get_raw(&coord), Some(vec![7])); } }