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) {}
}
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)
}
}
#[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);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Mutex};
struct RecordingMetrics {
durations: Mutex<Vec<(String, Duration)>>,
errors: Mutex<Vec<(String, String)>>,
exchanges: Mutex<Vec<String>>,
retries: Mutex<Vec<(String, String)>>,
rejections: Mutex<Vec<String>>,
pinned: Mutex<Vec<(&'static str, String, u64)>>,
allocator: Mutex<Vec<(AllocatorStat, u64)>>,
}
impl RecordingMetrics {
fn new() -> Self {
Self {
durations: Mutex::new(Vec::new()),
errors: Mutex::new(Vec::new()),
exchanges: Mutex::new(Vec::new()),
retries: Mutex::new(Vec::new()),
rejections: Mutex::new(Vec::new()),
pinned: Mutex::new(Vec::new()),
allocator: Mutex::new(Vec::new()),
}
}
}
impl MetricsCollector for RecordingMetrics {
fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
self.durations
.lock()
.expect("durations lock")
.push((route_id.to_string(), duration));
}
fn increment_errors(&self, route_id: &str, error_type: &str) {
self.errors
.lock()
.expect("errors lock")
.push((route_id.to_string(), error_type.to_string()));
}
fn increment_exchanges(&self, route_id: &str) {
self.exchanges
.lock()
.expect("exchanges lock")
.push(route_id.to_string());
}
fn set_queue_depth(&self, _queue: &str, _depth: usize) {}
fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {}
fn increment_retry_attempt(&self, scheme: &str, operation: &str) {
self.retries
.lock()
.expect("retries lock")
.push((scheme.to_string(), operation.to_string()));
}
fn increment_circuit_breaker_rejection(&self, route: &str) {
self.rejections
.lock()
.expect("rejections lock")
.push(route.to_string());
}
fn set_pinned_client_cache_size(&self, component: &str, entries: u64) {
self.pinned.lock().expect("pinned lock").push((
"set_pinned_client_cache_size",
component.to_string(),
entries,
));
}
fn increment_pinned_client_cache_hit(&self, component: &str) {
self.pinned.lock().expect("pinned lock").push((
"increment_pinned_client_cache_hit",
component.to_string(),
1,
));
}
fn increment_pinned_client_cache_miss(&self, component: &str) {
self.pinned.lock().expect("pinned lock").push((
"increment_pinned_client_cache_miss",
component.to_string(),
1,
));
}
fn set_allocator_memory(&self, stat: AllocatorStat, bytes: u64) {
self.allocator
.lock()
.expect("allocator lock")
.push((stat, bytes));
}
}
struct SurfaceProbe {
calls: Mutex<Vec<&'static str>>,
}
impl SurfaceProbe {
fn new() -> Self {
Self {
calls: Mutex::new(Vec::new()),
}
}
fn tag(&self, name: &'static str) {
self.calls.lock().expect("calls lock").push(name);
}
}
impl MetricsCollector for SurfaceProbe {
fn record_exchange_duration(&self, _route_id: &str, _duration: Duration) {
self.tag("record_exchange_duration");
}
fn increment_errors(&self, _route_id: &str, _error_type: &str) {
self.tag("increment_errors");
}
fn increment_exchanges(&self, _route_id: &str) {
self.tag("increment_exchanges");
}
fn set_queue_depth(&self, _queue: &str, _depth: usize) {
self.tag("set_queue_depth");
}
fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {
self.tag("record_circuit_breaker_change");
}
fn record_histogram(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {
self.tag("record_histogram");
}
fn record_counter(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {
self.tag("record_counter");
}
fn increment_retry_attempt(&self, _scheme: &str, _operation: &str) {
self.tag("increment_retry_attempt");
}
fn increment_circuit_breaker_rejection(&self, _route: &str) {
self.tag("increment_circuit_breaker_rejection");
}
fn set_route_state(&self, _route: &str, _state: &str) {
self.tag("set_route_state");
}
fn clear_route_state(&self, _route: &str) {
self.tag("clear_route_state");
}
fn record_build_info(&self, _version: &str, _git_sha: &str) {
self.tag("record_build_info");
}
fn record_uptime(&self, _seconds: f64) {
self.tag("record_uptime");
}
fn record_component_operation(&self, _component: &str, _operation: &str, _outcome: &str) {
self.tag("record_component_operation");
}
}
#[test]
fn test_noop_metrics_implements_trait() {
let metrics = NoOpMetrics;
let metrics_arc: Arc<dyn MetricsCollector> = Arc::new(metrics);
metrics_arc.record_exchange_duration("test-route", Duration::from_millis(100));
metrics_arc.increment_errors("test-route", "test-error");
metrics_arc.increment_exchanges("test-route");
metrics_arc.set_queue_depth("test-route", 5);
metrics_arc.record_circuit_breaker_change("test-route", "closed", "open");
}
#[test]
fn test_custom_metrics_collector() {
struct TestMetrics {
exchange_count: std::sync::atomic::AtomicU64,
}
impl MetricsCollector for TestMetrics {
fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
println!("Route {} took {}ms", route_id, duration.as_millis());
}
fn increment_errors(&self, route_id: &str, error_type: &str) {
println!("Route {} had error: {}", route_id, error_type);
}
fn increment_exchanges(&self, route_id: &str) {
self.exchange_count
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
println!("Route {} processed exchange", route_id);
}
fn set_queue_depth(&self, queue: &str, depth: usize) {
println!("Queue {queue} depth: {depth}");
}
fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str) {
println!("Route {} circuit breaker: {} -> {}", route_id, from, to);
}
}
let test_metrics = TestMetrics {
exchange_count: std::sync::atomic::AtomicU64::new(0),
};
let metrics_arc: Arc<dyn MetricsCollector> = Arc::new(test_metrics);
metrics_arc.record_exchange_duration("test-route", Duration::from_millis(100));
metrics_arc.increment_errors("test-route", "test-error");
metrics_arc.increment_exchanges("test-route");
metrics_arc.set_queue_depth("test-route", 5);
metrics_arc.record_circuit_breaker_change("test-route", "closed", "open");
}
#[test]
fn handle_delegates_to_stored_collector() {
let collector = Arc::new(RecordingMetrics::new());
let handle = MetricsHandle::new();
handle.register(collector.clone());
handle.record_exchange_duration("r", Duration::from_millis(1));
let recorded = collector.durations.lock().expect("durations lock").clone();
assert_eq!(recorded, vec![("r".to_string(), Duration::from_millis(1))]);
}
#[test]
fn second_registration_composes_both_observe() {
let a = Arc::new(RecordingMetrics::new());
let b = Arc::new(RecordingMetrics::new());
let handle = MetricsHandle::new();
handle.register(a.clone());
handle.register(b.clone());
handle.increment_errors("r", "x");
let a_errors = a.errors.lock().expect("errors lock").clone();
let b_errors = b.errors.lock().expect("errors lock").clone();
assert_eq!(a_errors, vec![("r".to_string(), "x".to_string())]);
assert_eq!(b_errors, vec![("r".to_string(), "x".to_string())]);
}
#[test]
fn register_same_arc_is_idempotent() {
let a = Arc::new(RecordingMetrics::new());
let handle = MetricsHandle::new();
handle.register(a.clone());
handle.register(a.clone());
handle.increment_exchanges("r");
let recorded = a.exchanges.lock().expect("exchanges lock").clone();
assert_eq!(recorded, vec!["r".to_string()]);
}
#[test]
fn handle_defaults_to_noop() {
let handle = MetricsHandle::new();
handle.record_exchange_duration("r", Duration::from_millis(1));
handle.increment_errors("r", "x");
handle.increment_exchanges("r");
handle.set_queue_depth("r", 5);
handle.record_circuit_breaker_change("r", "closed", "open");
handle.record_histogram("h", 1.0, &[("k", "v")]);
handle.record_counter("c", 1.0, &[("k", "v")]);
}
#[test]
fn composite_delegates_retry_and_rejection() {
let a = Arc::new(RecordingMetrics::new());
let b = Arc::new(RecordingMetrics::new());
let composite = CompositeMetricsCollector::new(vec![
Arc::clone(&a) as Arc<dyn MetricsCollector>,
Arc::clone(&b) as Arc<dyn MetricsCollector>,
]);
composite.increment_retry_attempt("kafka", "connect");
composite.increment_circuit_breaker_rejection("r1");
for member in [&a, &b] {
assert_eq!(
member.retries.lock().expect("retries lock").clone(),
vec![("kafka".to_string(), "connect".to_string())]
);
assert_eq!(
member.rejections.lock().expect("rejections lock").clone(),
vec!["r1".to_string()]
);
}
}
#[test]
fn noop_defaults_compile_and_do_nothing() {
let collector: Arc<dyn MetricsCollector> = Arc::new(NoOpMetrics);
collector.increment_retry_attempt("kafka", "connect");
collector.increment_circuit_breaker_rejection("r1");
}
#[test]
fn composite_delegates_full_trait_surface() {
let a = Arc::new(SurfaceProbe::new());
let b = Arc::new(SurfaceProbe::new());
let composite = CompositeMetricsCollector::new(vec![
Arc::clone(&a) as Arc<dyn MetricsCollector>,
Arc::clone(&b) as Arc<dyn MetricsCollector>,
]);
composite.record_exchange_duration("r", Duration::from_millis(1));
composite.increment_errors("r", "x");
composite.increment_exchanges("r");
composite.set_queue_depth("r", 1);
composite.record_circuit_breaker_change("r", "closed", "open");
composite.record_histogram("h", 1.0, &[("k", "v")]);
composite.record_counter("c", 1.0, &[("k", "v")]);
composite.increment_retry_attempt("kafka", "connect");
composite.increment_circuit_breaker_rejection("r1");
composite.set_route_state("r", "Started");
composite.clear_route_state("r");
composite.record_build_info("1.2.3", "abc1234");
composite.record_uptime(0.5);
composite.record_component_operation("redis", "command", "success");
let expected = vec![
"record_exchange_duration",
"increment_errors",
"increment_exchanges",
"set_queue_depth",
"record_circuit_breaker_change",
"record_histogram",
"record_counter",
"increment_retry_attempt",
"increment_circuit_breaker_rejection",
"set_route_state",
"clear_route_state",
"record_build_info",
"record_uptime",
"record_component_operation",
];
for member in [&a, &b] {
let calls = member.calls.lock().expect("calls lock").clone();
assert_eq!(calls, expected, "member missed part of the trait surface");
}
}
fn pinned_trio_expected() -> Vec<(&'static str, String, u64)> {
vec![
("set_pinned_client_cache_size", "camel-https".to_string(), 3),
(
"increment_pinned_client_cache_hit",
"camel-https".to_string(),
1,
),
(
"increment_pinned_client_cache_miss",
"camel-https".to_string(),
1,
),
]
}
#[test]
fn handle_forwards_pinned_cache_trio() {
let collector = Arc::new(RecordingMetrics::new());
let handle = MetricsHandle::new();
handle.register(collector.clone());
handle.set_pinned_client_cache_size("camel-https", 3);
handle.increment_pinned_client_cache_hit("camel-https");
handle.increment_pinned_client_cache_miss("camel-https");
let captured = collector.pinned.lock().expect("pinned lock").clone();
assert_eq!(captured, pinned_trio_expected());
let bystander = Arc::new(RecordingMetrics::new());
let unwired = MetricsHandle::new();
unwired.set_pinned_client_cache_size("camel-https", 3);
unwired.increment_pinned_client_cache_hit("camel-https");
unwired.increment_pinned_client_cache_miss("camel-https");
unwired.register(bystander.clone());
assert!(bystander.pinned.lock().expect("pinned lock").is_empty());
}
#[test]
fn composite_forwards_pinned_cache_trio_to_all_collectors() {
let a = Arc::new(RecordingMetrics::new());
let b = Arc::new(RecordingMetrics::new());
let composite = CompositeMetricsCollector::new(vec![
Arc::clone(&a) as Arc<dyn MetricsCollector>,
Arc::clone(&b) as Arc<dyn MetricsCollector>,
]);
composite.set_pinned_client_cache_size("camel-https", 3);
composite.increment_pinned_client_cache_hit("camel-https");
composite.increment_pinned_client_cache_miss("camel-https");
for member in [&a, &b] {
let captured = member.pinned.lock().expect("pinned lock").clone();
assert_eq!(
captured,
pinned_trio_expected(),
"member missed part of the pinned-cache trio"
);
}
}
#[test]
fn allocator_stat_as_str_image_is_closed_set() {
let image: std::collections::BTreeSet<&'static str> = [
AllocatorStat::Allocated,
AllocatorStat::Resident,
AllocatorStat::Active,
AllocatorStat::Mapped,
]
.iter()
.map(|stat| stat.as_str())
.collect();
let expected: std::collections::BTreeSet<&'static str> =
["active", "allocated", "mapped", "resident"]
.into_iter()
.collect();
assert_eq!(image, expected);
}
#[test]
fn handle_and_composite_forward_set_allocator_memory() {
let expected = vec![(AllocatorStat::Resident, 4096)];
let handle_collector = Arc::new(RecordingMetrics::new());
let handle = MetricsHandle::new();
handle.register(handle_collector.clone());
handle.set_allocator_memory(AllocatorStat::Resident, 4096);
assert_eq!(
handle_collector
.allocator
.lock()
.expect("allocator lock")
.clone(),
expected,
"wired handle must forward exactly one allocator emission"
);
let composite_collector = Arc::new(RecordingMetrics::new());
let composite = CompositeMetricsCollector::new(vec![
composite_collector.clone() as Arc<dyn MetricsCollector>
]);
composite.set_allocator_memory(AllocatorStat::Resident, 4096);
assert_eq!(
composite_collector
.allocator
.lock()
.expect("allocator lock")
.clone(),
expected,
"composite must forward exactly one allocator emission"
);
let bystander = Arc::new(RecordingMetrics::new());
let unwired = MetricsHandle::new();
unwired.set_allocator_memory(AllocatorStat::Resident, 4096);
unwired.register(bystander.clone());
assert!(
bystander
.allocator
.lock()
.expect("allocator lock")
.is_empty(),
"unwired-handle emissions are dropped, not replayed"
);
}
}