feat(cubesys): Task 4 reader/writer sharding (Mutex -> RwLock)
ConcurrentStore.inner is now Arc<RwLock<CubeStore>>: all read paths take the read side, all mutations + checkpoint take the write side. Readers no longer exclude each other and overlap an active writer (verified by concurrent_reads_dont_block_on_writer + cube-bench Task 4 section). WAL, checkpoint, and coordinate encoding are untouched, so durability/replay is unchanged (.check green). Honest finding recorded in docs/task4-reader-writer-sharding.md: on this 8-core host std RwLock removes reader-vs-reader exclusion (correct) but shows no wall-clock speedup for short reads (cache-line bounce on one shared lock). Real read-throughput scaling would need sharded/lock-free storage, left as a follow-up decision rather than invented.
This commit is contained in:
+101
-28
@@ -25,7 +25,7 @@ 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};
|
||||
use std::sync::{Arc, Mutex, RwLock};
|
||||
use std::thread::{self, JoinHandle};
|
||||
use std::time::Duration;
|
||||
|
||||
@@ -285,8 +285,15 @@ fn spawn_group(wal: Arc<Wal>) -> JoinHandle<()> {
|
||||
}
|
||||
|
||||
/// 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<Mutex<CubeStore<HashBackend>>>,
|
||||
inner: Arc<RwLock<CubeStore<HashBackend>>>,
|
||||
wal: Arc<Wal>,
|
||||
db_path: PathBuf,
|
||||
cp_seq_path: PathBuf,
|
||||
@@ -300,7 +307,7 @@ impl ConcurrentStore {
|
||||
/// No WAL, no checkpoint thread.
|
||||
pub fn memory() -> Self {
|
||||
ConcurrentStore {
|
||||
inner: Arc::new(Mutex::new(CubeStore::new(HashBackend::new()))),
|
||||
inner: Arc::new(RwLock::new(CubeStore::new(HashBackend::new()))),
|
||||
wal: Wal::memory(),
|
||||
db_path: PathBuf::new(),
|
||||
cp_seq_path: PathBuf::new(),
|
||||
@@ -355,7 +362,7 @@ impl ConcurrentStore {
|
||||
append_recovery_log(&rec_p, applied);
|
||||
}
|
||||
|
||||
let inner = Arc::new(Mutex::new(store));
|
||||
let inner = Arc::new(RwLock::new(store));
|
||||
let cs = ConcurrentStore {
|
||||
inner,
|
||||
wal,
|
||||
@@ -386,37 +393,37 @@ impl ConcurrentStore {
|
||||
Ok(cs)
|
||||
}
|
||||
|
||||
// ---- Reads (lock the inner store) ----
|
||||
// ---- Reads (read-lock the inner store; many readers run in parallel) ----
|
||||
|
||||
/// Raw backend get.
|
||||
pub fn get_raw(&self, key: &Czyx) -> Option<Vec<u8>> {
|
||||
self.inner.lock().unwrap().get_raw(key)
|
||||
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<u8>)> {
|
||||
self.inner.lock().unwrap().get_record(key)
|
||||
self.inner.read().unwrap().get_record(key)
|
||||
}
|
||||
|
||||
/// Every coordinate present.
|
||||
pub fn keys(&self) -> Vec<Czyx> {
|
||||
self.inner.lock().unwrap().keys()
|
||||
self.inner.read().unwrap().keys()
|
||||
}
|
||||
|
||||
/// Coordinates under a `C`/`Z`/`Y` prefix.
|
||||
pub fn scan_prefix(&self, c: u8, z: Option<u8>, y: Option<u8>) -> Vec<Czyx> {
|
||||
self.inner.lock().unwrap().scan_prefix(c, z, y)
|
||||
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<Czyx> {
|
||||
self.inner.lock().unwrap().linked_to(target)
|
||||
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<Czyx> {
|
||||
let g = self.inner.lock().unwrap();
|
||||
let g = self.inner.read().unwrap();
|
||||
let mut out = Vec::new();
|
||||
for k in g.keys() {
|
||||
if let Some((h, _)) = g.get_record(&k) {
|
||||
@@ -430,17 +437,26 @@ impl ConcurrentStore {
|
||||
}
|
||||
|
||||
/// A consistent point-in-time snapshot of the whole store. Used by the VM
|
||||
/// and cubefs, which take a `CubeStore` by value.
|
||||
/// 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<HashBackend> {
|
||||
self.inner.lock().unwrap().clone()
|
||||
self.inner.read().unwrap().clone()
|
||||
}
|
||||
|
||||
// ---- Writes (lock the inner store + log to WAL) ----
|
||||
/// 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<HashBackend>> {
|
||||
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<u8>) {
|
||||
let v = {
|
||||
let mut g = self.inner.lock().unwrap();
|
||||
let mut g = self.inner.write().unwrap();
|
||||
g.put_raw(key, value);
|
||||
g.get_raw(&key).unwrap_or_default()
|
||||
};
|
||||
@@ -449,14 +465,14 @@ impl ConcurrentStore {
|
||||
|
||||
/// Raw backend delete, durability-logged.
|
||||
pub fn delete_raw(&self, key: &Czyx) {
|
||||
self.inner.lock().unwrap().delete_raw(key);
|
||||
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.lock().unwrap();
|
||||
let mut g = self.inner.write().unwrap();
|
||||
g.put_record(key, header, body);
|
||||
g.get_raw(&key).unwrap_or_default()
|
||||
};
|
||||
@@ -465,7 +481,7 @@ impl ConcurrentStore {
|
||||
|
||||
/// Associate `src -> dst` (PDF Package 2 link), durability-logged.
|
||||
pub fn associate(&self, src: Czyx, dst: Czyx) -> bool {
|
||||
let ok = self.inner.lock().unwrap().associate(src, dst);
|
||||
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);
|
||||
@@ -474,12 +490,12 @@ impl ConcurrentStore {
|
||||
ok
|
||||
}
|
||||
|
||||
/// Run a closure with exclusive 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.
|
||||
/// 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<R>(&self, f: impl FnOnce(&mut CubeStore<HashBackend>) -> R) -> R {
|
||||
let mut g = self.inner.lock().unwrap();
|
||||
let mut g = self.inner.write().unwrap();
|
||||
f(&mut g)
|
||||
}
|
||||
|
||||
@@ -558,7 +574,7 @@ impl Drop for ConcurrentStore {
|
||||
///
|
||||
/// 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<Mutex<CubeStore<HashBackend>>>,
|
||||
inner: &Arc<RwLock<CubeStore<HashBackend>>>,
|
||||
db_path: &Path,
|
||||
cp_seq_path: &Path,
|
||||
wal: &Arc<Wal>,
|
||||
@@ -615,8 +631,11 @@ fn checkpoint_store(
|
||||
|
||||
/// Snapshot the current delta (only coordinates modified since the last base
|
||||
/// snapshot) as a list of `(coord, raw_value)` pairs.
|
||||
fn incremental(inner: &Arc<Mutex<CubeStore<HashBackend>>>, wal: &Arc<Wal>) -> Vec<(Czyx, Vec<u8>)> {
|
||||
let g = inner.lock().unwrap();
|
||||
fn incremental(
|
||||
inner: &Arc<RwLock<CubeStore<HashBackend>>>,
|
||||
wal: &Arc<Wal>,
|
||||
) -> Vec<(Czyx, Vec<u8>)> {
|
||||
let g = inner.read().unwrap();
|
||||
let base = wal.base_seq.load(Ordering::SeqCst);
|
||||
let now = wal.seq.load(Ordering::SeqCst);
|
||||
if now <= base {
|
||||
@@ -639,13 +658,13 @@ fn incremental(inner: &Arc<Mutex<CubeStore<HashBackend>>>, wal: &Arc<Wal>) -> Ve
|
||||
/// 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<Mutex<CubeStore<HashBackend>>>,
|
||||
inner: &Arc<RwLock<CubeStore<HashBackend>>>,
|
||||
db_path: &Path,
|
||||
cp_seq_path: &Path,
|
||||
delta_path: &Path,
|
||||
wal: &Arc<Wal>,
|
||||
) {
|
||||
let snap = inner.lock().unwrap().clone();
|
||||
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);
|
||||
@@ -1077,4 +1096,58 @@ mod tests {
|
||||
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]));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user