use alloc::collections::BTreeMap;
use alloc::string::{String, ToString};
use alloc::sync::Arc;
use alloc::vec;
use alloc::vec::Vec;
use core::fmt::Write as FmtWrite;
use core::sync::atomic::Ordering;
use super::{
MensuraArboris, MensuraEffectus, MensuraFibrae, MetricSnapshot, MetricType, MetricValue,
};
struct SpinRwLock<T> {
data: core::cell::UnsafeCell<T>,
state: core::sync::atomic::AtomicUsize,
}
unsafe impl<T: Send> Send for SpinRwLock<T> {}
unsafe impl<T: Send + Sync> Sync for SpinRwLock<T> {}
impl<T> SpinRwLock<T> {
const fn new(data: T) -> Self {
SpinRwLock {
data: core::cell::UnsafeCell::new(data),
state: core::sync::atomic::AtomicUsize::new(0),
}
}
fn read(&self) -> SpinRwLockReadGuard<'_, T> {
loop {
let state = self.state.load(Ordering::Acquire);
if state != usize::MAX
&& self
.state
.compare_exchange_weak(state, state + 1, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
{
return SpinRwLockReadGuard {
lock: self,
_marker: core::marker::PhantomData,
};
}
core::hint::spin_loop();
}
}
fn write(&self) -> SpinRwLockWriteGuard<'_, T> {
loop {
if self
.state
.compare_exchange_weak(0, usize::MAX, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
{
return SpinRwLockWriteGuard {
lock: self,
_marker: core::marker::PhantomData,
};
}
core::hint::spin_loop();
}
}
}
struct SpinRwLockReadGuard<'a, T> {
lock: &'a SpinRwLock<T>,
_marker: core::marker::PhantomData<T>,
}
unsafe impl<T: Sync> Sync for SpinRwLockReadGuard<'_, T> {}
unsafe impl<T: Sync> Send for SpinRwLockReadGuard<'_, T> {}
impl<T> core::ops::Deref for SpinRwLockReadGuard<'_, T> {
type Target = T;
fn deref(&self) -> &T {
unsafe { &*self.lock.data.get() }
}
}
impl<T> Drop for SpinRwLockReadGuard<'_, T> {
fn drop(&mut self) {
self.lock.state.fetch_sub(1, Ordering::Release);
}
}
struct SpinRwLockWriteGuard<'a, T> {
lock: &'a SpinRwLock<T>,
_marker: core::marker::PhantomData<T>,
}
unsafe impl<T: Sync> Sync for SpinRwLockWriteGuard<'_, T> {}
unsafe impl<T: Send + Sync> Send for SpinRwLockWriteGuard<'_, T> {}
impl<T> core::ops::Deref for SpinRwLockWriteGuard<'_, T> {
type Target = T;
fn deref(&self) -> &T {
unsafe { &*self.lock.data.get() }
}
}
impl<T> core::ops::DerefMut for SpinRwLockWriteGuard<'_, T> {
fn deref_mut(&mut self) -> &mut T {
unsafe { &mut *self.lock.data.get() }
}
}
impl<T> Drop for SpinRwLockWriteGuard<'_, T> {
fn drop(&mut self) {
self.lock.state.store(0, Ordering::Release);
}
}
impl<'a, T> SpinRwLockWriteGuard<'a, T> {
#[inline]
pub(crate) fn downgrade(self) -> SpinRwLockReadGuard<'a, T> {
self.lock.state.store(1, Ordering::Release);
let lock = self.lock;
core::mem::forget(self);
SpinRwLockReadGuard {
lock,
_marker: core::marker::PhantomData,
}
}
}
pub struct RegistrumMensurarum {
effects: SpinRwLock<BTreeMap<u64, Arc<MensuraEffectus>>>,
fibers: Arc<MensuraFibrae>,
supervisors: SpinRwLock<BTreeMap<String, Arc<MensuraArboris>>>,
}
impl RegistrumMensurarum {
pub fn new() -> Self {
RegistrumMensurarum {
effects: SpinRwLock::new(BTreeMap::new()),
fibers: Arc::new(MensuraFibrae::new()),
supervisors: SpinRwLock::new(BTreeMap::new()),
}
}
#[inline]
pub fn effect_metrics(&self, effect_id: u64, name: &str) -> Arc<MensuraEffectus> {
{
let effects = self.effects.read();
if let Some(metrics) = effects.get(&effect_id) {
return metrics.clone();
}
}
let mut effects = self.effects.write();
effects
.entry(effect_id)
.or_insert_with(|| Arc::new(MensuraEffectus::new(effect_id, name)));
let effects = effects.downgrade();
effects.get(&effect_id).unwrap().clone()
}
#[inline]
pub fn fiber_metrics(&self) -> Arc<MensuraFibrae> {
self.fibers.clone()
}
#[inline]
pub fn supervisor_metrics(&self, name: &str) -> Arc<MensuraArboris> {
{
let supervisors = self.supervisors.read();
if let Some(metrics) = supervisors.get(name) {
return metrics.clone();
}
}
let mut supervisors = self.supervisors.write();
supervisors
.entry(name.into())
.or_insert_with(|| Arc::new(MensuraArboris::new(name)))
.clone()
}
#[inline]
pub fn effect_summaries(&self) -> Vec<super::EffectMetricsSummary> {
let effects = self.effects.read();
effects.values().map(|m| m.summary()).collect()
}
#[inline]
pub fn fiber_summary(&self) -> super::FiberMetricsSummary {
self.fibers.summary()
}
#[allow(clippy::too_many_lines)]
pub fn export(&self) -> Vec<MetricSnapshot> {
let mut snapshots = Vec::with_capacity(32);
let effect_metrics: Vec<Arc<MensuraEffectus>> = {
let effects = self.effects.read();
effects.values().cloned().collect()
};
for metrics in &effect_metrics {
snapshots.push(MetricSnapshot {
name: "effect_operations_total".into(),
metric_type: MetricType::Counter,
value: MetricValue::Counter(metrics.total_operations()),
labels: vec![
("effect_id".into(), metrics.effect_id().to_string()),
("effect_name".into(), metrics.effect_name().into()),
],
});
snapshots.push(MetricSnapshot {
name: "effect_successes_total".into(),
metric_type: MetricType::Counter,
value: MetricValue::Counter(metrics.success_count()),
labels: vec![
("effect_id".into(), metrics.effect_id().to_string()),
("effect_name".into(), metrics.effect_name().into()),
],
});
snapshots.push(MetricSnapshot {
name: "effect_failures_total".into(),
metric_type: MetricType::Counter,
value: MetricValue::Counter(metrics.failure_count()),
labels: vec![
("effect_id".into(), metrics.effect_id().to_string()),
("effect_name".into(), metrics.effect_name().into()),
],
});
snapshots.push(MetricSnapshot {
name: "effect_in_flight".into(),
metric_type: MetricType::Gauge,
value: MetricValue::Gauge(metrics.in_flight()),
labels: vec![
("effect_id".into(), metrics.effect_id().to_string()),
("effect_name".into(), metrics.effect_name().into()),
],
});
let hist = metrics.latency_histogram();
let buckets: Vec<(u64, u64)> = hist
.boundaries()
.iter()
.zip(hist.bucket_counts())
.map(|(&b, c)| (b, c))
.collect();
snapshots.push(MetricSnapshot {
name: "effect_latency_nanoseconds".into(),
metric_type: MetricType::Histogram,
value: MetricValue::Histogram {
count: hist.count(),
sum: hist.sum(),
buckets,
},
labels: vec![
("effect_id".into(), metrics.effect_id().to_string()),
("effect_name".into(), metrics.effect_name().into()),
],
});
}
let fibers = &self.fibers;
snapshots.push(MetricSnapshot {
name: "fibers_spawned_total".into(),
metric_type: MetricType::Counter,
value: MetricValue::Counter(fibers.total_spawned()),
labels: vec![],
});
snapshots.push(MetricSnapshot {
name: "fibers_completed_total".into(),
metric_type: MetricType::Counter,
value: MetricValue::Counter(fibers.completed_count()),
labels: vec![],
});
snapshots.push(MetricSnapshot {
name: "fibers_failed_total".into(),
metric_type: MetricType::Counter,
value: MetricValue::Counter(fibers.failed_count()),
labels: vec![],
});
snapshots.push(MetricSnapshot {
name: "fibers_cancelled_total".into(),
metric_type: MetricType::Counter,
value: MetricValue::Counter(fibers.cancelled_count()),
labels: vec![],
});
snapshots.push(MetricSnapshot {
name: "fibers_active".into(),
metric_type: MetricType::Gauge,
value: MetricValue::Gauge(fibers.active_count()),
labels: vec![],
});
snapshots.push(MetricSnapshot {
name: "fibers_peak_active".into(),
metric_type: MetricType::Gauge,
value: MetricValue::Gauge(fibers.peak_active()),
labels: vec![],
});
snapshots
}
fn effuge_valorem(v: &str) -> String {
let mut out = String::with_capacity(v.len());
for c in v.chars() {
match c {
'\\' => out.push_str("\\\\"),
'"' => out.push_str("\\\""),
'\n' => out.push_str("\\n"),
_ => out.push(c),
}
}
out
}
pub fn prometheus_export(&self) -> String {
let mut output = String::new();
for snapshot in self.export() {
let type_str = match snapshot.metric_type {
MetricType::Counter => "counter",
MetricType::Gauge => "gauge",
MetricType::Histogram => "histogram",
};
writeln!(output, "# TYPE {} {}", snapshot.name, type_str)
.expect("writing to String is infallible");
let labels = if snapshot.labels.is_empty() {
String::new()
} else {
let label_str: Vec<String> = snapshot
.labels
.iter()
.map(|(k, v)| alloc::format!("{k}=\"{}\"", Self::effuge_valorem(v)))
.collect();
alloc::format!("{{{}}}", label_str.join(","))
};
match snapshot.value {
MetricValue::Counter(v) => {
writeln!(output, "{}{} {}", snapshot.name, labels, v)
.expect("writing to String is infallible");
}
MetricValue::Gauge(v) => {
writeln!(output, "{}{} {}", snapshot.name, labels, v)
.expect("writing to String is infallible");
}
MetricValue::Histogram {
count,
sum,
ref buckets,
} => {
for (boundary, bucket_count) in buckets {
let boundary_str = Self::effuge_valorem(&boundary.to_string());
let bucket_labels = if labels.is_empty() {
alloc::format!("{{le=\"{boundary_str}\"}}")
} else {
let inner = &labels[1..labels.len() - 1];
alloc::format!("{{{inner},le=\"{boundary_str}\"}}")
};
writeln!(
output,
"{}_bucket{} {}",
snapshot.name, bucket_labels, bucket_count
)
.expect("writing to String never fails");
}
writeln!(output, "{}_sum{} {}", snapshot.name, labels, sum)
.expect("writing to String never fails");
writeln!(output, "{}_count{} {}", snapshot.name, labels, count)
.expect("writing to String never fails");
}
}
}
output
}
}
impl Default for RegistrumMensurarum {
fn default() -> Self {
Self::new()
}
}
pub fn global_registry() -> &'static RegistrumMensurarum {
use alloc::boxed::Box;
use core::ptr;
use core::sync::atomic::AtomicPtr;
static REGISTRY: AtomicPtr<RegistrumMensurarum> = AtomicPtr::new(ptr::null_mut());
let mut ptr = REGISTRY.load(Ordering::Acquire);
if ptr.is_null() {
let new_registry = Box::new(RegistrumMensurarum::new());
let new_ptr = Box::into_raw(new_registry);
match REGISTRY.compare_exchange(
ptr::null_mut(),
new_ptr,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => ptr = new_ptr,
Err(existing) => {
unsafe {
drop(Box::from_raw(new_ptr));
}
ptr = existing;
}
}
}
unsafe { &*ptr }
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_registry_new() {
let registry = RegistrumMensurarum::new();
assert!(registry.effect_summaries().is_empty());
}
#[test]
fn test_registry_effect_metrics() {
let registry = RegistrumMensurarum::new();
let metrics1 = registry.effect_metrics(1, "StateEffect");
let metrics2 = registry.effect_metrics(1, "StateEffect");
assert!(Arc::ptr_eq(&metrics1, &metrics2));
}
#[test]
fn test_registry_different_effects() {
let registry = RegistrumMensurarum::new();
let metrics1 = registry.effect_metrics(1, "StateEffect");
let metrics2 = registry.effect_metrics(2, "IOEffect");
assert!(!Arc::ptr_eq(&metrics1, &metrics2));
}
#[test]
fn test_registry_fiber_metrics() {
let registry = RegistrumMensurarum::new();
let fibers1 = registry.fiber_metrics();
let fibers2 = registry.fiber_metrics();
assert!(Arc::ptr_eq(&fibers1, &fibers2));
}
#[test]
fn test_registry_export() {
let registry = RegistrumMensurarum::new();
let effect = registry.effect_metrics(1, "Test");
effect.operation_start();
effect.record_success(1000);
let snapshots = registry.export();
assert!(!snapshots.is_empty());
}
#[test]
fn test_registry_prometheus_export() {
let registry = RegistrumMensurarum::new();
let effect = registry.effect_metrics(1, "Test");
effect.operation_start();
effect.record_success(1000);
let output = registry.prometheus_export();
assert!(output.contains("effect_operations_total"));
assert!(output.contains("effect_successes_total"));
}
#[test]
fn prometheus_export_escapes_label_values() {
let registry = RegistrumMensurarum::new();
let effect = registry.effect_metrics(1, "weird\"name\\with\nnewline");
effect.operation_start();
effect.record_success(1000);
let output = registry.prometheus_export();
assert!(
output.contains(r#"\"name"#),
"expected escaped quote: {output}"
);
assert!(
output.contains(r"\\with"),
"expected escaped backslash: {output}"
);
assert!(
output.contains(r"\nnewline"),
"expected escaped newline: {output}"
);
for line in output.lines() {
let opens = line.matches('{').count();
let closes = line.matches('}').count();
assert_eq!(
opens, closes,
"unbalanced braces (raw newline split a label?): {line}"
);
}
}
#[test]
fn test_global_registry() {
let registry1 = global_registry();
let registry2 = global_registry();
assert!(core::ptr::eq(registry1, registry2));
}
#[test]
fn test_spinrwlock_downgrade() {
let lock = SpinRwLock::new(42);
let write_guard = lock.write();
assert_eq!(*write_guard, 42);
let read_guard = write_guard.downgrade();
assert_eq!(*read_guard, 42);
drop(read_guard);
let read_guard2 = lock.read();
assert_eq!(*read_guard2, 42);
}
#[test]
fn test_spinrwlock_downgrade_allows_concurrent_reads() {
use alloc::sync::Arc;
let lock = Arc::new(SpinRwLock::new(100));
let write_guard = lock.write();
let read_guard1 = write_guard.downgrade();
let read_guard2 = lock.read();
assert_eq!(*read_guard1, 100);
assert_eq!(*read_guard2, 100);
}
}