1use std::sync::{Arc, Mutex};
2use std::time::Duration;
3
4use arc_swap::ArcSwap;
5
6#[derive(Clone, Copy, PartialEq, Eq, Debug)]
17pub enum AllocatorStat {
18 Allocated,
20 Resident,
22 Active,
24 Mapped,
26}
27
28impl AllocatorStat {
29 pub fn as_str(&self) -> &'static str {
31 match self {
32 AllocatorStat::Allocated => "allocated",
33 AllocatorStat::Resident => "resident",
34 AllocatorStat::Active => "active",
35 AllocatorStat::Mapped => "mapped",
36 }
37 }
38}
39
40pub trait MetricsCollector: Send + Sync {
43 fn record_exchange_duration(&self, route_id: &str, duration: Duration);
45
46 fn increment_errors(&self, route_id: &str, error_type: &str);
48
49 fn increment_exchanges(&self, route_id: &str);
51
52 fn set_queue_depth(&self, queue: &str, depth: usize);
57
58 fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str);
60
61 fn record_histogram(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {}
64
65 fn record_counter(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {}
68
69 fn increment_retry_attempt(&self, _scheme: &str, _operation: &str) {}
73
74 fn increment_circuit_breaker_rejection(&self, _route: &str) {}
79
80 fn set_route_state(&self, _route: &str, _state: &str) {}
88
89 fn clear_route_state(&self, _route: &str) {}
92
93 fn record_build_info(&self, _version: &str, _git_sha: &str) {}
97
98 fn record_uptime(&self, _seconds: f64) {}
102
103 fn record_component_operation(&self, _component: &str, _operation: &str, _outcome: &str) {}
109
110 fn set_pinned_client_cache_size(&self, _component: &str, _entries: u64) {}
116
117 fn increment_pinned_client_cache_hit(&self, _component: &str) {}
122
123 fn increment_pinned_client_cache_miss(&self, _component: &str) {}
128
129 fn set_allocator_memory(&self, _stat: AllocatorStat, _bytes: u64) {}
134
135 fn set_master_leadership(&self, _lock: &str, _leader: bool) {}
142}
143
144pub struct NoOpMetrics;
146
147impl MetricsCollector for NoOpMetrics {
148 fn record_exchange_duration(&self, _route_id: &str, _duration: Duration) {}
149 fn increment_errors(&self, _route_id: &str, _error_type: &str) {}
150 fn increment_exchanges(&self, _route_id: &str) {}
151 fn set_queue_depth(&self, _queue: &str, _depth: usize) {}
152 fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {}
153}
154
155struct CollectorSlot(Arc<dyn MetricsCollector>);
161
162pub struct MetricsHandle {
178 inner: ArcSwap<CollectorSlot>,
179 members: Mutex<Vec<Arc<dyn MetricsCollector>>>,
183}
184
185impl MetricsHandle {
186 pub fn new() -> Self {
189 Self {
190 inner: ArcSwap::from_pointee(CollectorSlot(Arc::new(NoOpMetrics))),
191 members: Mutex::new(Vec::new()),
192 }
193 }
194
195 pub fn register(&self, collector: Arc<dyn MetricsCollector>) {
200 let mut members = self
201 .members
202 .lock()
203 .expect("metrics members lock poisoned by a panicked register"); if members.iter().any(|m| Arc::ptr_eq(m, &collector)) {
205 return;
206 }
207 let first = members.is_empty();
208 members.push(Arc::clone(&collector));
209 if first {
210 self.inner.store(Arc::new(CollectorSlot(collector)));
213 return;
214 }
215 let prev = Arc::clone(&self.inner.load().0);
216 self.inner.store(Arc::new(CollectorSlot(Arc::new(
217 CompositeMetricsCollector::new(vec![prev, collector]),
218 ))));
219 }
220}
221
222impl Default for MetricsHandle {
223 fn default() -> Self {
224 Self::new()
225 }
226}
227
228impl MetricsCollector for MetricsHandle {
229 fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
230 self.inner
231 .load()
232 .0
233 .record_exchange_duration(route_id, duration)
234 }
235
236 fn increment_errors(&self, route_id: &str, error_type: &str) {
237 self.inner.load().0.increment_errors(route_id, error_type)
238 }
239
240 fn increment_exchanges(&self, route_id: &str) {
241 self.inner.load().0.increment_exchanges(route_id)
242 }
243
244 fn set_queue_depth(&self, queue: &str, depth: usize) {
245 self.inner.load().0.set_queue_depth(queue, depth)
246 }
247
248 fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str) {
249 self.inner
250 .load()
251 .0
252 .record_circuit_breaker_change(route_id, from, to)
253 }
254
255 fn record_histogram(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
256 self.inner.load().0.record_histogram(name, value, labels)
257 }
258
259 fn record_counter(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
260 self.inner.load().0.record_counter(name, value, labels)
261 }
262
263 fn increment_retry_attempt(&self, scheme: &str, operation: &str) {
264 self.inner
265 .load()
266 .0
267 .increment_retry_attempt(scheme, operation)
268 }
269
270 fn increment_circuit_breaker_rejection(&self, route: &str) {
271 self.inner
272 .load()
273 .0
274 .increment_circuit_breaker_rejection(route)
275 }
276
277 fn set_route_state(&self, route: &str, state: &str) {
278 self.inner.load().0.set_route_state(route, state)
279 }
280
281 fn clear_route_state(&self, route: &str) {
282 self.inner.load().0.clear_route_state(route)
283 }
284
285 fn record_build_info(&self, version: &str, git_sha: &str) {
286 self.inner.load().0.record_build_info(version, git_sha)
287 }
288
289 fn record_uptime(&self, seconds: f64) {
290 self.inner.load().0.record_uptime(seconds)
291 }
292
293 fn record_component_operation(&self, component: &str, operation: &str, outcome: &str) {
294 self.inner
295 .load()
296 .0
297 .record_component_operation(component, operation, outcome)
298 }
299
300 fn set_pinned_client_cache_size(&self, component: &str, entries: u64) {
301 self.inner
302 .load()
303 .0
304 .set_pinned_client_cache_size(component, entries)
305 }
306
307 fn increment_pinned_client_cache_hit(&self, component: &str) {
308 self.inner
309 .load()
310 .0
311 .increment_pinned_client_cache_hit(component)
312 }
313
314 fn increment_pinned_client_cache_miss(&self, component: &str) {
315 self.inner
316 .load()
317 .0
318 .increment_pinned_client_cache_miss(component)
319 }
320
321 fn set_allocator_memory(&self, stat: AllocatorStat, bytes: u64) {
322 self.inner.load().0.set_allocator_memory(stat, bytes)
323 }
324
325 fn set_master_leadership(&self, lock: &str, leader: bool) {
326 self.inner.load().0.set_master_leadership(lock, leader)
327 }
328}
329
330#[doc(hidden)]
343pub struct CompositeMetricsCollector {
344 collectors: Vec<Arc<dyn MetricsCollector>>,
345}
346
347impl CompositeMetricsCollector {
348 #[doc(hidden)]
355 pub fn new(collectors: Vec<Arc<dyn MetricsCollector>>) -> Self {
356 Self { collectors }
357 }
358}
359
360impl MetricsCollector for CompositeMetricsCollector {
361 fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
362 for collector in &self.collectors {
363 collector.record_exchange_duration(route_id, duration);
364 }
365 }
366
367 fn increment_errors(&self, route_id: &str, error_type: &str) {
368 for collector in &self.collectors {
369 collector.increment_errors(route_id, error_type);
370 }
371 }
372
373 fn increment_exchanges(&self, route_id: &str) {
374 for collector in &self.collectors {
375 collector.increment_exchanges(route_id);
376 }
377 }
378
379 fn set_queue_depth(&self, queue: &str, depth: usize) {
380 for collector in &self.collectors {
381 collector.set_queue_depth(queue, depth);
382 }
383 }
384
385 fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str) {
386 for collector in &self.collectors {
387 collector.record_circuit_breaker_change(route_id, from, to);
388 }
389 }
390
391 fn record_histogram(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
392 for collector in &self.collectors {
393 collector.record_histogram(name, value, labels);
394 }
395 }
396
397 fn record_counter(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
398 for collector in &self.collectors {
399 collector.record_counter(name, value, labels);
400 }
401 }
402
403 fn increment_retry_attempt(&self, scheme: &str, operation: &str) {
404 for collector in &self.collectors {
405 collector.increment_retry_attempt(scheme, operation);
406 }
407 }
408
409 fn increment_circuit_breaker_rejection(&self, route: &str) {
410 for collector in &self.collectors {
411 collector.increment_circuit_breaker_rejection(route);
412 }
413 }
414
415 fn set_route_state(&self, route: &str, state: &str) {
416 for collector in &self.collectors {
417 collector.set_route_state(route, state);
418 }
419 }
420
421 fn clear_route_state(&self, route: &str) {
422 for collector in &self.collectors {
423 collector.clear_route_state(route);
424 }
425 }
426
427 fn record_build_info(&self, version: &str, git_sha: &str) {
428 for collector in &self.collectors {
429 collector.record_build_info(version, git_sha);
430 }
431 }
432
433 fn record_uptime(&self, seconds: f64) {
434 for collector in &self.collectors {
435 collector.record_uptime(seconds);
436 }
437 }
438
439 fn record_component_operation(&self, component: &str, operation: &str, outcome: &str) {
440 for collector in &self.collectors {
441 collector.record_component_operation(component, operation, outcome);
442 }
443 }
444
445 fn set_pinned_client_cache_size(&self, component: &str, entries: u64) {
446 for collector in &self.collectors {
447 collector.set_pinned_client_cache_size(component, entries);
448 }
449 }
450
451 fn increment_pinned_client_cache_hit(&self, component: &str) {
452 for collector in &self.collectors {
453 collector.increment_pinned_client_cache_hit(component);
454 }
455 }
456
457 fn increment_pinned_client_cache_miss(&self, component: &str) {
458 for collector in &self.collectors {
459 collector.increment_pinned_client_cache_miss(component);
460 }
461 }
462
463 fn set_allocator_memory(&self, stat: AllocatorStat, bytes: u64) {
464 for collector in &self.collectors {
465 collector.set_allocator_memory(stat, bytes);
466 }
467 }
468
469 fn set_master_leadership(&self, lock: &str, leader: bool) {
470 for collector in &self.collectors {
471 collector.set_master_leadership(lock, leader);
472 }
473 }
474}
475
476#[cfg(test)]
477#[path = "metrics_tests.rs"]
478mod tests;