use std::os::raw::c_void;
use std::sync::RwLock;
use abi_stable::std_types::{RSlice, RStr};
#[repr(u8)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MetricKind {
Counter = 0,
Gauge = 1,
Histogram = 2,
}
pub type MetricLabel<'a> = (RStr<'a>, RStr<'a>);
pub type RecordMetricFn = extern "C" fn(
ctx: *const c_void,
name: RStr<'_>,
labels: RSlice<'_, MetricLabel<'_>>,
value: f64,
kind: MetricKind,
);
pub struct MetricRecorder {
inner: RwLock<Option<RecorderInner>>,
}
struct RecorderInner {
record_fn: RecordMetricFn,
ctx: *const c_void,
}
unsafe impl Send for RecorderInner {}
unsafe impl Sync for RecorderInner {}
impl MetricRecorder {
#[must_use]
pub const fn new() -> Self {
Self {
inner: RwLock::new(None),
}
}
#[doc(hidden)]
pub fn install(&self, record_fn: RecordMetricFn, ctx: *const c_void) {
let new = RecorderInner { record_fn, ctx };
match self.inner.write() {
Ok(mut guard) => *guard = Some(new),
Err(poisoned) => *poisoned.into_inner() = Some(new),
}
}
pub fn counter(&self, name: &str, labels: &[(&str, &str)], increment: u64) {
#[allow(clippy::cast_precision_loss)]
self.emit(name, labels, increment as f64, MetricKind::Counter);
}
pub fn gauge(&self, name: &str, labels: &[(&str, &str)], value: f64) {
self.emit(name, labels, value, MetricKind::Gauge);
}
pub fn histogram(&self, name: &str, labels: &[(&str, &str)], value: f64) {
self.emit(name, labels, value, MetricKind::Histogram);
}
fn emit(&self, name: &str, labels: &[(&str, &str)], value: f64, kind: MetricKind) {
let guard = match self.inner.read() {
Ok(g) => g,
Err(poisoned) => poisoned.into_inner(),
};
let Some(ref inner) = *guard else {
return;
};
let ffi_labels: Vec<MetricLabel<'_>> = labels
.iter()
.map(|(k, v)| (RStr::from(*k), RStr::from(*v)))
.collect();
(inner.record_fn)(
inner.ctx,
RStr::from(name),
RSlice::from(ffi_labels.as_slice()),
value,
kind,
);
}
}
impl Default for MetricRecorder {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Debug for MetricRecorder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let installed = match self.inner.try_read() {
Ok(g) => Some(g.is_some()),
Err(_) => None,
};
f.debug_struct("MetricRecorder")
.field("installed", &installed)
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::ptr;
use std::sync::atomic::{AtomicU64, Ordering};
static EMITTED_COUNT: AtomicU64 = AtomicU64::new(0);
extern "C" fn test_recorder(
_ctx: *const c_void,
_name: RStr<'_>,
_labels: RSlice<'_, MetricLabel<'_>>,
_value: f64,
_kind: MetricKind,
) {
EMITTED_COUNT.fetch_add(1, Ordering::Relaxed);
}
#[test]
fn emit_without_install_is_silent_noop() {
let r = MetricRecorder::new();
r.counter("test", &[], 1);
r.gauge("test", &[], 1.0);
r.histogram("test", &[], 1.0);
}
#[test]
fn install_then_emit_routes_to_recorder() {
let before = EMITTED_COUNT.load(Ordering::Relaxed);
let r = MetricRecorder::new();
r.install(test_recorder, ptr::null());
r.counter("dispatches_total", &[("kind", "ok")], 1);
r.gauge("in_flight", &[], 7.0);
r.histogram("latency_seconds", &[("path", "/x")], 0.012);
assert_eq!(EMITTED_COUNT.load(Ordering::Relaxed) - before, 3);
}
#[test]
fn install_replaces_previous_recorder() {
static EMITTED_VIA_CTX: [AtomicU64; 2] = [AtomicU64::new(0), AtomicU64::new(0)];
extern "C" fn recorder_route_by_ctx(
ctx: *const c_void,
_name: RStr<'_>,
_labels: RSlice<'_, MetricLabel<'_>>,
_value: f64,
_kind: MetricKind,
) {
let idx = ctx as usize;
if idx < EMITTED_VIA_CTX.len() {
EMITTED_VIA_CTX[idx].fetch_add(1, Ordering::Relaxed);
}
}
let r = MetricRecorder::new();
r.install(recorder_route_by_ctx, std::ptr::null::<c_void>());
r.counter("c1", &[], 1);
assert_eq!(EMITTED_VIA_CTX[0].load(Ordering::Relaxed), 1);
assert_eq!(EMITTED_VIA_CTX[1].load(Ordering::Relaxed), 0);
r.install(recorder_route_by_ctx, std::ptr::dangling::<c_void>());
r.counter("c2", &[], 1);
assert_eq!(EMITTED_VIA_CTX[0].load(Ordering::Relaxed), 1);
assert_eq!(EMITTED_VIA_CTX[1].load(Ordering::Relaxed), 1);
}
#[test]
fn install_during_concurrent_emit_does_not_tear() {
use std::sync::Arc;
use std::thread;
use std::time::Duration;
static OK_CALLS: AtomicU64 = AtomicU64::new(0);
extern "C" fn ok_recorder(
_ctx: *const c_void,
_name: RStr<'_>,
_labels: RSlice<'_, MetricLabel<'_>>,
_value: f64,
_kind: MetricKind,
) {
OK_CALLS.fetch_add(1, Ordering::Relaxed);
}
let r = Arc::new(MetricRecorder::new());
r.install(ok_recorder, ptr::null());
let emitter = {
let r = Arc::clone(&r);
thread::spawn(move || {
for _ in 0..10_000 {
r.counter("dispatches_total", &[], 1);
}
})
};
let installer = {
let r = Arc::clone(&r);
thread::spawn(move || {
for i in 0..1_000 {
r.install(ok_recorder, i as *const c_void);
thread::sleep(Duration::from_micros(10));
}
})
};
emitter.join().unwrap();
installer.join().unwrap();
assert_eq!(OK_CALLS.load(Ordering::Relaxed), 10_000);
}
}