use std::{
fs,
future::Future,
process::Command,
sync::{
OnceLock,
atomic::{AtomicBool, AtomicU64, Ordering},
},
time::Instant,
};
#[derive(Debug, Clone, Copy)]
pub enum ProfileBucket {
SocketFetch,
ResponseParse,
RecordDecode,
PayloadCopy,
OffsetBookkeeping,
ChannelHandoff,
SourceEmit,
DatumConsume,
CommitBookkeeping,
}
#[derive(Debug, Clone, Copy, Default)]
pub struct ProfileBucketSnapshot {
pub calls: u64,
pub wall_ns: u64,
pub cpu_ns: u64,
}
#[derive(Debug, Clone, Copy, Default)]
pub struct ProfileSnapshot {
pub socket_fetch: ProfileBucketSnapshot,
pub response_parse: ProfileBucketSnapshot,
pub record_decode: ProfileBucketSnapshot,
pub payload_copy: ProfileBucketSnapshot,
pub offset_bookkeeping: ProfileBucketSnapshot,
pub channel_handoff: ProfileBucketSnapshot,
pub source_emit: ProfileBucketSnapshot,
pub datum_consume: ProfileBucketSnapshot,
pub commit_bookkeeping: ProfileBucketSnapshot,
}
#[derive(Debug, Default)]
struct BucketCounters {
calls: AtomicU64,
wall_ns: AtomicU64,
cpu_ns: AtomicU64,
}
#[derive(Debug, Default)]
struct ProfileCounters {
socket_fetch: BucketCounters,
response_parse: BucketCounters,
record_decode: BucketCounters,
payload_copy: BucketCounters,
offset_bookkeeping: BucketCounters,
channel_handoff: BucketCounters,
source_emit: BucketCounters,
datum_consume: BucketCounters,
commit_bookkeeping: BucketCounters,
}
static ENABLED: AtomicBool = AtomicBool::new(false);
static COUNTERS: ProfileCounters = ProfileCounters {
socket_fetch: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
response_parse: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
record_decode: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
payload_copy: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
offset_bookkeeping: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
channel_handoff: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
source_emit: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
datum_consume: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
commit_bookkeeping: BucketCounters {
calls: AtomicU64::new(0),
wall_ns: AtomicU64::new(0),
cpu_ns: AtomicU64::new(0),
},
};
pub fn set_enabled(enabled: bool) {
ENABLED.store(enabled, Ordering::Relaxed);
}
pub fn reset() {
COUNTERS.reset();
}
#[must_use]
pub fn enabled() -> bool {
ENABLED.load(Ordering::Relaxed)
}
pub fn measure<T, F>(bucket: ProfileBucket, f: F) -> T
where
F: FnOnce() -> T,
{
if !enabled() {
return f();
}
let wall_start = Instant::now();
let cpu_start = thread_cpu_ns();
let result = f();
record(
bucket,
wall_start.elapsed().as_nanos() as u64,
thread_cpu_ns().saturating_sub(cpu_start),
);
result
}
pub async fn measure_async<T, F>(bucket: ProfileBucket, future: F) -> T
where
F: Future<Output = T>,
{
if !enabled() {
return future.await;
}
let wall_start = Instant::now();
let cpu_start = thread_cpu_ns();
let result = future.await;
record(
bucket,
wall_start.elapsed().as_nanos() as u64,
thread_cpu_ns().saturating_sub(cpu_start),
);
result
}
pub fn snapshot() -> ProfileSnapshot {
ProfileSnapshot {
socket_fetch: snapshot_bucket(&COUNTERS.socket_fetch),
response_parse: snapshot_bucket(&COUNTERS.response_parse),
record_decode: snapshot_bucket(&COUNTERS.record_decode),
payload_copy: snapshot_bucket(&COUNTERS.payload_copy),
offset_bookkeeping: snapshot_bucket(&COUNTERS.offset_bookkeeping),
channel_handoff: snapshot_bucket(&COUNTERS.channel_handoff),
source_emit: snapshot_bucket(&COUNTERS.source_emit),
datum_consume: snapshot_bucket(&COUNTERS.datum_consume),
commit_bookkeeping: snapshot_bucket(&COUNTERS.commit_bookkeeping),
}
}
fn record(bucket: ProfileBucket, wall_ns: u64, cpu_ns: u64) {
let counters = match bucket {
ProfileBucket::SocketFetch => &COUNTERS.socket_fetch,
ProfileBucket::ResponseParse => &COUNTERS.response_parse,
ProfileBucket::RecordDecode => &COUNTERS.record_decode,
ProfileBucket::PayloadCopy => &COUNTERS.payload_copy,
ProfileBucket::OffsetBookkeeping => &COUNTERS.offset_bookkeeping,
ProfileBucket::ChannelHandoff => &COUNTERS.channel_handoff,
ProfileBucket::SourceEmit => &COUNTERS.source_emit,
ProfileBucket::DatumConsume => &COUNTERS.datum_consume,
ProfileBucket::CommitBookkeeping => &COUNTERS.commit_bookkeeping,
};
counters.calls.fetch_add(1, Ordering::Relaxed);
counters.wall_ns.fetch_add(wall_ns, Ordering::Relaxed);
counters.cpu_ns.fetch_add(cpu_ns, Ordering::Relaxed);
}
fn snapshot_bucket(counters: &BucketCounters) -> ProfileBucketSnapshot {
ProfileBucketSnapshot {
calls: counters.calls.load(Ordering::Relaxed),
wall_ns: counters.wall_ns.load(Ordering::Relaxed),
cpu_ns: counters.cpu_ns.load(Ordering::Relaxed),
}
}
impl ProfileCounters {
fn reset(&self) {
self.socket_fetch.reset();
self.response_parse.reset();
self.record_decode.reset();
self.payload_copy.reset();
self.offset_bookkeeping.reset();
self.channel_handoff.reset();
self.source_emit.reset();
self.datum_consume.reset();
self.commit_bookkeeping.reset();
}
}
impl BucketCounters {
fn reset(&self) {
self.calls.store(0, Ordering::Relaxed);
self.wall_ns.store(0, Ordering::Relaxed);
self.cpu_ns.store(0, Ordering::Relaxed);
}
}
fn thread_cpu_ns() -> u64 {
let Ok(stat) = fs::read_to_string("/proc/thread-self/stat") else {
return 0;
};
let Some(close) = stat.rfind(')') else {
return 0;
};
let fields = stat[close + 1..].split_whitespace().collect::<Vec<_>>();
if fields.len() <= 12 {
return 0;
}
let utime = fields[11].parse::<u64>().unwrap_or(0);
let stime = fields[12].parse::<u64>().unwrap_or(0);
((utime + stime) as f64 * 1_000_000_000.0 / clock_ticks_per_second()) as u64
}
fn clock_ticks_per_second() -> f64 {
static TICKS: OnceLock<f64> = OnceLock::new();
*TICKS.get_or_init(|| {
Command::new("getconf")
.arg("CLK_TCK")
.output()
.ok()
.and_then(|output| String::from_utf8(output.stdout).ok())
.and_then(|value| value.trim().parse::<f64>().ok())
.unwrap_or(100.0)
})
}