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)]
477mod tests {
478 use super::*;
479 use std::sync::{Arc, Mutex};
480
481 struct RecordingMetrics {
483 durations: Mutex<Vec<(String, Duration)>>,
484 errors: Mutex<Vec<(String, String)>>,
485 exchanges: Mutex<Vec<String>>,
486 retries: Mutex<Vec<(String, String)>>,
487 rejections: Mutex<Vec<String>>,
488 pinned: Mutex<Vec<(&'static str, String, u64)>>,
489 allocator: Mutex<Vec<(AllocatorStat, u64)>>,
490 leadership: Mutex<Vec<(String, bool)>>,
491 }
492
493 impl RecordingMetrics {
494 fn new() -> Self {
495 Self {
496 durations: Mutex::new(Vec::new()),
497 errors: Mutex::new(Vec::new()),
498 exchanges: Mutex::new(Vec::new()),
499 retries: Mutex::new(Vec::new()),
500 rejections: Mutex::new(Vec::new()),
501 pinned: Mutex::new(Vec::new()),
502 allocator: Mutex::new(Vec::new()),
503 leadership: Mutex::new(Vec::new()),
504 }
505 }
506 }
507
508 impl MetricsCollector for RecordingMetrics {
509 fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
510 self.durations
511 .lock()
512 .expect("durations lock")
513 .push((route_id.to_string(), duration));
514 }
515
516 fn increment_errors(&self, route_id: &str, error_type: &str) {
517 self.errors
518 .lock()
519 .expect("errors lock")
520 .push((route_id.to_string(), error_type.to_string()));
521 }
522
523 fn increment_exchanges(&self, route_id: &str) {
524 self.exchanges
525 .lock()
526 .expect("exchanges lock")
527 .push(route_id.to_string());
528 }
529
530 fn set_queue_depth(&self, _queue: &str, _depth: usize) {}
531
532 fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {}
533
534 fn increment_retry_attempt(&self, scheme: &str, operation: &str) {
535 self.retries
536 .lock()
537 .expect("retries lock")
538 .push((scheme.to_string(), operation.to_string()));
539 }
540
541 fn increment_circuit_breaker_rejection(&self, route: &str) {
542 self.rejections
543 .lock()
544 .expect("rejections lock")
545 .push(route.to_string());
546 }
547
548 fn set_pinned_client_cache_size(&self, component: &str, entries: u64) {
549 self.pinned.lock().expect("pinned lock").push((
550 "set_pinned_client_cache_size",
551 component.to_string(),
552 entries,
553 ));
554 }
555
556 fn increment_pinned_client_cache_hit(&self, component: &str) {
557 self.pinned.lock().expect("pinned lock").push((
558 "increment_pinned_client_cache_hit",
559 component.to_string(),
560 1,
561 ));
562 }
563
564 fn increment_pinned_client_cache_miss(&self, component: &str) {
565 self.pinned.lock().expect("pinned lock").push((
566 "increment_pinned_client_cache_miss",
567 component.to_string(),
568 1,
569 ));
570 }
571
572 fn set_allocator_memory(&self, stat: AllocatorStat, bytes: u64) {
573 self.allocator
574 .lock()
575 .expect("allocator lock")
576 .push((stat, bytes));
577 }
578
579 fn set_master_leadership(&self, lock: &str, leader: bool) {
580 self.leadership
581 .lock()
582 .expect("leadership lock")
583 .push((lock.to_string(), leader));
584 }
585 }
586
587 struct SurfaceProbe {
590 calls: Mutex<Vec<&'static str>>,
591 }
592
593 impl SurfaceProbe {
594 fn new() -> Self {
595 Self {
596 calls: Mutex::new(Vec::new()),
597 }
598 }
599
600 fn tag(&self, name: &'static str) {
601 self.calls.lock().expect("calls lock").push(name);
602 }
603 }
604
605 impl MetricsCollector for SurfaceProbe {
606 fn record_exchange_duration(&self, _route_id: &str, _duration: Duration) {
607 self.tag("record_exchange_duration");
608 }
609 fn increment_errors(&self, _route_id: &str, _error_type: &str) {
610 self.tag("increment_errors");
611 }
612 fn increment_exchanges(&self, _route_id: &str) {
613 self.tag("increment_exchanges");
614 }
615 fn set_queue_depth(&self, _queue: &str, _depth: usize) {
616 self.tag("set_queue_depth");
617 }
618 fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {
619 self.tag("record_circuit_breaker_change");
620 }
621 fn record_histogram(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {
622 self.tag("record_histogram");
623 }
624 fn record_counter(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {
625 self.tag("record_counter");
626 }
627 fn increment_retry_attempt(&self, _scheme: &str, _operation: &str) {
628 self.tag("increment_retry_attempt");
629 }
630 fn increment_circuit_breaker_rejection(&self, _route: &str) {
631 self.tag("increment_circuit_breaker_rejection");
632 }
633 fn set_route_state(&self, _route: &str, _state: &str) {
634 self.tag("set_route_state");
635 }
636
637 fn clear_route_state(&self, _route: &str) {
638 self.tag("clear_route_state");
639 }
640 fn record_build_info(&self, _version: &str, _git_sha: &str) {
641 self.tag("record_build_info");
642 }
643 fn record_uptime(&self, _seconds: f64) {
644 self.tag("record_uptime");
645 }
646 fn record_component_operation(&self, _component: &str, _operation: &str, _outcome: &str) {
647 self.tag("record_component_operation");
648 }
649
650 fn set_master_leadership(&self, _lock: &str, _leader: bool) {
651 self.tag("set_master_leadership");
652 }
653 }
654
655 #[test]
656 fn test_noop_metrics_implements_trait() {
657 let metrics = NoOpMetrics;
658 let metrics_arc: Arc<dyn MetricsCollector> = Arc::new(metrics);
659
660 metrics_arc.record_exchange_duration("test-route", Duration::from_millis(100));
662 metrics_arc.increment_errors("test-route", "test-error");
663 metrics_arc.increment_exchanges("test-route");
664 metrics_arc.set_queue_depth("test-route", 5);
665 metrics_arc.record_circuit_breaker_change("test-route", "closed", "open");
666 }
667
668 #[test]
669 fn test_custom_metrics_collector() {
670 struct TestMetrics {
671 exchange_count: std::sync::atomic::AtomicU64,
672 }
673
674 impl MetricsCollector for TestMetrics {
675 fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
676 println!("Route {} took {}ms", route_id, duration.as_millis());
678 }
679
680 fn increment_errors(&self, route_id: &str, error_type: &str) {
681 println!("Route {} had error: {}", route_id, error_type);
683 }
684
685 fn increment_exchanges(&self, route_id: &str) {
686 self.exchange_count
688 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
689 println!("Route {} processed exchange", route_id);
690 }
691
692 fn set_queue_depth(&self, queue: &str, depth: usize) {
693 println!("Queue {queue} depth: {depth}");
695 }
696
697 fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str) {
698 println!("Route {} circuit breaker: {} -> {}", route_id, from, to);
700 }
701 }
702
703 let test_metrics = TestMetrics {
704 exchange_count: std::sync::atomic::AtomicU64::new(0),
705 };
706 let metrics_arc: Arc<dyn MetricsCollector> = Arc::new(test_metrics);
707
708 metrics_arc.record_exchange_duration("test-route", Duration::from_millis(100));
710 metrics_arc.increment_errors("test-route", "test-error");
711 metrics_arc.increment_exchanges("test-route");
712 metrics_arc.set_queue_depth("test-route", 5);
713 metrics_arc.record_circuit_breaker_change("test-route", "closed", "open");
714
715 }
718
719 #[test]
720 fn handle_delegates_to_stored_collector() {
721 let collector = Arc::new(RecordingMetrics::new());
722 let handle = MetricsHandle::new();
723 handle.register(collector.clone());
724
725 handle.record_exchange_duration("r", Duration::from_millis(1));
726
727 let recorded = collector.durations.lock().expect("durations lock").clone();
728 assert_eq!(recorded, vec![("r".to_string(), Duration::from_millis(1))]);
729 }
730
731 #[test]
732 fn second_registration_composes_both_observe() {
733 let a = Arc::new(RecordingMetrics::new());
734 let b = Arc::new(RecordingMetrics::new());
735 let handle = MetricsHandle::new();
736 handle.register(a.clone());
737 handle.register(b.clone());
738
739 handle.increment_errors("r", "x");
740
741 let a_errors = a.errors.lock().expect("errors lock").clone();
742 let b_errors = b.errors.lock().expect("errors lock").clone();
743 assert_eq!(a_errors, vec![("r".to_string(), "x".to_string())]);
744 assert_eq!(b_errors, vec![("r".to_string(), "x".to_string())]);
745 }
746
747 #[test]
748 fn register_same_arc_is_idempotent() {
749 let a = Arc::new(RecordingMetrics::new());
750 let handle = MetricsHandle::new();
751 handle.register(a.clone());
752 handle.register(a.clone());
753
754 handle.increment_exchanges("r");
755
756 let recorded = a.exchanges.lock().expect("exchanges lock").clone();
757 assert_eq!(recorded, vec!["r".to_string()]);
758 }
759
760 #[test]
761 fn handle_defaults_to_noop() {
762 let handle = MetricsHandle::new();
763 handle.record_exchange_duration("r", Duration::from_millis(1));
764 handle.increment_errors("r", "x");
765 handle.increment_exchanges("r");
766 handle.set_queue_depth("r", 5);
767 handle.record_circuit_breaker_change("r", "closed", "open");
768 handle.record_histogram("h", 1.0, &[("k", "v")]);
769 handle.record_counter("c", 1.0, &[("k", "v")]);
770 }
771
772 #[test]
773 fn composite_delegates_retry_and_rejection() {
774 let a = Arc::new(RecordingMetrics::new());
775 let b = Arc::new(RecordingMetrics::new());
776 let composite = CompositeMetricsCollector::new(vec![
777 Arc::clone(&a) as Arc<dyn MetricsCollector>,
778 Arc::clone(&b) as Arc<dyn MetricsCollector>,
779 ]);
780
781 composite.increment_retry_attempt("kafka", "connect");
782 composite.increment_circuit_breaker_rejection("r1");
783
784 for member in [&a, &b] {
785 assert_eq!(
786 member.retries.lock().expect("retries lock").clone(),
787 vec![("kafka".to_string(), "connect".to_string())]
788 );
789 assert_eq!(
790 member.rejections.lock().expect("rejections lock").clone(),
791 vec!["r1".to_string()]
792 );
793 }
794 }
795
796 #[test]
797 fn noop_defaults_compile_and_do_nothing() {
798 let collector: Arc<dyn MetricsCollector> = Arc::new(NoOpMetrics);
799 collector.increment_retry_attempt("kafka", "connect");
801 collector.increment_circuit_breaker_rejection("r1");
802 }
803
804 #[test]
807 fn composite_delegates_full_trait_surface() {
808 let a = Arc::new(SurfaceProbe::new());
809 let b = Arc::new(SurfaceProbe::new());
810 let composite = CompositeMetricsCollector::new(vec![
811 Arc::clone(&a) as Arc<dyn MetricsCollector>,
812 Arc::clone(&b) as Arc<dyn MetricsCollector>,
813 ]);
814
815 composite.record_exchange_duration("r", Duration::from_millis(1));
816 composite.increment_errors("r", "x");
817 composite.increment_exchanges("r");
818 composite.set_queue_depth("r", 1);
819 composite.record_circuit_breaker_change("r", "closed", "open");
820 composite.record_histogram("h", 1.0, &[("k", "v")]);
821 composite.record_counter("c", 1.0, &[("k", "v")]);
822 composite.increment_retry_attempt("kafka", "connect");
823 composite.increment_circuit_breaker_rejection("r1");
824 composite.set_route_state("r", "Started");
825 composite.clear_route_state("r");
826 composite.record_build_info("1.2.3", "abc1234");
827 composite.record_uptime(0.5);
828 composite.record_component_operation("redis", "command", "success");
829 composite.set_master_leadership("lock-a", true);
830
831 let expected = vec![
832 "record_exchange_duration",
833 "increment_errors",
834 "increment_exchanges",
835 "set_queue_depth",
836 "record_circuit_breaker_change",
837 "record_histogram",
838 "record_counter",
839 "increment_retry_attempt",
840 "increment_circuit_breaker_rejection",
841 "set_route_state",
842 "clear_route_state",
843 "record_build_info",
844 "record_uptime",
845 "record_component_operation",
846 "set_master_leadership",
847 ];
848 for member in [&a, &b] {
849 let calls = member.calls.lock().expect("calls lock").clone();
850 assert_eq!(calls, expected, "member missed part of the trait surface");
851 }
852 }
853
854 fn pinned_trio_expected() -> Vec<(&'static str, String, u64)> {
857 vec![
858 ("set_pinned_client_cache_size", "camel-https".to_string(), 3),
859 (
860 "increment_pinned_client_cache_hit",
861 "camel-https".to_string(),
862 1,
863 ),
864 (
865 "increment_pinned_client_cache_miss",
866 "camel-https".to_string(),
867 1,
868 ),
869 ]
870 }
871
872 #[test]
873 fn handle_forwards_pinned_cache_trio() {
874 let collector = Arc::new(RecordingMetrics::new());
875 let handle = MetricsHandle::new();
876 handle.register(collector.clone());
877
878 handle.set_pinned_client_cache_size("camel-https", 3);
879 handle.increment_pinned_client_cache_hit("camel-https");
880 handle.increment_pinned_client_cache_miss("camel-https");
881
882 let captured = collector.pinned.lock().expect("pinned lock").clone();
883 assert_eq!(captured, pinned_trio_expected());
884
885 let bystander = Arc::new(RecordingMetrics::new());
889 let unwired = MetricsHandle::new();
890 unwired.set_pinned_client_cache_size("camel-https", 3);
891 unwired.increment_pinned_client_cache_hit("camel-https");
892 unwired.increment_pinned_client_cache_miss("camel-https");
893 unwired.register(bystander.clone());
894 assert!(bystander.pinned.lock().expect("pinned lock").is_empty());
895 }
896
897 #[test]
898 fn composite_forwards_pinned_cache_trio_to_all_collectors() {
899 let a = Arc::new(RecordingMetrics::new());
900 let b = Arc::new(RecordingMetrics::new());
901 let composite = CompositeMetricsCollector::new(vec![
902 Arc::clone(&a) as Arc<dyn MetricsCollector>,
903 Arc::clone(&b) as Arc<dyn MetricsCollector>,
904 ]);
905
906 composite.set_pinned_client_cache_size("camel-https", 3);
907 composite.increment_pinned_client_cache_hit("camel-https");
908 composite.increment_pinned_client_cache_miss("camel-https");
909
910 for member in [&a, &b] {
911 let captured = member.pinned.lock().expect("pinned lock").clone();
912 assert_eq!(
913 captured,
914 pinned_trio_expected(),
915 "member missed part of the pinned-cache trio"
916 );
917 }
918 }
919
920 #[test]
923 fn allocator_stat_as_str_image_is_closed_set() {
924 let image: std::collections::BTreeSet<&'static str> = [
925 AllocatorStat::Allocated,
926 AllocatorStat::Resident,
927 AllocatorStat::Active,
928 AllocatorStat::Mapped,
929 ]
930 .iter()
931 .map(|stat| stat.as_str())
932 .collect();
933 let expected: std::collections::BTreeSet<&'static str> =
934 ["active", "allocated", "mapped", "resident"]
935 .into_iter()
936 .collect();
937 assert_eq!(image, expected);
938 }
939
940 #[test]
944 fn handle_and_composite_forward_set_allocator_memory() {
945 let expected = vec![(AllocatorStat::Resident, 4096)];
946
947 let handle_collector = Arc::new(RecordingMetrics::new());
948 let handle = MetricsHandle::new();
949 handle.register(handle_collector.clone());
950 handle.set_allocator_memory(AllocatorStat::Resident, 4096);
951 assert_eq!(
952 handle_collector
953 .allocator
954 .lock()
955 .expect("allocator lock")
956 .clone(),
957 expected,
958 "wired handle must forward exactly one allocator emission"
959 );
960
961 let composite_collector = Arc::new(RecordingMetrics::new());
962 let composite = CompositeMetricsCollector::new(vec![
963 composite_collector.clone() as Arc<dyn MetricsCollector>
964 ]);
965 composite.set_allocator_memory(AllocatorStat::Resident, 4096);
966 assert_eq!(
967 composite_collector
968 .allocator
969 .lock()
970 .expect("allocator lock")
971 .clone(),
972 expected,
973 "composite must forward exactly one allocator emission"
974 );
975
976 let bystander = Arc::new(RecordingMetrics::new());
977 let unwired = MetricsHandle::new();
978 unwired.set_allocator_memory(AllocatorStat::Resident, 4096);
979 unwired.register(bystander.clone());
980 assert!(
981 bystander
982 .allocator
983 .lock()
984 .expect("allocator lock")
985 .is_empty(),
986 "unwired-handle emissions are dropped, not replayed"
987 );
988 }
989
990 #[test]
995 fn handle_and_composite_forward_set_master_leadership() {
996 let expected = vec![("lock-a".to_string(), true), ("lock-a".to_string(), false)];
997
998 let handle_collector = Arc::new(RecordingMetrics::new());
999 let handle = MetricsHandle::new();
1000 handle.register(handle_collector.clone());
1001 handle.set_master_leadership("lock-a", true);
1002 handle.set_master_leadership("lock-a", false);
1003 assert_eq!(
1004 handle_collector
1005 .leadership
1006 .lock()
1007 .expect("leadership lock")
1008 .clone(),
1009 expected,
1010 "wired handle must forward every leadership edge"
1011 );
1012
1013 let composite_collector = Arc::new(RecordingMetrics::new());
1014 let composite = CompositeMetricsCollector::new(vec![
1015 composite_collector.clone() as Arc<dyn MetricsCollector>
1016 ]);
1017 composite.set_master_leadership("lock-a", true);
1018 composite.set_master_leadership("lock-a", false);
1019 assert_eq!(
1020 composite_collector
1021 .leadership
1022 .lock()
1023 .expect("leadership lock")
1024 .clone(),
1025 expected,
1026 "composite must forward every leadership edge"
1027 );
1028
1029 let bystander = Arc::new(RecordingMetrics::new());
1030 let unwired = MetricsHandle::new();
1031 unwired.set_master_leadership("lock-a", true);
1032 unwired.register(bystander.clone());
1033 assert!(
1034 bystander
1035 .leadership
1036 .lock()
1037 .expect("leadership lock")
1038 .is_empty(),
1039 "unwired-handle emissions are dropped, not replayed"
1040 );
1041 }
1042}