use std::sync::{Arc, Mutex};
use std::time::Duration;
use arc_swap::ArcSwap;
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum AllocatorStat {
Allocated,
Resident,
Active,
Mapped,
}
impl AllocatorStat {
pub fn as_str(&self) -> &'static str {
match self {
AllocatorStat::Allocated => "allocated",
AllocatorStat::Resident => "resident",
AllocatorStat::Active => "active",
AllocatorStat::Mapped => "mapped",
}
}
}
pub trait MetricsCollector: Send + Sync {
fn record_exchange_duration(&self, route_id: &str, duration: Duration);
fn increment_errors(&self, route_id: &str, error_type: &str);
fn increment_exchanges(&self, route_id: &str);
fn set_queue_depth(&self, queue: &str, depth: usize);
fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str);
fn record_histogram(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {}
fn record_counter(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {}
fn increment_retry_attempt(&self, _scheme: &str, _operation: &str) {}
fn increment_circuit_breaker_rejection(&self, _route: &str) {}
fn set_route_state(&self, _route: &str, _state: &str) {}
fn clear_route_state(&self, _route: &str) {}
fn record_build_info(&self, _version: &str, _git_sha: &str) {}
fn record_uptime(&self, _seconds: f64) {}
fn record_component_operation(&self, _component: &str, _operation: &str, _outcome: &str) {}
fn set_pinned_client_cache_size(&self, _component: &str, _entries: u64) {}
fn increment_pinned_client_cache_hit(&self, _component: &str) {}
fn increment_pinned_client_cache_miss(&self, _component: &str) {}
fn set_allocator_memory(&self, _stat: AllocatorStat, _bytes: u64) {}
fn set_master_leadership(&self, _lock: &str, _leader: bool) {}
}
pub struct NoOpMetrics;
impl MetricsCollector for NoOpMetrics {
fn record_exchange_duration(&self, _route_id: &str, _duration: Duration) {}
fn increment_errors(&self, _route_id: &str, _error_type: &str) {}
fn increment_exchanges(&self, _route_id: &str) {}
fn set_queue_depth(&self, _queue: &str, _depth: usize) {}
fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {}
}
struct CollectorSlot(Arc<dyn MetricsCollector>);
pub struct MetricsHandle {
inner: ArcSwap<CollectorSlot>,
members: Mutex<Vec<Arc<dyn MetricsCollector>>>,
}
impl MetricsHandle {
pub fn new() -> Self {
Self {
inner: ArcSwap::from_pointee(CollectorSlot(Arc::new(NoOpMetrics))),
members: Mutex::new(Vec::new()),
}
}
pub fn register(&self, collector: Arc<dyn MetricsCollector>) {
let mut members = self
.members
.lock()
.expect("metrics members lock poisoned by a panicked register"); if members.iter().any(|m| Arc::ptr_eq(m, &collector)) {
return;
}
let first = members.is_empty();
members.push(Arc::clone(&collector));
if first {
self.inner.store(Arc::new(CollectorSlot(collector)));
return;
}
let prev = Arc::clone(&self.inner.load().0);
self.inner.store(Arc::new(CollectorSlot(Arc::new(
CompositeMetricsCollector::new(vec![prev, collector]),
))));
}
}
impl Default for MetricsHandle {
fn default() -> Self {
Self::new()
}
}
impl MetricsCollector for MetricsHandle {
fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
self.inner
.load()
.0
.record_exchange_duration(route_id, duration)
}
fn increment_errors(&self, route_id: &str, error_type: &str) {
self.inner.load().0.increment_errors(route_id, error_type)
}
fn increment_exchanges(&self, route_id: &str) {
self.inner.load().0.increment_exchanges(route_id)
}
fn set_queue_depth(&self, queue: &str, depth: usize) {
self.inner.load().0.set_queue_depth(queue, depth)
}
fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str) {
self.inner
.load()
.0
.record_circuit_breaker_change(route_id, from, to)
}
fn record_histogram(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
self.inner.load().0.record_histogram(name, value, labels)
}
fn record_counter(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
self.inner.load().0.record_counter(name, value, labels)
}
fn increment_retry_attempt(&self, scheme: &str, operation: &str) {
self.inner
.load()
.0
.increment_retry_attempt(scheme, operation)
}
fn increment_circuit_breaker_rejection(&self, route: &str) {
self.inner
.load()
.0
.increment_circuit_breaker_rejection(route)
}
fn set_route_state(&self, route: &str, state: &str) {
self.inner.load().0.set_route_state(route, state)
}
fn clear_route_state(&self, route: &str) {
self.inner.load().0.clear_route_state(route)
}
fn record_build_info(&self, version: &str, git_sha: &str) {
self.inner.load().0.record_build_info(version, git_sha)
}
fn record_uptime(&self, seconds: f64) {
self.inner.load().0.record_uptime(seconds)
}
fn record_component_operation(&self, component: &str, operation: &str, outcome: &str) {
self.inner
.load()
.0
.record_component_operation(component, operation, outcome)
}
fn set_pinned_client_cache_size(&self, component: &str, entries: u64) {
self.inner
.load()
.0
.set_pinned_client_cache_size(component, entries)
}
fn increment_pinned_client_cache_hit(&self, component: &str) {
self.inner
.load()
.0
.increment_pinned_client_cache_hit(component)
}
fn increment_pinned_client_cache_miss(&self, component: &str) {
self.inner
.load()
.0
.increment_pinned_client_cache_miss(component)
}
fn set_allocator_memory(&self, stat: AllocatorStat, bytes: u64) {
self.inner.load().0.set_allocator_memory(stat, bytes)
}
fn set_master_leadership(&self, lock: &str, leader: bool) {
self.inner.load().0.set_master_leadership(lock, leader)
}
}
#[doc(hidden)]
pub struct CompositeMetricsCollector {
collectors: Vec<Arc<dyn MetricsCollector>>,
}
impl CompositeMetricsCollector {
#[doc(hidden)]
pub fn new(collectors: Vec<Arc<dyn MetricsCollector>>) -> Self {
Self { collectors }
}
}
impl MetricsCollector for CompositeMetricsCollector {
fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
for collector in &self.collectors {
collector.record_exchange_duration(route_id, duration);
}
}
fn increment_errors(&self, route_id: &str, error_type: &str) {
for collector in &self.collectors {
collector.increment_errors(route_id, error_type);
}
}
fn increment_exchanges(&self, route_id: &str) {
for collector in &self.collectors {
collector.increment_exchanges(route_id);
}
}
fn set_queue_depth(&self, queue: &str, depth: usize) {
for collector in &self.collectors {
collector.set_queue_depth(queue, depth);
}
}
fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str) {
for collector in &self.collectors {
collector.record_circuit_breaker_change(route_id, from, to);
}
}
fn record_histogram(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
for collector in &self.collectors {
collector.record_histogram(name, value, labels);
}
}
fn record_counter(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
for collector in &self.collectors {
collector.record_counter(name, value, labels);
}
}
fn increment_retry_attempt(&self, scheme: &str, operation: &str) {
for collector in &self.collectors {
collector.increment_retry_attempt(scheme, operation);
}
}
fn increment_circuit_breaker_rejection(&self, route: &str) {
for collector in &self.collectors {
collector.increment_circuit_breaker_rejection(route);
}
}
fn set_route_state(&self, route: &str, state: &str) {
for collector in &self.collectors {
collector.set_route_state(route, state);
}
}
fn clear_route_state(&self, route: &str) {
for collector in &self.collectors {
collector.clear_route_state(route);
}
}
fn record_build_info(&self, version: &str, git_sha: &str) {
for collector in &self.collectors {
collector.record_build_info(version, git_sha);
}
}
fn record_uptime(&self, seconds: f64) {
for collector in &self.collectors {
collector.record_uptime(seconds);
}
}
fn record_component_operation(&self, component: &str, operation: &str, outcome: &str) {
for collector in &self.collectors {
collector.record_component_operation(component, operation, outcome);
}
}
fn set_pinned_client_cache_size(&self, component: &str, entries: u64) {
for collector in &self.collectors {
collector.set_pinned_client_cache_size(component, entries);
}
}
fn increment_pinned_client_cache_hit(&self, component: &str) {
for collector in &self.collectors {
collector.increment_pinned_client_cache_hit(component);
}
}
fn increment_pinned_client_cache_miss(&self, component: &str) {
for collector in &self.collectors {
collector.increment_pinned_client_cache_miss(component);
}
}
fn set_allocator_memory(&self, stat: AllocatorStat, bytes: u64) {
for collector in &self.collectors {
collector.set_allocator_memory(stat, bytes);
}
}
fn set_master_leadership(&self, lock: &str, leader: bool) {
for collector in &self.collectors {
collector.set_master_leadership(lock, leader);
}
}
}
#[cfg(test)]
#[path = "metrics_tests.rs"]
mod tests;