#![forbid(unsafe_code)]
pub mod trace;
use std::collections::HashMap;
use std::sync::Mutex;
use std::time::Instant;
pub fn wall_ns() -> u64 {
static START: std::sync::OnceLock<Instant> = std::sync::OnceLock::new();
let start = *START.get_or_init(Instant::now);
start.elapsed().as_nanos() as u64
}
const MAX_SAMPLES: usize = 4096;
thread_local! {
static CURRENT_REQUEST: std::cell::Cell<Option<u64>> = const {
std::cell::Cell::new(None)
};
static DETACH_DEPTH: std::cell::Cell<u32> = const {
std::cell::Cell::new(0)
};
}
#[derive(Debug, Default)]
struct Phase {
nanos_total: u128,
count: u64,
samples: Vec<u64>,
}
#[derive(Debug, Default)]
pub struct Timings {
phases: Mutex<HashMap<&'static str, Phase>>,
requests: Mutex<HashMap<u64, RequestAcc>>,
completed: Mutex<Vec<RequestResult>>,
next_request: std::sync::atomic::AtomicU64,
}
#[derive(Debug)]
struct RequestAcc {
name: &'static str,
t0: Instant,
phases: HashMap<&'static str, u64>,
}
#[derive(Debug, Clone)]
pub struct RequestResult {
pub name: &'static str,
pub total_ns: u64,
pub phases: Vec<(&'static str, u64)>,
pub residual_ns: i128,
}
#[derive(Debug, Clone)]
pub struct ReconRow {
pub phase: &'static str,
pub total_ms: f64,
pub share: f64,
}
#[derive(Debug, Clone)]
pub struct Reconciliation {
pub requests: u64,
pub total_ms: f64,
pub rows: Vec<ReconRow>,
pub residual_ms: f64,
pub residual_share: f64,
pub overlap: bool,
}
pub struct RequestGuard<'a> {
timings: &'a Timings,
id: Option<u64>,
}
impl Drop for RequestGuard<'_> {
fn drop(&mut self) {
let Some(id) = self.id else {
return;
};
CURRENT_REQUEST.set(None);
let mut m = self.timings.requests.lock().expect("requests poisoned");
let acc = m.remove(&id).expect("request envelope must exist");
let mut phases: Vec<(&'static str, u64)> = acc.phases.into_iter().collect();
phases.sort_by_key(|&(_, ns)| std::cmp::Reverse(ns));
let total_ns = acc.t0.elapsed().as_nanos() as u64;
let sum: u64 = phases.iter().map(|(_, ns)| *ns).sum();
drop(m);
self.timings
.completed
.lock()
.expect("completed poisoned")
.push(RequestResult {
name: acc.name,
total_ns,
phases,
residual_ns: total_ns as i128 - sum as i128,
});
}
}
#[derive(Debug, Clone, Copy)]
pub struct TimingRow {
pub phase: &'static str,
pub count: u64,
pub total_ms: f64,
pub p50_us: f64,
pub p95_us: f64,
pub p99_us: f64,
}
fn percentile(sorted: &[u64], q: f64) -> f64 {
if sorted.is_empty() {
return 0.0;
}
let idx = ((sorted.len() - 1) as f64 * q).round() as usize;
sorted[idx] as f64
}
impl Timings {
pub fn record(&self, phase: &'static str, nanos: u64) {
let mut m = self.phases.lock().expect("timings poisoned");
let p = m.entry(phase).or_default();
p.nanos_total += nanos as u128;
p.count += 1;
if p.samples.len() >= MAX_SAMPLES {
p.samples.remove(0);
}
p.samples.push(nanos);
}
pub fn time<T>(&self, phase: &'static str, f: impl FnOnce() -> T) -> T {
let t = Instant::now();
let out = f();
self.record(phase, t.elapsed().as_nanos() as u64);
out
}
pub fn snapshot(&self) -> Vec<TimingRow> {
let m = self.phases.lock().expect("timings poisoned");
let mut rows: Vec<TimingRow> = m
.iter()
.map(|(name, p)| {
let mut s = p.samples.clone();
s.sort_unstable();
TimingRow {
phase: name,
count: p.count,
total_ms: p.nanos_total as f64 / 1e6,
p50_us: percentile(&s, 0.50) / 1e3,
p95_us: percentile(&s, 0.95) / 1e3,
p99_us: percentile(&s, 0.99) / 1e3,
}
})
.collect();
rows.sort_by(|a, b| b.total_ms.total_cmp(&a.total_ms));
rows
}
pub fn render(&self) -> String {
let mut out = String::new();
out.push_str("phase timings (cumulative ms, per-sample us p50/p95/p99):\n");
for r in self.snapshot() {
out.push_str(&format!(
" {:<28} n={:>8} total={:>10.2} ms p50={:>9.1} p95={:>9.1} p99={:>9.1} us\n",
r.phase, r.count, r.total_ms, r.p50_us, r.p95_us, r.p99_us
));
}
out
}
pub fn request(&self, name: &'static str) -> RequestGuard<'_> {
if CURRENT_REQUEST.get().is_some() {
return RequestGuard {
timings: self,
id: None,
};
}
let id = self
.next_request
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
self.requests.lock().expect("requests poisoned").insert(
id,
RequestAcc {
name,
t0: Instant::now(),
phases: HashMap::new(),
},
);
CURRENT_REQUEST.set(Some(id));
RequestGuard {
timings: self,
id: Some(id),
}
}
pub fn time_request<T>(&self, phase: &'static str, f: impl FnOnce() -> T) -> T {
let t = Instant::now();
let out = f();
let ns = t.elapsed().as_nanos() as u64;
self.record(phase, ns);
if DETACH_DEPTH.get() == 0 {
if let Some(id) = CURRENT_REQUEST.get() {
let mut m = self.requests.lock().expect("requests poisoned");
if let Some(acc) = m.get_mut(&id) {
*acc.phases.entry(phase).or_default() += ns;
}
}
}
out
}
pub fn detach<T>(&self, f: impl FnOnce() -> T) -> T {
DETACH_DEPTH.set(DETACH_DEPTH.get() + 1);
let out = f();
DETACH_DEPTH.set(DETACH_DEPTH.get() - 1);
out
}
pub fn results(&self) -> Vec<RequestResult> {
self.completed.lock().expect("completed poisoned").clone()
}
pub fn clear(&self) {
self.phases.lock().expect("timings poisoned").clear();
self.requests.lock().expect("requests poisoned").clear();
self.completed.lock().expect("completed poisoned").clear();
}
pub fn reconcile(&self) -> Reconciliation {
let completed = self.completed.lock().expect("completed poisoned");
let mut total: u128 = 0;
let mut sums: HashMap<&'static str, u128> = HashMap::new();
for r in completed.iter() {
total += r.total_ns as u128;
for (p, ns) in &r.phases {
*sums.entry(p).or_default() += *ns as u128;
}
}
let phase_sum: u128 = sums.values().sum();
let residual = total as i128 - phase_sum as i128;
let mut rows: Vec<ReconRow> = sums
.iter()
.map(|(p, ns)| ReconRow {
phase: p,
total_ms: *ns as f64 / 1e6,
share: if total > 0 {
*ns as f64 / total as f64
} else {
0.0
},
})
.collect();
rows.sort_by(|a, b| b.total_ms.total_cmp(&a.total_ms));
Reconciliation {
requests: completed.len() as u64,
total_ms: total as f64 / 1e6,
rows,
residual_ms: residual as f64 / 1e6,
residual_share: if total > 0 {
residual as f64 / total as f64
} else {
0.0
},
overlap: residual < 0,
}
}
pub fn render_reconciled(&self) -> String {
let r = self.reconcile();
let mut out = String::new();
out.push_str(&format!(
"request reconciliation (n={} requests, {:.2} ms total):\n",
r.requests, r.total_ms
));
out.push_str(&format!(
" {:<26} {:>12} {:>9}\n",
"phase", "total ms", "share"
));
for row in &r.rows {
out.push_str(&format!(
" {:<26} {:>12.2} {:>8.1}%\n",
row.phase,
row.total_ms,
row.share * 100.0
));
}
out.push_str(&format!(
" {:<26} {:>12.2} {:>8.1}% <- residual (fuse/scheduler/other)\n",
"unaccounted",
r.residual_ms,
r.residual_share * 100.0
));
out.push_str(&format!(
" {:<26} {:>12.2} {:>8.1}% <- sum(phases) + residual == total {}\n",
"total",
r.total_ms,
100.0,
if r.overlap {
"OVERLAP! (nested phase rows)"
} else {
"OK"
}
));
out
}
}
pub const WRITE_BUCKETS: [&str; 7] = ["<4K", "4-16K", "16-64K", "64-256K", "256K-1M", "=1M", ">1M"];
fn write_bucket(len: usize) -> usize {
match len {
0..=4095 => 0,
4096..=16383 => 1,
16384..=65535 => 2,
65536..=262143 => 3,
262144..=1048575 => 4,
1048576 => 5,
_ => 6,
}
}
#[derive(Debug, Default)]
struct OpStat {
count: u64,
nanos_total: u128,
samples: Vec<u64>,
}
#[derive(Debug, Default)]
pub struct FuseStats {
ops: Mutex<HashMap<&'static str, OpStat>>,
write_sizes: Mutex<[u64; 7]>,
in_flight: std::sync::atomic::AtomicU64,
max_in_flight: std::sync::atomic::AtomicU64,
}
pub struct InFlight<'a> {
stats: &'a FuseStats,
}
impl<'a> InFlight<'a> {
pub fn begin(stats: &'a FuseStats) -> Self {
let now = stats
.in_flight
.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
+ 1;
stats
.max_in_flight
.fetch_max(now, std::sync::atomic::Ordering::Relaxed);
Self { stats }
}
}
impl Drop for InFlight<'_> {
fn drop(&mut self) {
self.stats
.in_flight
.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
}
}
impl FuseStats {
pub fn record_op(&self, op: &'static str, nanos: u64) {
let mut m = self.ops.lock().expect("fuse stats poisoned");
let s = m.entry(op).or_default();
s.count += 1;
s.nanos_total += nanos as u128;
if s.samples.len() >= MAX_SAMPLES {
s.samples.remove(0);
}
s.samples.push(nanos);
}
pub fn record_write_size(&self, len: usize) {
let mut h = self.write_sizes.lock().expect("fuse stats poisoned");
h[write_bucket(len)] += 1;
}
pub fn max_concurrency(&self) -> u64 {
self.max_in_flight
.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn snapshot(&self) -> Vec<(String, u64, f64, f64, f64, f64)> {
let m = self.ops.lock().expect("fuse stats poisoned");
let mut rows: Vec<(String, u64, f64, f64, f64, f64)> = m
.iter()
.map(|(name, s)| {
let mut v = s.samples.clone();
v.sort_unstable();
(
name.to_string(),
s.count,
s.nanos_total as f64 / 1e6,
percentile(&v, 0.50) / 1e3,
percentile(&v, 0.95) / 1e3,
percentile(&v, 0.99) / 1e3,
)
})
.collect();
rows.sort_by(|a, b| b.2.total_cmp(&a.2));
rows
}
pub fn render(&self) -> String {
let mut out = String::new();
out.push_str(&format!(
"fuse requests: max concurrency {}\n",
self.max_concurrency()
));
for (op, count, total_ms, p50, p95, p99) in self.snapshot() {
out.push_str(&format!(
" {:<12} n={:>8} total={:>10.2} ms p50={:>9.1} p95={:>9.1} p99={:>9.1} us\n",
op, count, total_ms, p50, p95, p99
));
}
let h = self.write_sizes.lock().expect("fuse stats poisoned");
out.push_str("write request size histogram:\n");
for (i, label) in WRITE_BUCKETS.iter().enumerate() {
out.push_str(&format!(" {label:>10}: {}\n", h[i]));
}
out
}
}