#![forbid(unsafe_code)]
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use crate::optimizer::foreground::ForegroundPolicy;
use crate::optimizer::policy::OptimizeOptions;
use crate::store::transaction::CrashHooks;
use crate::store::{AttrUpdate, NewEntry, Store, StoreConfig};
fn create_store(dir: &tempfile::TempDir) -> Store {
let cfg = StoreConfig {
segment_size: 1024 * 1024,
..Default::default()
};
Store::create(dir.path(), &cfg, [0x44; 16]).unwrap()
}
fn file_bytes(i: u64) -> Vec<u8> {
let body = blake3::hash(&i.to_le_bytes());
let mut out = Vec::with_capacity(599);
out.extend_from_slice(b"# generated config\nhost = node-");
out.extend_from_slice(i.to_string().as_bytes());
out.extend_from_slice(b"\nport = 8000\nuser = svc\npayload = ");
out.extend_from_slice(&body.as_bytes()[..32]);
out.resize(599, 0x5A);
out
}
fn run_court(
dir: &tempfile::TempDir,
files: u64,
writers: usize,
setattr_iters: usize,
checkpoint_iters: usize,
) -> usize {
let store = Arc::new(create_store(dir));
let root = store
.create_entry(
1,
b"root",
NewEntry::dir(0o755, 1000, 1000),
&CrashHooks::none(),
)
.unwrap();
let inos: Vec<u64> = (0..files)
.map(|i| {
store
.create_entry(
root,
format!("f{i:04}.cfg").as_bytes(),
NewEntry::file(0o644, 1000, 1000),
&CrashHooks::none(),
)
.unwrap()
})
.collect();
let refs: Vec<Vec<u8>> = (0..files).map(file_bytes).collect();
let stop = Arc::new(AtomicBool::new(false));
let counter = AtomicUsize::new(0);
let store_w = Arc::clone(&store);
let inos_a = Arc::new(inos);
let refs_a = Arc::new(refs);
std::thread::scope(|s| {
for _ in 0..writers {
let store = Arc::clone(&store_w);
let inos = Arc::clone(&inos_a);
let refs = Arc::clone(&refs_a);
let counter = &counter;
s.spawn(move || {
let opts = OptimizeOptions::default();
let fg = ForegroundPolicy::full();
loop {
let i = counter.fetch_add(1, Ordering::Relaxed);
if i >= files as usize {
break;
}
store
.epoch_write(inos[i], 0, &refs[i], opts, fg, &CrashHooks::none())
.unwrap_or_else(|e| panic!("file {i}: epoch_write failed: {e}"));
}
});
}
let store = Arc::clone(&store_w);
let inos = Arc::clone(&inos_a);
let stop_s = Arc::clone(&stop);
s.spawn(move || {
let mut i = 0usize;
while !stop_s.load(Ordering::Relaxed) && i < setattr_iters {
let ino = inos[i % inos.len()];
let mode = Some(0o644 + (i % 2) as u32);
store
.epoch_setattr(
ino,
&AttrUpdate {
mode,
..Default::default()
},
&CrashHooks::none(),
)
.unwrap_or_else(|e| panic!("setattr: {e}"));
i += 1;
}
});
let store = Arc::clone(&store_w);
let stop_c = Arc::clone(&stop);
s.spawn(move || {
let mut i = 0usize;
while !stop_c.load(Ordering::Relaxed) && i < checkpoint_iters {
store
.epoch_checkpoint(&CrashHooks::none())
.unwrap_or_else(|e| panic!("checkpoint: {e}"));
i += 1;
}
});
while counter.load(Ordering::Relaxed) < files as usize {
std::thread::yield_now();
}
stop.store(true, Ordering::Relaxed);
});
drop(store);
drop(store_w);
let reopened = Store::open(dir.path(), &StoreConfig::default()).unwrap();
let mut bad = 0;
for i in 0..files as usize {
let got = reopened.read_file(inos_a[i], 0, 599).unwrap();
if got != refs_a[i] {
bad += 1;
}
}
bad
}
#[test]
fn concurrent_writes_with_checkpointing_stay_byte_exact() {
let dir = tempfile::TempDir::new().unwrap();
let bad = run_court(&dir, 64, 4, 256, 64);
assert_eq!(
bad, 0,
"concurrent writes + setattr + checkpoint lost extents (read-back mismatches: {bad})"
);
}
#[test]
fn concurrent_writes_without_setattr_stay_byte_exact() {
let dir = tempfile::TempDir::new().unwrap();
let bad = run_court(&dir, 64, 4, 0, 64);
assert_eq!(
bad, 0,
"concurrent writes + checkpoint lost extents (read-back mismatches: {bad})"
);
}
#[test]
fn write_then_rename_in_one_epoch_survives_reopen() {
let dir = tempfile::TempDir::new().unwrap();
let cfg = StoreConfig {
segment_size: 1024 * 1024,
..Default::default()
};
let store = Store::create(dir.path(), &cfg, [0x55; 16]).unwrap();
let hooks = CrashHooks::none();
let ino = store
.epoch_create(1, b"tmp", NewEntry::file(0o600, 1000, 1000), &hooks)
.unwrap();
store
.epoch_write(
ino,
0,
b"log-only blob",
OptimizeOptions::default(),
ForegroundPolicy::full(),
&hooks,
)
.unwrap();
store.epoch_rename(1, b"tmp", 1, b"final", &hooks).unwrap();
drop(store);
let reopened = Store::open(dir.path(), &StoreConfig::default()).unwrap();
let out = reopened.read_file(ino, 0, 64).unwrap();
assert_eq!(
out, b"log-only blob",
"rename replay clobbered the moved file's extent root"
);
let entry = reopened
.dir_lookup(1, b"final")
.unwrap()
.expect("final entry exists");
assert_eq!(entry.ino, ino);
}