#![allow(dead_code)]
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc;
use std::thread;
use std::time::{Duration, Instant};
pub fn bench_dir() -> PathBuf {
if let Some(dir) = std::env::var_os("RINGFIRE_BENCH_DIR") {
return PathBuf::from(dir);
}
let shm = Path::new("/dev/shm");
if shm.is_dir() {
shm.to_path_buf()
} else {
std::env::temp_dir()
}
}
pub struct TempShm(PathBuf);
impl TempShm {
pub fn new(name: &str) -> Self {
static NEXT: AtomicU64 = AtomicU64::new(0);
let n = NEXT.fetch_add(1, Ordering::Relaxed);
let path = bench_dir().join(format!(
"ringfire_bench_{}_{}_{}.shm",
name,
std::process::id(),
n
));
let _ = std::fs::remove_file(&path);
Self(path)
}
pub fn path(&self) -> &Path {
&self.0
}
}
impl Drop for TempShm {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.0);
}
}
pub struct AbortOnPanic(pub &'static str);
impl Drop for AbortOnPanic {
fn drop(&mut self) {
if std::thread::panicking() {
eprintln!("bench helper thread `{}` panicked; aborting", self.0);
std::process::abort();
}
}
}
pub struct SpinBound {
polls: u32,
since: Option<Instant>,
limit: Duration,
what: &'static str,
}
impl SpinBound {
pub const fn new(what: &'static str, limit: Duration) -> Self {
Self {
polls: 0,
since: None,
limit,
what,
}
}
#[inline(always)]
pub fn spin(&mut self) {
core::hint::spin_loop();
self.polls = self.polls.wrapping_add(1);
if self.polls & 0xFFFF == 0 {
self.check();
}
}
#[cold]
#[inline(never)]
fn check(&mut self) {
let now = Instant::now();
match self.since {
None => self.since = Some(now),
Some(since) if now - since > self.limit => {
panic!("{} did not complete within {:?}", self.what, self.limit)
}
Some(_) => {}
}
}
}
pub const REPLY_TIMEOUT: Duration = Duration::from_secs(10);
pub fn timed_chunks(
iters: u64,
chunk: u64,
prepare: impl FnMut(u64),
run: impl FnMut(u64),
) -> Duration {
timed_chunks_with_clock(iters, chunk, prepare, run, Instant::now, Instant::elapsed)
}
pub fn timed_chunks_with_clock<C>(
iters: u64,
chunk: u64,
mut prepare: impl FnMut(u64),
mut run: impl FnMut(u64),
mut now: impl FnMut() -> C,
mut elapsed: impl FnMut(&C) -> Duration,
) -> Duration {
assert!(chunk > 0);
let mut total = Duration::ZERO;
let mut left = iters;
while left > 0 {
let n = left.min(chunk);
prepare(n);
let start = now();
run(n);
total += elapsed(&start);
left -= n;
}
total
}
pub struct Watchdog {
cancel: Option<mpsc::Sender<()>>,
worker: Option<thread::JoinHandle<()>>,
}
impl Watchdog {
pub fn start(what: impl Into<String>, limit: Duration) -> Self {
let what = what.into();
let (tx, rx) = mpsc::channel();
let worker = thread::spawn(move || {
if matches!(rx.recv_timeout(limit), Err(mpsc::RecvTimeoutError::Timeout)) {
eprintln!(
"{} exceeded the whole-case deadline {:?}; aborting",
what, limit
);
std::process::abort();
}
});
Self {
cancel: Some(tx),
worker: Some(worker),
}
}
pub fn default_limit() -> Duration {
let seconds = std::env::var("RINGFIRE_BENCH_WATCHDOG_SECS")
.map(|s| {
s.parse::<u64>()
.expect("RINGFIRE_BENCH_WATCHDOG_SECS must be a positive integer")
})
.unwrap_or(900);
assert!(seconds > 0, "watchdog deadline must be positive");
Duration::from_secs(seconds)
}
}
impl Drop for Watchdog {
fn drop(&mut self) {
drop(self.cancel.take());
if let Some(worker) = self.worker.take() {
worker.join().unwrap();
}
}
}