#![forbid(unsafe_code)]
use std::time::Instant;
use crate::store::inode::Inode;
use crate::store::transaction::CrashHooks;
use crate::store::{Store, StoreConfig};
fn env_or(key: &str, default: &str) -> String {
std::env::var(key).unwrap_or_else(|_| default.to_string())
}
fn thread_cpu_seconds() -> f64 {
let body = std::fs::read_to_string("/proc/self/stat").unwrap_or_default();
let Some(rest) = body.split_once(')') else {
return 0.0;
};
let fields: Vec<&str> = rest.1[1..].split_whitespace().collect();
let utime: f64 = fields.get(11).and_then(|v| v.parse().ok()).unwrap_or(0.0);
let stime: f64 = fields.get(12).and_then(|v| v.parse().ok()).unwrap_or(0.0);
(utime + stime) / 100.0 }
fn write_phase(store: &Store, ino: u64, mib: u64, fsync_every: u64) -> (f64, f64, f64, u64) {
let options = crate::optimizer::policy::OptimizeOptions::default();
let cpu_start = thread_cpu_seconds();
let total = mib * 1024 * 1024;
let chunk = 64 * 1024usize;
let batches = total / chunk as u64;
let wall = Instant::now();
let mut written = 0u64;
let mut barriers = 0u64;
let mut flushes = 0u64;
let mut batch: Vec<(u64, Vec<u8>)> = Vec::new();
for b in 0..batches {
let pat = (b / 16) % 4;
let mut data = Vec::with_capacity(chunk);
match pat {
0 => {
for i in 0..chunk {
data.push(b'a' + (i % 26) as u8);
}
}
1 => data.resize(chunk, 0),
2 => {
for i in 0..chunk {
data.push((i % 7) as u8);
}
}
_ => {
for i in 0..chunk {
data.push(((i * 31 + b as usize) % 251) as u8);
}
}
}
batch.push((written, data));
written += chunk as u64;
if batch.len() >= 32 {
store
.write_region_batch(ino, &batch, options)
.expect("write batch");
batch.clear();
flushes += 1;
if fsync_every > 0 && flushes.is_multiple_of(fsync_every) {
store
.durability_barrier(&CrashHooks::none())
.expect("barrier");
barriers += 1;
}
}
}
if !batch.is_empty() {
store
.write_region_batch(ino, &batch, options)
.expect("write tail");
}
if fsync_every == 0 {
store
.durability_barrier(&CrashHooks::none())
.expect("final barrier");
barriers += 1;
}
let wall_s = wall.elapsed().as_secs_f64();
let cpu_s = thread_cpu_seconds() - cpu_start;
(
total as f64 / wall_s / 1024.0 / 1024.0,
wall_s,
cpu_s,
barriers,
)
}
fn read_phase_seq(store: &Store, ino: u64, mib: u64) -> (f64, f64, f64, f64, f64) {
let chunk = 64 * 1024usize;
let total = mib * 1024 * 1024;
let cpu_start = thread_cpu_seconds();
let mut samples: Vec<u64> = Vec::new();
let wall = Instant::now();
let mut off = 0u64;
while off < total {
let t = Instant::now();
let data = store.read_file(ino, off, chunk as u64).expect("read");
samples.push(t.elapsed().as_micros() as u64);
assert_eq!(data.len(), chunk.min((total - off) as usize));
off += chunk as u64;
}
let wall_s = wall.elapsed().as_secs_f64();
let cpu_s = thread_cpu_seconds() - cpu_start;
let mbps = total as f64 / wall_s / 1024.0 / 1024.0;
let (p50, p95, p99) = percentiles(&mut samples);
(mbps, cpu_s, p50, p95, p99)
}
fn read_phase_random(store: &Store, ino: u64, mib: u64) -> (f64, f64, f64, f64, f64) {
let span = (mib * 1024 * 1024).saturating_sub(4096);
let samples_n = 4096u64;
let cpu_start = thread_cpu_seconds();
let wall = Instant::now();
let mut x: u64 = 0x9E37_79B9_7F4A_7C15; let mut samples: Vec<u64> = Vec::with_capacity(samples_n as usize);
for _ in 0..samples_n {
x = x
.wrapping_mul(6364136223846793005)
.wrapping_add(1442695040888963407);
let off = (x >> 16) % span;
let t = Instant::now();
let data = store.read_file(ino, off, 4096).expect("random read");
samples.push(t.elapsed().as_micros() as u64);
assert_eq!(data.len(), 4096);
}
let wall_s = wall.elapsed().as_secs_f64();
let cpu_s = thread_cpu_seconds() - cpu_start;
let mbps = (samples_n * 4096) as f64 / wall_s / 1024.0 / 1024.0;
let (p50, p95, p99) = percentiles(&mut samples);
(mbps, cpu_s, p50, p95, p99)
}
fn percentiles(samples: &mut [u64]) -> (f64, f64, f64) {
if samples.is_empty() {
return (0.0, 0.0, 0.0);
}
samples.sort_unstable();
let p = |q: f64| {
let i = ((samples.len() - 1) as f64 * q).round() as usize;
samples[i] as f64
};
(p(0.50), p(0.95), p(0.99))
}
fn mixed_phase(store: &Store, ino: u64, mib: u64) -> (f64, f64) {
let chunk = 64 * 1024usize;
let ops = mib * 1024 * 1024 / chunk as u64;
let cpu_start = thread_cpu_seconds();
let wall = Instant::now();
for i in 0..ops {
let off = (i % 16) * chunk as u64;
if i % 2 == 0 {
let mut data = vec![(i % 251) as u8; chunk];
data[0] = i as u8;
store.write_region(ino, off, &data).expect("mixed write");
} else {
store.read_file(ino, off, chunk as u64).expect("mixed read");
}
}
let wall_s = wall.elapsed().as_secs_f64();
let cpu_s = thread_cpu_seconds() - cpu_start;
(
(ops * chunk as u64) as f64 / wall_s / 1024.0 / 1024.0,
cpu_s,
)
}
#[test]
fn transport_real_court() {
let device = env_or("TRANSPORT_DEVICE", "/dev/shm");
let backend_s = env_or("TRANSPORT_BACKEND", "sync");
let backend = crate::store::io::IoBackendKind::parse(&backend_s).expect("backend");
let work_mib: u64 = env_or("TRANSPORT_WORK_MIB", "256").parse().expect("mib");
let dir = std::path::PathBuf::from(&device).join(format!(
"efs-transport-{}-{}",
backend_s,
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("device dir");
let config = StoreConfig {
io_backend: backend,
..Default::default()
};
let store = Store::create(&dir, &config, [0x55; 16]).expect("create");
{
let mut tx = store.begin_tx().expect("tx");
let inode = Inode::new_file(1000, 1000, 0o644);
Store::put_inode_in_tx(&mut tx, 7, &inode).expect("inode 7");
Store::put_inode_in_tx(&mut tx, 8, &inode).expect("inode 8");
tx.commit(&CrashHooks::none()).expect("commit");
}
let (w_mbps, w_wall, w_cpu, w_bar) = write_phase(&store, 7, work_mib, 0);
let (f_mbps, _, f_cpu, f_bar) = write_phase(&store, 8, work_mib / 4, 1);
let (r_mbps, r_cpu, r50, r95, r99) = read_phase_seq(&store, 7, work_mib / 2);
let (rr_mbps, rr_cpu, rr50, rr95, rr99) = read_phase_random(&store, 7, work_mib / 4);
let (m_mbps, m_cpu) = mixed_phase(&store, 7, work_mib / 8);
let phases: Vec<(String, u64, f64, f64, f64, f64)> = store
.perf()
.snapshot()
.into_iter()
.map(|row| {
(
row.phase.to_string(),
row.count,
row.total_ms,
row.p50_us,
row.p95_us,
row.p99_us,
)
})
.collect();
let result = serde_json::json!({
"device": device,
"backend": backend_s,
"work_mib": work_mib,
"write_mbps": w_mbps,
"write_wall_s": w_wall,
"write_cpu_s": w_cpu,
"write_barriers": w_bar,
"write_fsync_every_mbps": f_mbps,
"write_fsync_cpu_s": f_cpu,
"write_fsync_barriers": f_bar,
"read_mbps": r_mbps,
"read_cpu_s": r_cpu,
"read_p50_us": r50,
"read_p95_us": r95,
"read_p99_us": r99,
"random_read_mbps": rr_mbps,
"random_read_cpu_s": rr_cpu,
"random_read_p50_us": rr50,
"random_read_p95_us": rr95,
"random_read_p99_us": rr99,
"mixed_mbps": m_mbps,
"mixed_cpu_s": m_cpu,
"phases": phases,
});
println!("TRANSPORT_RESULT {}", result);
eprintln!(
"transport: device={device} backend={backend_s} write={w_mbps:.1} MiB/s \
(cpu {w_cpu:.2}s/{w_wall:.2}s) fsync={f_mbps:.1} MiB/s read={r_mbps:.1} MiB/s \
p50={r50:.0} p95={r95:.0} p99={r99:.0} µs random={rr_mbps:.1} MiB/s \
mixed={m_mbps:.1} MiB/s"
);
drop(store);
let _ = std::fs::remove_dir_all(&dir);
}