use std::sync::Arc;
use std::time::Instant;
#[cfg(target_os = "linux")]
use goosefs_sdk::cache::store::UringPageStore;
use goosefs_sdk::cache::store::{LocalPageStore, PageStore};
use goosefs_sdk::cache::PageId;
fn env_or<T: std::str::FromStr>(key: &str, default: T) -> T {
std::env::var(key)
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(default)
}
struct BenchResult {
label: &'static str,
ops_per_sec: f64,
p50_ns: u64,
p99_ns: u64,
total_ns: u64,
}
async fn bench_single_threaded(
store: &Arc<dyn PageStore>,
page_id: &PageId,
page_size: usize,
iterations: usize,
label: &'static str,
) -> BenchResult {
let mut dst = vec![0u8; page_size];
match tokio::time::timeout(
std::time::Duration::from_secs(5),
store.get(page_id, 0, &mut dst),
)
.await
{
Ok(Ok(_)) => {}
Ok(Err(e)) => panic!("{label} warm-up failed: {e}"),
Err(_) => panic!(
"{label} warm-up TIMED OUT after 5s — io_uring backend may be broken; \
check RUST_LOG=trace output and `dmesg | tail` for kernel errors"
),
}
let mut latencies_ns: Vec<u64> = Vec::with_capacity(iterations);
let start = Instant::now();
for _ in 0..iterations {
let op_start = Instant::now();
let n = store.get(page_id, 0, &mut dst).await.expect("get failed");
debug_assert_eq!(n, page_size, "short read in benchmark");
latencies_ns.push(op_start.elapsed().as_nanos() as u64);
}
let total = start.elapsed();
latencies_ns.sort_unstable();
let p50 = latencies_ns[latencies_ns.len() / 2];
let p99 = latencies_ns[latencies_ns.len() * 99 / 100];
let ops_per_sec = iterations as f64 / total.as_secs_f64().max(1e-9);
BenchResult {
label,
ops_per_sec,
p50_ns: p50,
p99_ns: p99,
total_ns: total.as_nanos() as u64,
}
}
async fn bench_concurrent(
store: Arc<dyn PageStore>,
page_id: PageId,
page_size: usize,
concurrency: usize,
iterations_per_task: usize,
label: &'static str,
) -> BenchResult {
{
let mut dst = vec![0u8; page_size];
let _ = store.get(&page_id, 0, &mut dst).await;
}
let start = Instant::now();
let mut handles = Vec::with_capacity(concurrency);
for _ in 0..concurrency {
let store = Arc::clone(&store);
let pid = page_id.clone();
handles.push(tokio::spawn(async move {
let mut dst = vec![0u8; page_size];
let mut latencies: Vec<u64> = Vec::with_capacity(iterations_per_task);
for _ in 0..iterations_per_task {
let op_start = Instant::now();
let n = store.get(&pid, 0, &mut dst).await.expect("get failed");
debug_assert_eq!(n, page_size);
latencies.push(op_start.elapsed().as_nanos() as u64);
}
latencies
}));
}
let mut all_latencies: Vec<u64> = Vec::with_capacity(concurrency * iterations_per_task);
for h in handles {
all_latencies.extend(h.await.unwrap());
}
let total = start.elapsed();
all_latencies.sort_unstable();
let p50 = all_latencies[all_latencies.len() / 2];
let p99 = all_latencies[all_latencies.len() * 99 / 100];
let total_ops = concurrency * iterations_per_task;
let ops_per_sec = total_ops as f64 / total.as_secs_f64().max(1e-9);
BenchResult {
label,
ops_per_sec,
p50_ns: p50,
p99_ns: p99,
total_ns: total.as_nanos() as u64,
}
}
fn print_result(r: &BenchResult) {
println!(
" {:<16} {:>10.0} ops/s p50={:>6}µs p99={:>6}µs total={:.2}s",
r.label,
r.ops_per_sec,
r.p50_ns / 1000,
r.p99_ns / 1000,
r.total_ns as f64 / 1e9,
);
}
fn print_header(title: &str) {
println!("\n── {title} ────────────────────────────────────────");
}
async fn bench_concurrent_multi_file(
store: Arc<dyn PageStore>,
page_ids: Vec<PageId>,
page_size: usize,
concurrency: usize,
iterations_per_task: usize,
label: &'static str,
) -> BenchResult {
for id in &page_ids {
let mut dst = vec![0u8; page_size];
let _ = store.get(id, 0, &mut dst).await;
}
let start = Instant::now();
let mut handles = Vec::with_capacity(concurrency);
for task_id in 0..concurrency {
let store = Arc::clone(&store);
let page_ids = page_ids.clone();
handles.push(tokio::spawn(async move {
let mut dst = vec![0u8; page_size];
let mut latencies: Vec<u64> = Vec::with_capacity(iterations_per_task);
for i in 0..iterations_per_task {
let id = &page_ids[(i + task_id) % page_ids.len()];
let op_start = Instant::now();
let n = store.get(id, 0, &mut dst).await.expect("get failed");
debug_assert_eq!(n, page_size);
latencies.push(op_start.elapsed().as_nanos() as u64);
}
latencies
}));
}
let mut all_latencies: Vec<u64> = Vec::with_capacity(concurrency * iterations_per_task);
for h in handles {
all_latencies.extend(h.await.unwrap());
}
let total = start.elapsed();
all_latencies.sort_unstable();
let p50 = all_latencies[all_latencies.len() / 2];
let p99 = all_latencies[all_latencies.len() * 99 / 100];
let total_ops = concurrency * iterations_per_task;
let ops_per_sec = total_ops as f64 / total.as_secs_f64().max(1e-9);
BenchResult {
label,
ops_per_sec,
p50_ns: p50,
p99_ns: p99,
total_ns: total.as_nanos() as u64,
}
}
#[tokio::main]
async fn main() {
let iterations: usize = env_or("BENCH_ITERATIONS", 100_000);
let concurrency: usize = env_or("BENCH_CONCURRENCY", 32);
let concurrent_iterations: usize = env_or("BENCH_CONCURRENT_ITERATIONS", 10_000);
let page_size: usize = env_or("BENCH_PAGE_SIZE", 1024);
let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let base_dir = std::env::temp_dir().join(format!("gfs_uring_bench_{ts}"));
println!("╔══════════════════════════════════════════════════════════════╗");
println!("║ Page-Cache Backend Benchmark: io_uring vs tokio::fs ║");
println!("╚══════════════════════════════════════════════════════════════╝");
println!(" page_size={page_size}B iterations={iterations} concurrency={concurrency}×{concurrent_iterations}");
println!(" cache_dir={}", base_dir.display());
#[cfg(target_os = "linux")]
{
match goosefs_sdk::cache::store::is_uring_available() {
true => println!(" io_uring: available"),
false => panic!(
"io_uring is NOT available on this platform. \
Set GOOSEFS_USER_CLIENT_CACHE_URING_ENABLED=false to skip the io_uring benchmark."
),
}
}
let local_dir = base_dir.join("tokio_fs");
let local_store: Arc<dyn PageStore> = Arc::new(
LocalPageStore::create(&local_dir, page_size as u64)
.await
.expect("LocalPageStore create"),
);
#[cfg(target_os = "linux")]
let uring_store: Arc<dyn PageStore> = {
goosefs_sdk::cache::store::init_uring_config(16384, 0);
let uring_dir = base_dir.join("uring");
Arc::new(
UringPageStore::create(&uring_dir, page_size as u64)
.await
.expect("UringPageStore create"),
)
};
let page_id = PageId::new("bench-file", 0);
let page_data = vec![0x42u8; page_size];
local_store
.put(&page_id, &page_data)
.await
.expect("local put");
#[cfg(target_os = "linux")]
uring_store
.put(&page_id, &page_data)
.await
.expect("uring put");
print_header("Single-threaded cache-hit throughput");
let r_local =
bench_single_threaded(&local_store, &page_id, page_size, iterations, "tokio::fs").await;
print_result(&r_local);
#[cfg(target_os = "linux")]
let r_uring = Some(
bench_single_threaded(&uring_store, &page_id, page_size, iterations, "io_uring").await,
);
#[cfg(not(target_os = "linux"))]
let r_uring: Option<BenchResult> = None;
if let Some(r) = &r_uring {
print_result(r);
let speedup = r.ops_per_sec / r_local.ops_per_sec.max(1.0);
println!(" → io_uring speedup: {speedup:.2}×");
}
#[cfg(not(target_os = "linux"))]
{
println!(" (io_uring backend not available on this platform)");
}
print_header(&format!(
"Concurrent cache-hit throughput ({concurrency} tasks)"
));
let rc_local = bench_concurrent(
Arc::clone(&local_store),
page_id.clone(),
page_size,
concurrency,
concurrent_iterations,
"tokio::fs",
)
.await;
print_result(&rc_local);
#[cfg(target_os = "linux")]
let rc_uring = Some(
bench_concurrent(
Arc::clone(&uring_store),
page_id.clone(),
page_size,
concurrency,
concurrent_iterations,
"io_uring",
)
.await,
);
#[cfg(not(target_os = "linux"))]
let rc_uring: Option<BenchResult> = None;
if let Some(rc) = &rc_uring {
print_result(rc);
let speedup = rc.ops_per_sec / rc_local.ops_per_sec.max(1.0);
println!(" → io_uring speedup: {speedup:.2}×");
}
#[cfg(not(target_os = "linux"))]
{
println!(" (io_uring backend not available on this platform)");
}
let n_files = env_or("BENCH_MULTI_FILE_COUNT", 64);
print_header(&format!(
"Multi-file cache-hit throughput ({n_files} files, {concurrency} concurrent tasks)"
));
let multi_file_ids: Vec<PageId> = (0..n_files)
.map(|i| PageId::new(format!("bench-file-{i}"), 0))
.collect();
for id in &multi_file_ids {
local_store
.put(id, &page_data)
.await
.expect("local put multi");
#[cfg(target_os = "linux")]
uring_store
.put(id, &page_data)
.await
.expect("uring put multi");
}
let rc_local_multi = bench_concurrent_multi_file(
Arc::clone(&local_store),
multi_file_ids.clone(),
page_size,
concurrency,
concurrent_iterations,
"tokio::fs",
)
.await;
print_result(&rc_local_multi);
#[cfg(target_os = "linux")]
let rc_uring_multi = Some(
bench_concurrent_multi_file(
Arc::clone(&uring_store),
multi_file_ids.clone(),
page_size,
concurrency,
concurrent_iterations,
"io_uring",
)
.await,
);
#[cfg(not(target_os = "linux"))]
let rc_uring_multi: Option<BenchResult> = None;
if let Some(rc) = &rc_uring_multi {
print_result(rc);
let speedup = rc.ops_per_sec / rc_local_multi.ops_per_sec.max(1.0);
println!(" → io_uring speedup: {speedup:.2}×");
}
println!("\n═══════════════════════════════════════════════════════════════");
println!(" Summary (page_size={page_size}B)");
println!("───────────────────────────────────────────────────────────────");
println!(
" {:<16} {:>14} {:>14} {:>10} {:>10}",
"Backend", "Single ops/s", "Conc ops/s", "p99(1T)", "p99(32T)"
);
println!("───────────────────────────────────────────────────────────────");
println!(
" {:<16} {:>14.0} {:>14.0} {:>8}µs {:>8}µs",
"tokio::fs",
r_local.ops_per_sec,
rc_local.ops_per_sec,
r_local.p99_ns / 1000,
rc_local.p99_ns / 1000,
);
if let (Some(r), Some(rc)) = (&r_uring, &rc_uring) {
println!(
" {:<16} {:>14.0} {:>14.0} {:>8}µs {:>8}µs",
"io_uring",
r.ops_per_sec,
rc.ops_per_sec,
r.p99_ns / 1000,
rc.p99_ns / 1000,
);
}
println!("───────────────────────────────────────────────────────────────");
let _ = tokio::fs::remove_dir_all(&base_dir).await;
}