use crate::core::CacheLayer;
use dashmap::DashMap;
use once_cell::sync::Lazy;
use serde::Serialize;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
#[derive(Clone, Debug, Default)]
pub struct UnifiedMetrics {
inner: Arc<UnifiedMetricsInner>,
}
struct UnifiedMetricsInner {
counters: AtomicCounters,
dynamic_metrics: DashMap<String, MetricValue>,
service_operation_counters: DashMap<String, AtomicU64>,
service_label_overflow: AtomicU64,
config: MetricsConfig,
}
impl std::fmt::Debug for UnifiedMetricsInner {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("UnifiedMetricsInner")
.field("counters", &"<AtomicCounters>")
.field("dynamic_metrics", &self.dynamic_metrics)
.field("config", &self.config)
.finish()
}
}
impl Default for UnifiedMetricsInner {
fn default() -> Self {
Self {
counters: AtomicCounters::default(),
dynamic_metrics: DashMap::new(),
service_operation_counters: DashMap::new(),
service_label_overflow: AtomicU64::new(0),
config: MetricsConfig::default(),
}
}
}
#[derive(Debug)]
pub struct AtomicCounters {
pub l1_hits: AtomicU64,
pub l1_misses: AtomicU64,
pub l2_hits: AtomicU64,
pub l2_misses: AtomicU64,
pub l1_sets: AtomicU64,
pub l2_sets: AtomicU64,
pub l1_deletes: AtomicU64,
pub l2_deletes: AtomicU64,
pub total_operations: AtomicU64,
pub errors: AtomicU64,
pub prefetch_total: AtomicU64,
pub compression_total: AtomicU64,
pub compression_bytes_saved: AtomicU64,
pub l1_items: AtomicU64,
pub l1_capacity_used: AtomicU64,
pub l2_degraded: AtomicU64,
pub l2_retry_total: AtomicU64,
pub backfill_success: AtomicU64,
pub backfill_failed: AtomicU64,
pub evictions: AtomicU64,
}
#[derive(Debug, Clone, Serialize)]
pub enum MetricValue {
Counter(u64),
Gauge(f64),
Histogram(HistogramData),
Timer(TimerData),
Text(String),
}
#[derive(Debug, Clone, Serialize)]
pub struct HistogramData {
pub count: u64,
pub sum: f64,
pub min: f64,
pub max: f64,
pub buckets: Vec<(f64, u64)>,
}
#[derive(Debug, Clone, Serialize)]
pub struct TimerData {
pub count: u64,
pub total_duration: Duration,
pub min_duration: Duration,
pub max_duration: Duration,
}
#[derive(Debug, Clone)]
pub struct MetricsConfig {
pub detailed: bool,
pub histogram_buckets: Vec<f64>,
pub max_dynamic_metrics: usize,
pub max_service_labels: usize,
pub retention_period: Option<Duration>,
}
impl Default for AtomicCounters {
fn default() -> Self {
Self {
l1_hits: AtomicU64::new(0),
l1_misses: AtomicU64::new(0),
l2_hits: AtomicU64::new(0),
l2_misses: AtomicU64::new(0),
l1_sets: AtomicU64::new(0),
l2_sets: AtomicU64::new(0),
l1_deletes: AtomicU64::new(0),
l2_deletes: AtomicU64::new(0),
total_operations: AtomicU64::new(0),
errors: AtomicU64::new(0),
prefetch_total: AtomicU64::new(0),
compression_total: AtomicU64::new(0),
compression_bytes_saved: AtomicU64::new(0),
l1_items: AtomicU64::new(0),
l1_capacity_used: AtomicU64::new(0),
l2_degraded: AtomicU64::new(0),
l2_retry_total: AtomicU64::new(0),
backfill_success: AtomicU64::new(0),
backfill_failed: AtomicU64::new(0),
evictions: AtomicU64::new(0),
}
}
}
impl Default for MetricsConfig {
fn default() -> Self {
Self {
detailed: true,
histogram_buckets: vec![
0.1, 0.5, 1.0, 2.5, 5.0, 10.0, 25.0, 50.0, 100.0, 250.0, 500.0, 1000.0,
],
max_dynamic_metrics: 1000,
max_service_labels: 64,
retention_period: Some(Duration::from_secs(3600)), }
}
}
impl UnifiedMetrics {
pub fn new() -> Self {
Self::with_config(MetricsConfig::default())
}
pub fn with_config(config: MetricsConfig) -> Self {
Self {
inner: Arc::new(UnifiedMetricsInner {
counters: AtomicCounters::default(),
dynamic_metrics: DashMap::new(),
service_operation_counters: DashMap::new(),
service_label_overflow: AtomicU64::new(0),
config,
}),
}
}
pub fn record_operation(&self, operation: CacheOperation) {
match (&operation.layer, &operation.op_type, &operation.result) {
(CacheLayer::L1, CacheOpType::Get, CacheOpResult::Hit) => {
self.inner.counters.l1_hits.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(CacheLayer::L1, CacheOpType::Get, CacheOpResult::Miss) => {
self.inner
.counters
.l1_misses
.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(CacheLayer::L2, CacheOpType::Get, CacheOpResult::Hit) => {
self.inner.counters.l2_hits.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(CacheLayer::L2, CacheOpType::Get, CacheOpResult::Miss) => {
self.inner
.counters
.l2_misses
.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(CacheLayer::L1, CacheOpType::Set, CacheOpResult::Success) => {
self.inner.counters.l1_sets.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(CacheLayer::L2, CacheOpType::Set, CacheOpResult::Success) => {
self.inner.counters.l2_sets.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(CacheLayer::L1, CacheOpType::Delete, CacheOpResult::Success) => {
self.inner
.counters
.l1_deletes
.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(CacheLayer::L2, CacheOpType::Delete, CacheOpResult::Success) => {
self.inner
.counters
.l2_deletes
.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
(_, _, CacheOpResult::Error) => {
self.inner.counters.errors.fetch_add(1, Ordering::Relaxed);
self.inner
.counters
.total_operations
.fetch_add(1, Ordering::Relaxed);
}
_ => {}
}
if self.inner.config.detailed {
let key = format!(
"cache:operations:{}:{}:{}",
operation.layer, operation.op_type, operation.result
);
self.increment_counter(&key, 1);
}
}
pub fn record_duration(&self, operation: &CacheOperation, duration: Duration) {
if !self.inner.config.detailed {
return;
}
let key = format!("cache:duration:{}:{}", operation.layer, operation.op_type);
self.record_timer(&key, duration);
}
pub fn record_custom(&self, key: &str, value: MetricValue) {
if self.inner.dynamic_metrics.len() >= self.inner.config.max_dynamic_metrics {
if let Some(first_key) = self
.inner
.dynamic_metrics
.iter()
.next()
.map(|r| r.key().clone())
{
self.inner.dynamic_metrics.remove(&first_key);
}
}
self.inner.dynamic_metrics.insert(key.to_string(), value);
}
pub fn increment_counter(&self, key: &str, value: u64) {
self.inner
.dynamic_metrics
.entry(key.to_string())
.and_modify(|metric| {
if let MetricValue::Counter(count) = metric {
*count += value;
}
})
.or_insert(MetricValue::Counter(value));
}
pub fn record_operation_with_service(&self, operation: &CacheOperation, service: &str) {
self.record_operation(operation.clone());
if service.is_empty() {
return;
}
let cap = self.inner.config.max_service_labels;
if let Some(count) = self.inner.service_operation_counters.get(service) {
count.fetch_add(1, Ordering::Relaxed);
return;
}
if self.inner.service_operation_counters.len() < cap {
self.inner
.service_operation_counters
.entry(service.to_string())
.and_modify(|count| {
count.fetch_add(1, Ordering::Relaxed);
})
.or_insert_with(|| AtomicU64::new(1));
} else if let Some(count) = self.inner.service_operation_counters.get_mut(service) {
count.fetch_add(1, Ordering::Relaxed);
} else {
self.inner
.service_label_overflow
.fetch_add(1, Ordering::Relaxed);
}
}
pub fn record_backend_operation(&self, backend: &str) {
self.increment_counter(&format!("oxcache_backend_{backend}_operations_total"), 1);
}
pub fn backend_operation_counters(&self) -> Vec<(String, u64)> {
let mut out: Vec<(String, u64)> = self
.inner
.dynamic_metrics
.iter()
.filter_map(|entry| {
let key = entry.key();
if key.starts_with("oxcache_backend_")
&& key.ends_with("_operations_total")
&& let MetricValue::Counter(count) = entry.value()
{
return Some((key.clone(), *count));
}
None
})
.collect();
out.sort_by(|a, b| a.0.cmp(&b.0));
out
}
pub fn service_operation_counters(&self) -> Vec<(String, u64)> {
let mut out: Vec<(String, u64)> = self
.inner
.service_operation_counters
.iter()
.map(|entry| (entry.key().clone(), entry.value().load(Ordering::Relaxed)))
.collect();
out.sort_by(|a, b| a.0.cmp(&b.0));
out
}
pub fn service_label_overflow(&self) -> u64 {
self.inner.service_label_overflow.load(Ordering::Relaxed)
}
pub fn set_gauge(&self, key: &str, value: f64) {
self.inner
.dynamic_metrics
.insert(key.to_string(), MetricValue::Gauge(value));
}
pub fn record_histogram(&self, key: &str, value: f64) {
self.inner
.dynamic_metrics
.entry(key.to_string())
.and_modify(|metric| {
if let MetricValue::Histogram(hist) = metric {
hist.count += 1;
hist.sum += value;
hist.min = hist.min.min(value);
hist.max = hist.max.max(value);
for (boundary, count) in &mut hist.buckets {
if value <= *boundary {
*count += 1;
}
}
}
})
.or_insert_with(|| {
let buckets = self
.inner
.config
.histogram_buckets
.iter()
.map(|&boundary| (boundary, if value <= boundary { 1 } else { 0 }))
.collect();
MetricValue::Histogram(HistogramData {
count: 1,
sum: value,
min: value,
max: value,
buckets,
})
});
}
pub fn record_timer(&self, key: &str, duration: Duration) {
self.inner
.dynamic_metrics
.entry(key.to_string())
.and_modify(|metric| {
if let MetricValue::Timer(timer) = metric {
timer.count += 1;
timer.total_duration += duration;
timer.min_duration = timer.min_duration.min(duration);
timer.max_duration = timer.max_duration.max(duration);
}
})
.or_insert_with(|| {
MetricValue::Timer(TimerData {
count: 1,
total_duration: duration,
min_duration: duration,
max_duration: duration,
})
});
}
pub fn get_counters(&self) -> CounterSnapshot {
CounterSnapshot {
l1_hits: self.inner.counters.l1_hits.load(Ordering::Relaxed),
l1_misses: self.inner.counters.l1_misses.load(Ordering::Relaxed),
l2_hits: self.inner.counters.l2_hits.load(Ordering::Relaxed),
l2_misses: self.inner.counters.l2_misses.load(Ordering::Relaxed),
l1_sets: self.inner.counters.l1_sets.load(Ordering::Relaxed),
l2_sets: self.inner.counters.l2_sets.load(Ordering::Relaxed),
l1_deletes: self.inner.counters.l1_deletes.load(Ordering::Relaxed),
l2_deletes: self.inner.counters.l2_deletes.load(Ordering::Relaxed),
total_operations: self.inner.counters.total_operations.load(Ordering::Relaxed),
errors: self.inner.counters.errors.load(Ordering::Relaxed),
prefetch_total: self.inner.counters.prefetch_total.load(Ordering::Relaxed),
compression_total: self
.inner
.counters
.compression_total
.load(Ordering::Relaxed),
compression_bytes_saved: self
.inner
.counters
.compression_bytes_saved
.load(Ordering::Relaxed),
l1_items: self.inner.counters.l1_items.load(Ordering::Relaxed),
l1_capacity_used: self.inner.counters.l1_capacity_used.load(Ordering::Relaxed),
l2_degraded: self.inner.counters.l2_degraded.load(Ordering::Relaxed),
l2_retry_total: self.inner.counters.l2_retry_total.load(Ordering::Relaxed),
backfill_success: self.inner.counters.backfill_success.load(Ordering::Relaxed),
backfill_failed: self.inner.counters.backfill_failed.load(Ordering::Relaxed),
evictions: self.inner.counters.evictions.load(Ordering::Relaxed),
}
}
pub fn get_dynamic_metrics(&self) -> std::collections::HashMap<String, MetricValue> {
self.inner
.dynamic_metrics
.iter()
.map(|r| (r.key().clone(), r.value().clone()))
.collect()
}
pub fn snapshot(&self) -> MetricsSnapshot {
MetricsSnapshot {
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
counters: self.get_counters(),
dynamic_metrics: self.get_dynamic_metrics(),
service_operations: self.service_operation_counters().into_iter().collect(),
}
}
pub fn reset(&self) {
self.inner.counters.l1_hits.store(0, Ordering::Relaxed);
self.inner.counters.l1_misses.store(0, Ordering::Relaxed);
self.inner.counters.l2_hits.store(0, Ordering::Relaxed);
self.inner.counters.l2_misses.store(0, Ordering::Relaxed);
self.inner.counters.l1_sets.store(0, Ordering::Relaxed);
self.inner.counters.l2_sets.store(0, Ordering::Relaxed);
self.inner.counters.l1_deletes.store(0, Ordering::Relaxed);
self.inner.counters.l2_deletes.store(0, Ordering::Relaxed);
self.inner
.counters
.total_operations
.store(0, Ordering::Relaxed);
self.inner.counters.errors.store(0, Ordering::Relaxed);
self.inner
.counters
.prefetch_total
.store(0, Ordering::Relaxed);
self.inner
.counters
.compression_total
.store(0, Ordering::Relaxed);
self.inner
.counters
.compression_bytes_saved
.store(0, Ordering::Relaxed);
self.inner.counters.l1_items.store(0, Ordering::Relaxed);
self.inner
.counters
.l1_capacity_used
.store(0, Ordering::Relaxed);
self.inner.counters.l2_degraded.store(0, Ordering::Relaxed);
self.inner
.counters
.l2_retry_total
.store(0, Ordering::Relaxed);
self.inner
.counters
.backfill_success
.store(0, Ordering::Relaxed);
self.inner
.counters
.backfill_failed
.store(0, Ordering::Relaxed);
self.inner.counters.evictions.store(0, Ordering::Relaxed);
self.inner.dynamic_metrics.clear();
self.inner.service_operation_counters.clear();
self.inner
.service_label_overflow
.store(0, Ordering::Relaxed);
}
pub fn export_prometheus(&self) -> String {
let snapshot = self.snapshot();
snapshot.export_prometheus()
}
#[cfg(any(feature = "serialization", feature = "full"))]
pub fn export_json(&self) -> Result<String, serde_json::Error> {
let snapshot = self.snapshot();
serde_json::to_string_pretty(&snapshot)
}
pub fn hit_rates(&self) -> HitRates {
let counters = self.get_counters();
let l1_total = counters.l1_hits.saturating_add(counters.l1_misses);
let l2_total = counters.l2_hits.saturating_add(counters.l2_misses);
let total_hits = counters.l1_hits.saturating_add(counters.l2_hits);
let total_misses = counters.l1_misses.saturating_add(counters.l2_misses);
let overall_total = total_hits.saturating_add(total_misses);
HitRates {
l1_hit_rate: if l1_total > 0 {
counters.l1_hits as f64 / l1_total as f64
} else {
0.0
},
l2_hit_rate: if l2_total > 0 {
counters.l2_hits as f64 / l2_total as f64
} else {
0.0
},
overall_hit_rate: if overall_total > 0 {
total_hits as f64 / overall_total as f64
} else {
0.0
},
}
}
pub fn record_l2_retry(&self) {
self.inner
.counters
.l2_retry_total
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_l2_degraded(&self) {
self.inner
.counters
.l2_degraded
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_backfill_success(&self) {
self.inner
.counters
.backfill_success
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_backfill_failed(&self) {
self.inner
.counters
.backfill_failed
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_eviction(&self, count: u64) {
self.inner
.counters
.evictions
.fetch_add(count, Ordering::Relaxed);
}
pub fn eviction_count(&self) -> u64 {
self.inner.counters.evictions.load(Ordering::Relaxed)
}
pub fn record_latency_seconds(&self, seconds: f64) {
let value = if seconds.is_nan() || seconds.is_sign_negative() {
0.0
} else {
seconds
};
self.inner
.dynamic_metrics
.entry(OPERATION_LATENCY_HISTOGRAM.to_string())
.and_modify(|metric| {
if let MetricValue::Histogram(hist) = metric {
hist.count += 1;
hist.sum += value;
hist.min = hist.min.min(value);
hist.max = hist.max.max(value);
for (boundary, count) in &mut hist.buckets {
if value <= *boundary {
*count += 1;
}
}
}
})
.or_insert_with(|| {
let buckets = PROMETHEUS_LATENCY_BUCKETS
.iter()
.map(|&boundary| (boundary, u64::from(value <= boundary)))
.collect();
MetricValue::Histogram(HistogramData {
count: 1,
sum: value,
min: value,
max: value,
buckets,
})
});
}
pub fn export_prometheus_standard(&self) -> String {
let counters = self.get_counters();
let mut out = String::with_capacity(2048);
let counter_line = |out: &mut String, name: &str, help: &str| {
out.push_str(&format!("# HELP {name} {help}.\n"));
out.push_str(&format!("# TYPE {name} counter\n"));
};
counter_line(&mut out, "oxcache_hits_total", "Total cache hits");
out.push_str(&format!(
"oxcache_hits_total{{layer=\"l1\"}} {}\n",
counters.l1_hits
));
out.push_str(&format!(
"oxcache_hits_total{{layer=\"l2\"}} {}\n",
counters.l2_hits
));
counter_line(&mut out, "oxcache_misses_total", "Total cache misses");
out.push_str(&format!(
"oxcache_misses_total{{layer=\"l1\"}} {}\n",
counters.l1_misses
));
out.push_str(&format!(
"oxcache_misses_total{{layer=\"l2\"}} {}\n",
counters.l2_misses
));
counter_line(&mut out, "oxcache_sets_total", "Total cache writes");
out.push_str(&format!(
"oxcache_sets_total{{layer=\"l1\"}} {}\n",
counters.l1_sets
));
out.push_str(&format!(
"oxcache_sets_total{{layer=\"l2\"}} {}\n",
counters.l2_sets
));
counter_line(&mut out, "oxcache_deletes_total", "Total cache deletions");
out.push_str(&format!(
"oxcache_deletes_total{{layer=\"l1\"}} {}\n",
counters.l1_deletes
));
out.push_str(&format!(
"oxcache_deletes_total{{layer=\"l2\"}} {}\n",
counters.l2_deletes
));
counter_line(&mut out, "oxcache_evictions_total", "Total cache evictions");
out.push_str(&format!("oxcache_evictions_total {}\n", counters.evictions));
counter_line(
&mut out,
"oxcache_operations_total",
"Total cache operations",
);
out.push_str(&format!(
"oxcache_operations_total {}\n",
counters.total_operations
));
counter_line(&mut out, "oxcache_errors_total", "Total cache errors");
out.push_str(&format!("oxcache_errors_total {}\n", counters.errors));
for (key, count) in &self.backend_operation_counters() {
let backend = key
.trim_start_matches("oxcache_backend_")
.trim_end_matches("_operations_total");
out.push_str(&format!("# HELP {key} Operations on backend {backend}.\n"));
out.push_str(&format!("# TYPE {key} counter\n"));
out.push_str(&format!(
"{key}{{backend=\"{}\"}} {count}\n",
escape_prometheus_label_value(backend)
));
}
let service_counters = self.service_operation_counters();
if !service_counters.is_empty() {
out.push_str(
"# HELP oxcache_service_operations_total Operations attributed to service labels.\n",
);
out.push_str("# TYPE oxcache_service_operations_total counter\n");
for (service, count) in &service_counters {
out.push_str(&format!(
"oxcache_service_operations_total{{service=\"{}\"}} {count}\n",
escape_prometheus_label_value(service)
));
}
}
let overflow = self.service_label_overflow();
if overflow > 0 {
out.push_str(
"# HELP oxcache_service_labels_overflow_total Operations dropped after the service label cardinality cap.\n",
);
out.push_str("# TYPE oxcache_service_labels_overflow_total counter\n");
out.push_str(&format!(
"oxcache_service_labels_overflow_total {overflow}\n"
));
}
out.push_str(
"# HELP oxcache_operation_duration_seconds Cache operation latency in seconds.\n",
);
out.push_str("# TYPE oxcache_operation_duration_seconds histogram\n");
if let Some(MetricValue::Histogram(hist)) = self
.inner
.dynamic_metrics
.get(OPERATION_LATENCY_HISTOGRAM)
.map(|r| r.value().clone())
{
for (boundary, count) in &hist.buckets {
out.push_str(&format!(
"oxcache_operation_duration_seconds_bucket{{le=\"{boundary}\"}} {count}\n"
));
}
out.push_str(&format!(
"oxcache_operation_duration_seconds_bucket{{le=\"+Inf\"}} {}\n",
hist.count
));
out.push_str(&format!(
"oxcache_operation_duration_seconds_sum {}\n",
hist.sum
));
out.push_str(&format!(
"oxcache_operation_duration_seconds_count {}\n",
hist.count
));
} else {
out.push_str("oxcache_operation_duration_seconds_bucket{le=\"+Inf\"} 0\n");
out.push_str("oxcache_operation_duration_seconds_sum 0\n");
out.push_str("oxcache_operation_duration_seconds_count 0\n");
}
out
}
}
pub const OPERATION_LATENCY_HISTOGRAM: &str = "oxcache_operation_duration_seconds";
pub(crate) fn escape_prometheus_label_value(value: &str) -> String {
let mut out = String::with_capacity(value.len());
for ch in value.chars() {
match ch {
'\\' => out.push_str("\\\\"),
'"' => out.push_str("\\\""),
'\n' => out.push_str("\\n"),
'\r' => out.push_str("\\r"),
other => out.push(other),
}
}
out
}
pub const PROMETHEUS_LATENCY_BUCKETS: &[f64] = &[
0.000_1, 0.000_5, 0.001, 0.005, 0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0,
];
#[derive(Debug, Clone)]
pub struct CacheOperation {
pub layer: CacheLayer,
pub op_type: CacheOpType,
pub result: CacheOpResult,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CacheOpType {
Get,
Set,
Delete,
Clear,
}
impl std::fmt::Display for CacheOpType {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
CacheOpType::Get => write!(f, "Get"),
CacheOpType::Set => write!(f, "Set"),
CacheOpType::Delete => write!(f, "Delete"),
CacheOpType::Clear => write!(f, "Clear"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CacheOpResult {
Hit,
Miss,
Success,
Error,
}
impl std::fmt::Display for CacheOpResult {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
CacheOpResult::Hit => write!(f, "Hit"),
CacheOpResult::Miss => write!(f, "Miss"),
CacheOpResult::Success => write!(f, "Success"),
CacheOpResult::Error => write!(f, "Error"),
}
}
}
#[derive(Debug, Clone, Serialize)]
pub struct CounterSnapshot {
pub l1_hits: u64,
pub l1_misses: u64,
pub l2_hits: u64,
pub l2_misses: u64,
pub l1_sets: u64,
pub l2_sets: u64,
pub l1_deletes: u64,
pub l2_deletes: u64,
pub total_operations: u64,
pub errors: u64,
pub prefetch_total: u64,
pub compression_total: u64,
pub compression_bytes_saved: u64,
pub l1_items: u64,
pub l1_capacity_used: u64,
pub l2_degraded: u64,
pub l2_retry_total: u64,
pub backfill_success: u64,
pub backfill_failed: u64,
pub evictions: u64,
}
#[derive(Debug, Clone, Serialize)]
pub struct MetricsSnapshot {
pub timestamp: u64,
pub counters: CounterSnapshot,
pub dynamic_metrics: std::collections::HashMap<String, MetricValue>,
#[serde(skip_serializing_if = "std::collections::HashMap::is_empty")]
pub service_operations: std::collections::HashMap<String, u64>,
}
#[derive(Debug, Clone)]
pub struct HitRates {
pub l1_hit_rate: f64,
pub l2_hit_rate: f64,
pub overall_hit_rate: f64,
}
impl MetricsSnapshot {
pub fn export_prometheus(&self) -> String {
let mut output = String::new();
output.push_str("# Cache Metrics Snapshot\n");
output.push_str(&format!("# Generated at: {}\n", self.timestamp));
output.push_str(&format!("cache_l1_hits_total {}\n", self.counters.l1_hits));
output.push_str(&format!(
"cache_l1_misses_total {}\n",
self.counters.l1_misses
));
output.push_str(&format!("cache_l2_hits_total {}\n", self.counters.l2_hits));
output.push_str(&format!(
"cache_l2_misses_total {}\n",
self.counters.l2_misses
));
output.push_str(&format!("cache_l1_sets_total {}\n", self.counters.l1_sets));
output.push_str(&format!("cache_l2_sets_total {}\n", self.counters.l2_sets));
output.push_str(&format!(
"cache_l1_deletes_total {}\n",
self.counters.l1_deletes
));
output.push_str(&format!(
"cache_l2_deletes_total {}\n",
self.counters.l2_deletes
));
output.push_str(&format!(
"cache_operations_total {}\n",
self.counters.total_operations
));
output.push_str(&format!("cache_errors_total {}\n", self.counters.errors));
output.push_str(&format!(
"cache_evictions_total {}\n",
self.counters.evictions
));
output.push_str(&format!("cache_l1_items {}\n", self.counters.l1_items));
output.push_str(&format!(
"cache_l1_capacity_used_bytes {}\n",
self.counters.l1_capacity_used
));
output.push_str(&format!(
"cache_l2_degraded_total {}\n",
self.counters.l2_degraded
));
output.push_str(&format!(
"cache_l2_retry_total {}\n",
self.counters.l2_retry_total
));
output.push_str(&format!(
"cache_backfill_success_total {}\n",
self.counters.backfill_success
));
output.push_str(&format!(
"cache_backfill_failed_total {}\n",
self.counters.backfill_failed
));
for (key, value) in &self.dynamic_metrics {
match value {
MetricValue::Counter(count) => {
output.push_str(&format!("{}_counter {}\n", key, count));
}
MetricValue::Gauge(value) => {
output.push_str(&format!("{}_gauge {}\n", key, value));
}
MetricValue::Histogram(hist) => {
output.push_str(&format!("{}_histogram_sum {}\n", key, hist.sum));
output.push_str(&format!("{}_histogram_count {}\n", key, hist.count));
for (boundary, count) in &hist.buckets {
output.push_str(&format!(
"{}_histogram_bucket{{le=\"{}\"}} {}\n",
key, boundary, count
));
}
}
MetricValue::Timer(timer) => {
output.push_str(&format!(
"{}_timer_seconds_sum {}\n",
key,
timer.total_duration.as_secs_f64()
));
output.push_str(&format!("{}_timer_seconds_count {}\n", key, timer.count));
}
MetricValue::Text(text) => {
output.push_str(&format!("{}_info \"{}\"\n", key, text));
}
}
}
output
}
}
pub static GLOBAL_UNIFIED_METRICS: Lazy<UnifiedMetrics> = Lazy::new(UnifiedMetrics::new);
pub mod convenience {
use super::*;
pub fn record_operation(operation: CacheOperation) {
GLOBAL_UNIFIED_METRICS.record_operation(operation);
}
pub fn hit_rates() -> HitRates {
GLOBAL_UNIFIED_METRICS.hit_rates()
}
pub fn export_prometheus() -> String {
GLOBAL_UNIFIED_METRICS.export_prometheus()
}
#[cfg(any(feature = "serialization", feature = "full"))]
pub fn export_json() -> Result<String, serde_json::Error> {
GLOBAL_UNIFIED_METRICS.export_json()
}
pub fn reset() {
GLOBAL_UNIFIED_METRICS.reset();
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_atomic_counters() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
});
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Miss,
});
let counters = metrics.get_counters();
assert_eq!(counters.l1_hits, 1);
assert_eq!(counters.l1_misses, 1);
assert_eq!(counters.total_operations, 2);
}
#[test]
fn test_hit_rates() {
let metrics = UnifiedMetrics::new();
for _ in 0..7 {
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
});
}
for _ in 0..3 {
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Miss,
});
}
let hit_rates = metrics.hit_rates();
assert_eq!(hit_rates.l1_hit_rate, 0.7);
assert_eq!(hit_rates.l2_hit_rate, 0.0);
assert_eq!(hit_rates.overall_hit_rate, 0.7);
}
#[test]
fn test_dynamic_metrics() {
let metrics = UnifiedMetrics::new();
metrics.increment_counter("test_counter", 5);
metrics.increment_counter("test_counter", 3);
metrics.set_gauge("test_gauge", 42.5);
metrics.record_histogram("test_histogram", 1.5);
metrics.record_histogram("test_histogram", 2.7);
metrics.record_timer("test_timer", Duration::from_millis(100));
metrics.record_timer("test_timer", Duration::from_millis(200));
let dynamic_metrics = metrics.get_dynamic_metrics();
let counter_metric = dynamic_metrics.get("test_counter");
assert!(
matches!(counter_metric, Some(MetricValue::Counter(count)) if *count == 8),
"Expected counter metric with value 8, got {:?}",
counter_metric
);
let gauge_metric = dynamic_metrics.get("test_gauge");
assert!(
matches!(gauge_metric, Some(MetricValue::Gauge(value)) if *value == 42.5),
"Expected gauge metric with value 42.5, got {:?}",
gauge_metric
);
}
#[test]
fn test_snapshot() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Set,
result: CacheOpResult::Success,
});
let snapshot = metrics.snapshot();
assert_eq!(snapshot.counters.l1_sets, 1);
assert!(
snapshot.dynamic_metrics.is_empty()
|| !snapshot.dynamic_metrics.contains_key("detailed")
);
let prometheus = snapshot.export_prometheus();
assert!(prometheus.contains("cache_l1_sets_total 1"));
let json = serde_json::to_string_pretty(&snapshot).unwrap();
assert!(json.contains("l1_sets"));
}
#[test]
#[serial_test::serial]
fn test_global_metrics() {
convenience::reset();
convenience::record_operation(CacheOperation {
layer: CacheLayer::L2,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
});
let hit_rates = convenience::hit_rates();
assert_eq!(hit_rates.l2_hit_rate, 1.0);
let prometheus = convenience::export_prometheus();
assert!(prometheus.contains("cache_l2_hits_total 1"));
}
#[test]
fn test_unified_metrics_default_impl() {
let metrics = UnifiedMetrics::default();
let counters = metrics.get_counters();
assert_eq!(counters.l1_hits, 0);
assert_eq!(counters.total_operations, 0);
}
#[test]
fn test_unified_metrics_debug_format() {
let metrics = UnifiedMetrics::new();
let debug_str = format!("{:?}", metrics);
assert!(debug_str.contains("UnifiedMetrics"));
assert!(debug_str.contains("UnifiedMetricsInner"));
assert!(debug_str.contains("dynamic_metrics"));
assert!(debug_str.contains("config"));
}
#[test]
fn test_l2_get_miss_operation() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L2,
op_type: CacheOpType::Get,
result: CacheOpResult::Miss,
});
let counters = metrics.get_counters();
assert_eq!(counters.l2_misses, 1);
assert_eq!(counters.total_operations, 1);
}
#[test]
fn test_l2_set_success_operation() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L2,
op_type: CacheOpType::Set,
result: CacheOpResult::Success,
});
let counters = metrics.get_counters();
assert_eq!(counters.l2_sets, 1);
assert_eq!(counters.total_operations, 1);
}
#[test]
fn test_l1_delete_success_operation() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Delete,
result: CacheOpResult::Success,
});
let counters = metrics.get_counters();
assert_eq!(counters.l1_deletes, 1);
assert_eq!(counters.total_operations, 1);
}
#[test]
fn test_l2_delete_success_operation() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L2,
op_type: CacheOpType::Delete,
result: CacheOpResult::Success,
});
let counters = metrics.get_counters();
assert_eq!(counters.l2_deletes, 1);
assert_eq!(counters.total_operations, 1);
}
#[test]
fn test_error_result_operation() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Error,
});
metrics.record_operation(CacheOperation {
layer: CacheLayer::L2,
op_type: CacheOpType::Set,
result: CacheOpResult::Error,
});
let counters = metrics.get_counters();
assert_eq!(counters.errors, 2);
assert_eq!(counters.total_operations, 2);
}
#[test]
fn test_record_duration_with_detailed_disabled() {
let config = MetricsConfig {
detailed: false,
..MetricsConfig::default()
};
let metrics = UnifiedMetrics::with_config(config);
let op = CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
};
metrics.record_duration(&op, Duration::from_millis(100));
let dynamic = metrics.get_dynamic_metrics();
assert!(!dynamic.keys().any(|k| k.contains("duration")));
}
#[test]
fn test_record_duration_with_detailed_enabled() {
let metrics = UnifiedMetrics::new();
let op = CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
};
metrics.record_duration(&op, Duration::from_millis(150));
let dynamic = metrics.get_dynamic_metrics();
let timer_key = "cache:duration:L1:Get";
assert!(
dynamic.contains_key(timer_key),
"Expected timer metric at key {}",
timer_key
);
match dynamic.get(timer_key) {
Some(MetricValue::Timer(t)) => {
assert_eq!(t.count, 1);
assert_eq!(t.total_duration, Duration::from_millis(150));
}
other => panic!("Expected Timer, got {:?}", other),
}
}
#[test]
fn test_record_custom_at_capacity_check() {
let config = MetricsConfig {
max_dynamic_metrics: 0,
..MetricsConfig::default()
};
let metrics = UnifiedMetrics::with_config(config);
metrics.record_custom("only_metric", MetricValue::Counter(1));
let dynamic = metrics.get_dynamic_metrics();
assert!(
dynamic.contains_key("only_metric"),
"metric should be present after insert"
);
assert_eq!(dynamic.len(), 1);
}
#[test]
fn test_record_custom_under_capacity_no_eviction() {
let config = MetricsConfig {
max_dynamic_metrics: 10,
..MetricsConfig::default()
};
let metrics = UnifiedMetrics::with_config(config);
metrics.record_custom("a", MetricValue::Counter(1));
metrics.record_custom("b", MetricValue::Counter(2));
let dynamic = metrics.get_dynamic_metrics();
assert_eq!(dynamic.len(), 2);
assert!(dynamic.contains_key("a"));
assert!(dynamic.contains_key("b"));
}
#[test]
fn test_export_json_method() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
});
let json = metrics.export_json().unwrap();
assert!(json.contains("counters"));
assert!(json.contains("l1_hits"));
assert!(json.contains("dynamic_metrics"));
}
#[test]
fn test_cache_op_type_clear_display() {
assert_eq!(CacheOpType::Clear.to_string(), "Clear");
}
#[test]
fn test_cache_op_type_all_display_variants() {
assert_eq!(CacheOpType::Get.to_string(), "Get");
assert_eq!(CacheOpType::Set.to_string(), "Set");
assert_eq!(CacheOpType::Delete.to_string(), "Delete");
assert_eq!(CacheOpType::Clear.to_string(), "Clear");
}
#[test]
fn test_cache_op_result_error_display() {
assert_eq!(CacheOpResult::Error.to_string(), "Error");
}
#[test]
fn test_cache_op_result_all_display_variants() {
assert_eq!(CacheOpResult::Hit.to_string(), "Hit");
assert_eq!(CacheOpResult::Miss.to_string(), "Miss");
assert_eq!(CacheOpResult::Success.to_string(), "Success");
assert_eq!(CacheOpResult::Error.to_string(), "Error");
}
#[test]
fn test_export_prometheus_gauge_metric() {
let metrics = UnifiedMetrics::new();
metrics.set_gauge("my_gauge", 42.5);
let prom = metrics.export_prometheus();
assert!(prom.contains("my_gauge_gauge 42.5"));
}
#[test]
fn test_export_prometheus_histogram_metric() {
let metrics = UnifiedMetrics::new();
metrics.record_histogram("my_histogram", 1.5);
metrics.record_histogram("my_histogram", 2.5);
let prom = metrics.export_prometheus();
assert!(prom.contains("my_histogram_histogram_sum 4"));
assert!(prom.contains("my_histogram_histogram_count 2"));
assert!(prom.contains("my_histogram_histogram_bucket"));
assert!(prom.contains("le=\""));
}
#[test]
fn test_export_prometheus_timer_metric() {
let metrics = UnifiedMetrics::new();
metrics.record_timer("my_timer", Duration::from_millis(250));
let prom = metrics.export_prometheus();
assert!(prom.contains("my_timer_timer_seconds_sum"));
assert!(prom.contains("my_timer_timer_seconds_count 1"));
}
#[test]
fn test_export_prometheus_text_metric() {
let metrics = UnifiedMetrics::new();
metrics.record_custom("my_info", MetricValue::Text("version1".to_string()));
let prom = metrics.export_prometheus();
assert!(prom.contains("my_info_info \"version1\""));
}
#[test]
fn test_export_prometheus_counter_dynamic_metric() {
let metrics = UnifiedMetrics::new();
metrics.increment_counter("my_dyn_counter", 7);
let prom = metrics.export_prometheus();
assert!(prom.contains("my_dyn_counter_counter 7"));
}
#[test]
#[serial_test::serial]
fn test_convenience_export_json() {
convenience::reset();
let json = convenience::export_json().unwrap();
assert!(json.contains("counters"));
assert!(json.contains("l1_hits"));
assert!(json.contains("dynamic_metrics"));
}
#[test]
fn test_metrics_config_default() {
let config = MetricsConfig::default();
assert!(config.detailed);
assert!(!config.histogram_buckets.is_empty());
assert_eq!(config.max_dynamic_metrics, 1000);
assert_eq!(config.retention_period, Some(Duration::from_secs(3600)));
}
#[test]
fn test_atomic_counters_default() {
let counters = AtomicCounters::default();
assert_eq!(counters.l1_hits.load(Ordering::Relaxed), 0);
assert_eq!(counters.l1_misses.load(Ordering::Relaxed), 0);
assert_eq!(counters.errors.load(Ordering::Relaxed), 0);
assert_eq!(counters.prefetch_total.load(Ordering::Relaxed), 0);
}
#[test]
fn test_reset_clears_all_metrics() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
});
metrics.increment_counter("custom", 5);
assert!(!metrics.get_dynamic_metrics().is_empty());
metrics.reset();
let counters = metrics.get_counters();
assert_eq!(counters.l1_hits, 0);
assert_eq!(counters.total_operations, 0);
assert!(metrics.get_dynamic_metrics().is_empty());
}
#[test]
fn test_hit_rates_overall_default_is_zero() {
let metrics = UnifiedMetrics::new();
let rates = metrics.hit_rates();
assert_eq!(rates.l1_hit_rate, 0.0);
assert_eq!(rates.l2_hit_rate, 0.0);
assert_eq!(rates.overall_hit_rate, 0.0);
}
fn l1_hit() -> CacheOperation {
CacheOperation {
layer: CacheLayer::L1,
op_type: CacheOpType::Get,
result: CacheOpResult::Hit,
}
}
#[test]
fn service_dimension_records_per_service_and_exports_prometheus() {
let metrics = UnifiedMetrics::new();
metrics.record_operation_with_service(&l1_hit(), "svc-a");
metrics.record_operation_with_service(&l1_hit(), "svc-a");
metrics.record_operation_with_service(&l1_hit(), "svc-b");
assert_eq!(metrics.get_counters().l1_hits, 3);
let export = metrics.export_prometheus_standard();
assert!(
export.contains("# HELP oxcache_service_operations_total"),
"HELP line missing: {export}"
);
assert!(
export.contains("# TYPE oxcache_service_operations_total counter"),
"TYPE line missing: {export}"
);
assert!(
export.contains("oxcache_service_operations_total{service=\"svc-a\"} 2"),
"svc-a line missing: {export}"
);
assert!(
export.contains("oxcache_service_operations_total{service=\"svc-b\"} 1"),
"svc-b line missing: {export}"
);
assert!(!export.contains("service_labels_overflow"));
}
#[test]
fn service_dimension_default_off() {
let metrics = UnifiedMetrics::new();
metrics.record_operation(l1_hit());
assert_eq!(metrics.get_counters().l1_hits, 1);
let export = metrics.export_prometheus_standard();
assert!(
!export.contains("oxcache_service_operations_total"),
"service lines must be absent by default: {export}"
);
assert!(metrics.service_operation_counters().is_empty());
#[cfg(feature = "serialization")]
{
let json = metrics.export_json().unwrap();
assert!(
!json.contains("\"service_operations\""),
"default snapshot JSON must omit service_operations: {json}"
);
}
}
#[test]
fn service_dimension_cap_prevents_label_explosion() {
let config = MetricsConfig {
max_service_labels: 2,
..MetricsConfig::default()
};
let metrics = UnifiedMetrics::with_config(config);
metrics.record_operation_with_service(&l1_hit(), "svc-1");
metrics.record_operation_with_service(&l1_hit(), "svc-2");
metrics.record_operation_with_service(&l1_hit(), "svc-3");
metrics.record_operation_with_service(&l1_hit(), "svc-3");
let counters = metrics.service_operation_counters();
assert_eq!(counters.len(), 2, "cap must hold: {counters:?}");
assert!(counters.iter().all(|(k, _)| k != "svc-3"));
let export = metrics.export_prometheus_standard();
assert!(
!export.contains("svc-3"),
"over-cap service must not appear: {export}"
);
assert!(
export.contains("oxcache_service_labels_overflow_total 2"),
"overflow counter must reflect dropped attributions: {export}"
);
assert_eq!(metrics.get_counters().l1_hits, 4);
}
#[cfg(feature = "serialization")]
#[test]
fn service_dimension_visible_in_json() {
let metrics = UnifiedMetrics::new();
metrics.record_operation_with_service(&l1_hit(), "orders");
metrics.record_operation_with_service(&l1_hit(), "orders");
let snapshot = metrics.snapshot();
assert_eq!(
snapshot.service_operations.get("orders"),
Some(&2),
"snapshot must carry per-service counters"
);
let json = metrics.export_json().unwrap();
assert!(json.contains("\"service_operations\""), "json key: {json}");
assert!(json.contains("\"orders\""), "json service name: {json}");
}
#[test]
fn service_dimension_reset_clears_counters() {
let metrics = UnifiedMetrics::new();
metrics.record_operation_with_service(&l1_hit(), "svc-reset");
metrics.reset();
assert!(metrics.service_operation_counters().is_empty());
let export = metrics.export_prometheus_standard();
assert!(!export.contains("svc-reset"));
}
#[test]
fn service_label_value_is_escaped_in_prometheus_export() {
let metrics = UnifiedMetrics::new();
let hostile = "svc\"quote\\back\nnewline";
metrics.record_operation_with_service(&l1_hit(), hostile);
metrics.record_operation_with_service(&l1_hit(), hostile);
let export = metrics.export_prometheus_standard();
let expected_line =
"oxcache_service_operations_total{service=\"svc\\\"quote\\\\back\\nnewline\"} 2";
assert!(
export.contains(expected_line),
"escaped service line missing: {export}"
);
assert_eq!(
export
.matches("oxcache_service_operations_total{service=")
.count(),
1,
"hostile label must not forge extra series: {export}"
);
assert!(!export.contains("svc\"quote"));
}
#[test]
fn service_dimension_cap_zero_keeps_overflow_visible() {
let config = MetricsConfig {
max_service_labels: 0,
..MetricsConfig::default()
};
let metrics = UnifiedMetrics::with_config(config);
metrics.record_operation_with_service(&l1_hit(), "svc-any");
metrics.record_operation_with_service(&l1_hit(), "svc-any");
assert!(metrics.service_operation_counters().is_empty());
assert_eq!(metrics.service_label_overflow(), 2);
let export = metrics.export_prometheus_standard();
assert!(
export.contains("oxcache_service_labels_overflow_total 2"),
"overflow must be exported even with empty service map: {export}"
);
assert_eq!(metrics.get_counters().l1_hits, 2);
}
#[test]
fn service_dimension_concurrent_first_touch_has_no_lost_counts() {
use std::sync::Arc;
let metrics = Arc::new(UnifiedMetrics::new());
let handles: Vec<_> = (0..8)
.map(|_| {
let metrics = Arc::clone(&metrics);
std::thread::spawn(move || {
for _ in 0..250 {
metrics.record_operation_with_service(&l1_hit(), "svc-race");
}
})
})
.collect();
for handle in handles {
handle.join().unwrap();
}
let counters = metrics.service_operation_counters();
assert_eq!(
counters.first().map(|(_, c)| *c),
Some(2_000),
"concurrent first-touch must not lose counts: {counters:?}"
);
}
}