Skip to main content

camel_api/
metrics.rs

1use std::sync::{Arc, Mutex};
2use std::time::Duration;
3
4use arc_swap::ArcSwap;
5
6/// The closed set of allocator memory statistics published through
7/// [`MetricsCollector::set_allocator_memory`].
8///
9/// # exhaustive-by-contract
10///
11/// exhaustive-by-contract: a closed 4-variant allocator stat set whose
12/// label values (`allocated | resident | active | mapped`) are fixed by the
13/// metrics spec; out-of-crate emitters (the camel-cli jemalloc sampler) match
14/// every variant, so adding one is a contract change, not a compatible
15/// extension.
16#[derive(Clone, Copy, PartialEq, Eq, Debug)]
17pub enum AllocatorStat {
18    /// Total bytes allocated by the allocator (in-use).
19    Allocated,
20    /// Resident bytes backed by physical pages (RSS contribution).
21    Resident,
22    /// Bytes in active pages.
23    Active,
24    /// Bytes in mapped virtual ranges.
25    Mapped,
26}
27
28impl AllocatorStat {
29    /// The Prometheus `stat` label value for this statistic.
30    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
40/// Trait for collecting metrics from the Camel runtime.
41/// Implementations can integrate with Prometheus, OpenTelemetry, etc.
42pub trait MetricsCollector: Send + Sync {
43    /// Record exchange processing time
44    fn record_exchange_duration(&self, route_id: &str, duration: Duration);
45
46    /// Increment error counter
47    fn increment_errors(&self, route_id: &str, error_type: &str);
48
49    /// Increment exchange counter
50    fn increment_exchanges(&self, route_id: &str);
51
52    /// Update the depth of a buffered stage's queue
53    /// (`camel_queue_depth{queue}`). The `queue` label is a closed set of
54    /// component-declared identifiers (`seda:<endpoint-name>`,
55    /// `aggregator:<route>`, `resequencer:<route>`).
56    fn set_queue_depth(&self, queue: &str, depth: usize);
57
58    /// Record circuit breaker state change
59    fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str);
60
61    /// Record a histogram observation (e.g., cost, latency distribution).
62    /// Default: no-op (backward-compatible).
63    fn record_histogram(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {}
64
65    /// Record a monotonically-increasing counter (e.g. `foo_total`).
66    /// Default: no-op (backward-compatible).
67    fn record_counter(&self, _name: &str, _value: f64, _labels: &[(&str, &str)]) {}
68
69    /// Increment the per-attempt retry counter (`camel_retry_attempts_total`,
70    /// labels scheme+operation). Called once per retry attempt, including the
71    /// first. Default: no-op (backward-compatible).
72    fn increment_retry_attempt(&self, _scheme: &str, _operation: &str) {}
73
74    /// Increment the circuit-breaker rejection counter
75    /// (`camel_circuit_breaker_rejections_total`, label route). Open-breaker
76    /// fast-fails count here, not as errors. Default: no-op
77    /// (backward-compatible).
78    fn increment_circuit_breaker_rejection(&self, _route: &str) {}
79
80    /// Publish a route lifecycle-state transition (`camel_route_state`,
81    /// labels route+state). `state` is the projection's state label — a
82    /// closed set by construction (`Registered`, `Starting`, `Started`,
83    /// `Suspended`, `Stopping`, `Stopped`, `Failed`). Implementations keep
84    /// the route's last-published state so a transition sets the new series
85    /// to 1 and zeroes the previous one. Default: no-op
86    /// (backward-compatible).
87    fn set_route_state(&self, _route: &str, _state: &str) {}
88
89    /// Drop a route's state series (route removed/undeployed) so a
90    /// scrape reflects only routes that exist.
91    fn clear_route_state(&self, _route: &str) {}
92
93    /// Publish build identification (`camel_build_info{git_sha,version}`,
94    /// value 1). Called once when the context is built. Default: no-op
95    /// (backward-compatible).
96    fn record_build_info(&self, _version: &str, _git_sha: &str) {}
97
98    /// Publish process uptime in seconds (`camel_uptime_seconds`),
99    /// refreshed periodically by the runtime. Default: no-op
100    /// (backward-compatible).
101    fn record_uptime(&self, _seconds: f64) {}
102
103    /// Increment the uniform component-operations counter
104    /// (`camel_component_operations_total`, labels component+operation+
105    /// outcome). `outcome` is a closed set — "success" or "failure"
106    /// only; callers derive it from a bool (see `ComponentMetrics`),
107    /// never pass free text. Default: no-op (backward-compatible).
108    fn record_component_operation(&self, _component: &str, _operation: &str, _outcome: &str) {}
109
110    /// Publish the pinned client cache size for a component
111    /// (`camel_pinned_client_cache_size{component}`, gauge, unit: entries).
112    /// Emitted by the owning component after each lookup, reflecting the
113    /// current (approximate) entry count. Default:
114    /// no-op (backward-compatible).
115    fn set_pinned_client_cache_size(&self, _component: &str, _entries: u64) {}
116
117    /// Increment the pinned client cache hit counter for a component
118    /// (`camel_pinned_client_cache_hits_total{component}`) — a pinned
119    /// lookup served by the cache without a rebuild. Default: no-op
120    /// (backward-compatible).
121    fn increment_pinned_client_cache_hit(&self, _component: &str) {}
122
123    /// Increment the pinned client cache miss counter for a component
124    /// (`camel_pinned_client_cache_misses_total{component}`) — a pinned
125    /// lookup that required a client rebuild. Default: no-op
126    /// (backward-compatible).
127    fn increment_pinned_client_cache_miss(&self, _component: &str) {}
128
129    /// Publish an allocator memory statistic
130    /// (`camel_allocator_memory_bytes{stat}`, gauge, unit: bytes). `stat`
131    /// is a closed [`AllocatorStat`] variant; the sampler refreshes the
132    /// current value periodically. Default: no-op (backward-compatible).
133    fn set_allocator_memory(&self, _stat: AllocatorStat, _bytes: u64) {}
134
135    /// Publish the leadership state for a master lock
136    /// (`camel_master_is_leader{lock}`, gauge): 1 while leadership is
137    /// held, 0 after it is lost. Emitted on the same observed state edges
138    /// as the `master_leadership_transitions_total` counter; the gauge
139    /// exists for steady-state readability ("who leads lock X now"), not
140    /// transition counting. Default: no-op (backward-compatible).
141    fn set_master_leadership(&self, _lock: &str, _leader: bool) {}
142}
143
144/// No-op metrics collector for default behavior
145pub 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
155/// Sized slot around `Arc<dyn MetricsCollector>`.
156///
157/// `ArcSwap`'s `RefCnt` implementation requires a `Sized` target, so a bare
158/// `ArcSwap<dyn MetricsCollector>` does not compile; this newtype restores
159/// `Sized`-ness without changing the stored pointee.
160struct CollectorSlot(Arc<dyn MetricsCollector>);
161
162/// A late-bound [`MetricsCollector`] cell.
163///
164/// Contract:
165///
166/// - **Late binding:** a `MetricsHandle` can be handed to consumers before any real
167///   collector exists; it seeds itself with [`NoOpMetrics`] so calls before (and
168///   without) registration are safe no-ops.
169/// - **Composition, not replacement:** each [`MetricsHandle::register`] composes the
170///   new collector *over* the currently stored one (see [`CompositeMetricsCollector`]);
171///   previously registered collectors keep observing.
172/// - **Same-Arc idempotence:** registering the same collector `Arc` twice is a no-op
173///   (detected via `Arc::ptr_eq` against the membership list), so a call site that
174///   wires the same collector through two builder paths does not double-count.
175/// - **Delegation cost:** each trait-method call costs one atomic load of the stored
176///   `Arc` (`ArcSwap::load`); the hot path never clones the `Arc`.
177pub struct MetricsHandle {
178    inner: ArcSwap<CollectorSlot>,
179    /// Membership list of every accepted collector, parallel to `inner`.
180    /// Kept because the stored `dyn` composite cannot be introspected for
181    /// `Arc::ptr_eq` dedupe.
182    members: Mutex<Vec<Arc<dyn MetricsCollector>>>,
183}
184
185impl MetricsHandle {
186    /// Creates a handle that delegates to [`NoOpMetrics`] until a collector is
187    /// registered.
188    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    /// Registers `collector`, composing it over whatever is currently stored.
196    ///
197    /// If the exact same `Arc` was already registered, this is a no-op
198    /// (see *same-Arc idempotence* in the type-level docs).
199    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"); // allow-unwrap
204        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            // Store directly — composing over the seeded NoOp would leave a
211            // permanent dead leg in every later composite chain.
212            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/// A [`MetricsCollector`] that fans every observation out to a list of collectors,
331/// in registration order.
332///
333/// Built by [`MetricsHandle::register`] — the second registration stores a
334/// composite of `[first, second]`; a third composes over that composite, so
335/// ordering and prior observation are preserved (composition, not replacement).
336///
337/// Internal type, hidden from the published docs. Out-of-tree code must not
338/// construct composites directly: registering an externally built composite
339/// plus its inner collector double-counts (the handle's opaque-Arc dedupe
340/// cannot see inside a composite). Register collectors via
341/// [`MetricsHandle::register`] instead.
342#[doc(hidden)]
343pub struct CompositeMetricsCollector {
344    collectors: Vec<Arc<dyn MetricsCollector>>,
345}
346
347impl CompositeMetricsCollector {
348    /// Creates a composite that delegates to `collectors` in order.
349    ///
350    /// Internal constructor, hidden from the published docs. Prefer
351    /// [`MetricsHandle::register`], which composes while deduplicating by
352    /// `Arc` pointer identity; direct construction bypasses that dedupe and
353    /// can double-count.
354    #[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    /// Test double that records observations for later inspection.
482    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    /// Test double that tags every trait-method call by name, for
588    /// delegation-parity assertions over the full `MetricsCollector` surface.
589    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        // All methods should execute without panicking
661        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                // In a real implementation, this would record the duration
677                println!("Route {} took {}ms", route_id, duration.as_millis());
678            }
679
680            fn increment_errors(&self, route_id: &str, error_type: &str) {
681                // In a real implementation, this would increment an error counter
682                println!("Route {} had error: {}", route_id, error_type);
683            }
684
685            fn increment_exchanges(&self, route_id: &str) {
686                // In a real implementation, this would increment an exchange counter
687                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                // In a real implementation, this would update a gauge
694                println!("Queue {queue} depth: {depth}");
695            }
696
697            fn record_circuit_breaker_change(&self, route_id: &str, from: &str, to: &str) {
698                // In a real implementation, this would record the state change
699                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        // Test that all methods work
709        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        // Note: We can't easily test the counter value without additional accessors
716        // This is just to verify the trait implementation works
717    }
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        // Both new methods must exist as no-op defaults: compile + no panic.
800        collector.increment_retry_attempt("kafka", "connect");
801        collector.increment_circuit_breaker_rejection("r1");
802    }
803
804    /// Delegation parity: the composite fans the full `MetricsCollector`
805    /// surface out to every member.
806    #[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    /// Expected pinned-cache triple captures for one call of each method
855    /// with component `"camel-https"` and entries `3` (counters record 1).
856    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        // An unwired handle delegates to the seeded NoOp: neither panics
886        // nor records into any collector double. Emissions made before
887        // registration are dropped, not buffered and replayed.
888        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    /// The `as_str()` image of `AllocatorStat` is the closed label-value set
921    /// (spec: `allocated | resident | active | mapped`).
922    #[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    /// `set_allocator_memory` forwards through a wired `MetricsHandle` and a
941    /// `CompositeMetricsCollector` (exactly one capture each); an unwired
942    /// handle neither panics nor records into a later-registered double.
943    #[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    /// `set_master_leadership` forwards through a wired `MetricsHandle` and
991    /// a `CompositeMetricsCollector` (exactly one capture each); an
992    /// unwired handle neither panics nor records into a later-registered
993    /// double.
994    #[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}