#![forbid(unsafe_code)]
type DurableLedger = std::sync::Arc<std::sync::Mutex<Vec<(u64, u64, Vec<u8>)>>>;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Instant;
use crate::optimizer::policy::OptimizeOptions;
use crate::store::transaction::{CrashHooks, CrashPoint};
use crate::store::{NewEntry, Store, StoreConfig};
use tempfile::TempDir;
fn create_store(dir: &TempDir) -> Arc<Store> {
let cfg = StoreConfig {
segment_size: 128 * 1024 * 1024,
..Default::default()
};
Arc::new(Store::create(dir.path(), &cfg, [0x55; 16]).unwrap())
}
fn create_file(store: &Store, name: &str) -> u64 {
store
.create_entry(
store.current_root().root_dir_ino,
name.as_bytes(),
NewEntry::file(0o644, 1000, 1000),
&CrashHooks::none(),
)
.unwrap()
}
fn stream_for(writer: u64, cycle: u64) -> Vec<u8> {
let mut out = Vec::with_capacity(65536);
let mut state = 0x12b_c0deu64;
state ^= writer.wrapping_mul(0x9e37_79b9_7f4a_7c15);
state ^= cycle.wrapping_mul(0xbf58_476d_1ce4_e5b9);
for _ in 0..65536 {
state = state
.wrapping_mul(6364136223846793005)
.wrapping_add(1442695040888963407);
out.push((state >> 33) as u8);
}
out
}
fn run_one(point: CrashPoint) {
let dir = TempDir::new().unwrap();
let store = create_store(&dir);
let fg = store.foreground_policy();
let opts = OptimizeOptions::default();
let crashed = Arc::new(AtomicBool::new(false));
let durable: DurableLedger = Arc::new(Mutex::new(Vec::new()));
let mut inos = Vec::new();
for w in 0..8u64 {
inos.push(create_file(&store, &format!("f{w}")));
}
let crash_ino = create_file(&store, "crash");
std::thread::scope(|s| {
for w in 0..8u64 {
let store = Arc::clone(&store);
let crashed = Arc::clone(&crashed);
let durable = Arc::clone(&durable);
let inos = &inos;
s.spawn(move || {
let mut cycle = 0u64;
while !crashed.load(Ordering::Relaxed) && cycle < 12 {
let data = stream_for(w, cycle);
store
.epoch_write(
inos[w as usize],
cycle * 65536,
&data,
opts,
fg,
&CrashHooks::none(),
)
.unwrap();
match store.durability_barrier(&CrashHooks::none()) {
Ok(()) => {
durable
.lock()
.unwrap()
.push((inos[w as usize], cycle * 65536, data))
}
Err(_) => {
crashed.store(true, Ordering::Relaxed);
break;
}
}
cycle += 1;
}
});
}
let store = Arc::clone(&store);
let crashed = Arc::clone(&crashed);
let durable = Arc::clone(&durable);
s.spawn(move || {
let mut cycle = 0u64;
let t0 = Instant::now();
while !crashed.load(Ordering::Relaxed) {
let data = stream_for(999, cycle);
store
.epoch_write(
crash_ino,
cycle * 65536,
&data,
opts,
fg,
&CrashHooks::none(),
)
.unwrap();
match store.durability_barrier(&CrashHooks::crash_at(point)) {
Ok(()) => {
durable
.lock()
.unwrap()
.push((crash_ino, cycle * 65536, data));
}
Err(_) => {
crashed.store(true, Ordering::Relaxed);
break;
}
}
cycle += 1;
assert!(
t0.elapsed().as_secs() < 30,
"crash point {point:?} never fired"
);
}
});
});
assert!(
crashed.load(Ordering::Relaxed),
"crash point {point:?} must fire"
);
drop(store);
let cfg = StoreConfig {
segment_size: 128 * 1024 * 1024,
..Default::default()
};
let store2 = Store::open(dir.path(), &cfg).unwrap();
let durable = durable.lock().unwrap();
assert!(
!durable.is_empty(),
"crash point {point:?}: at least one fsync must have returned Ok before the crash"
);
let mut missing: Vec<(u64, u64)> = Vec::new();
for (ino, off, data) in durable.iter() {
let got = store2.read_file(*ino, *off, data.len() as u64).unwrap();
if &got != data {
missing.push((*ino, *off));
}
}
assert!(
missing.is_empty(),
"crash point {point:?}: {} RETURNED fsyncs missing after recovery: {missing:?} of {} recorded",
missing.len(),
durable.len()
);
drop(store2);
let report = crate::fsck::fsck(dir.path(), &crate::fsck::FsckOptions::default()).unwrap();
assert!(
report.is_clean(),
"crash point {point:?}: fsck must be clean after recovery ({})",
report.error_count()
);
}
#[test]
fn group_barrier_crash_at_every_stage_keeps_returned_fsyncs() {
for point in [
CrashPoint::AfterRecordAppend,
CrashPoint::AfterSegmentFdatasync,
CrashPoint::AfterSegmentDirFsync,
CrashPoint::AfterSuperblockWrite,
CrashPoint::AfterSuperblockFsync,
] {
run_one(point);
}
}