Skip to main content

camel_core/shared/components/domain/
registry.rs

1use std::collections::HashMap;
2use std::sync::Arc;
3
4use camel_api::CamelError;
5use camel_api::component_metadata::ComponentMetadata;
6use camel_component_api::Component;
7use tokio_util::sync::CancellationToken;
8
9/// Registry that stores components by their URI scheme.
10///
11/// Also harvests and indexes [`ComponentMetadata`] for each registered
12/// component, so the metadata can be queried through a
13/// [`ComponentMetadataCatalog`](camel_api::component_metadata::ComponentMetadataCatalog)
14/// without re-invoking the component.
15pub struct Registry {
16    components: HashMap<String, Arc<dyn Component>>,
17    metadata: HashMap<String, ComponentMetadata>,
18}
19
20impl Registry {
21    /// Create an empty registry.
22    pub fn new() -> Self {
23        Self {
24            components: HashMap::new(),
25            metadata: HashMap::new(),
26        }
27    }
28
29    /// Register a component. Replaces any existing component with the same scheme.
30    ///
31    /// Harvests the component's [`ComponentMetadata`] and indexes it by scheme
32    /// in parallel with the component insertion. Validates that the metadata's
33    /// scheme matches the component's scheme, normalizing on mismatch with a
34    /// warning log.
35    pub fn register(&mut self, component: Arc<dyn Component>) {
36        let scheme = component.scheme().to_string();
37        let mut metadata = component.metadata();
38        if let Err(e) = metadata.validate_scheme(&scheme) {
39            tracing::warn!(scheme = %scheme, error = %e, "metadata scheme mismatch, normalizing");
40            metadata.scheme = scheme.clone();
41        }
42        self.metadata.insert(scheme.clone(), metadata);
43        self.components.insert(scheme, component);
44    }
45
46    /// Look up a component by scheme.
47    pub fn get(&self, scheme: &str) -> Option<Arc<dyn Component>> {
48        self.components.get(scheme).cloned()
49    }
50
51    /// Look up a component by scheme, returning an error if not found.
52    pub fn get_or_err(&self, scheme: &str) -> Result<Arc<dyn Component>, CamelError> {
53        self.get(scheme)
54            .ok_or_else(|| CamelError::ComponentNotFound(scheme.to_string()))
55    }
56
57    /// Look up harvested metadata for a component by scheme.
58    pub fn get_metadata(&self, scheme: &str) -> Option<ComponentMetadata> {
59        self.metadata.get(scheme).cloned()
60    }
61
62    /// Return metadata for every registered component.
63    pub fn all_metadata(&self) -> Vec<ComponentMetadata> {
64        self.metadata.values().cloned().collect()
65    }
66
67    /// Return the schemes of every registered component's metadata.
68    pub fn metadata_schemes(&self) -> Vec<String> {
69        self.metadata.keys().cloned().collect()
70    }
71
72    /// Returns the number of registered components.
73    pub fn len(&self) -> usize {
74        self.components.len()
75    }
76
77    /// Returns true if no components are registered.
78    pub fn is_empty(&self) -> bool {
79        self.components.is_empty()
80    }
81}
82
83impl Default for Registry {
84    fn default() -> Self {
85        Self::new()
86    }
87}
88
89/// Adapter that lets `Registry` participate as a `ComponentContext`.
90///
91/// Wraps the shared `Arc<Mutex<Registry>>` and delegates `resolve_component`
92/// to `Registry::get`. The metrics collector is threaded in from the
93/// composition root (camel-cli) at construction; it is the ADR-0066
94/// late-bound handle, not a backend snapshot, so late registrations flow
95/// through [`camel_component_api::ComponentContext::metrics`] without
96/// re-snapshotting. The
97/// components-lever snapshot gates only the component-operations family —
98/// the error family is never lever-gated. When no collector is wired
99/// (e.g. compile-time security scan, standalone examples without a
100/// live context), construction resolves
101/// [`camel_api::NoOpMetrics`].
102///
103/// The registry reference is non-owning: the context holds only a
104/// [`std::sync::Weak`] to the registry, so it never keeps the registry
105/// (and, transitively, the `CamelContext` that owns it) alive.
106/// Component resolution works only while a strong registry reference
107/// exists elsewhere — the owning `CamelContext`, or an anchor the
108/// embedder retains. Once the last strong reference drops, resolution
109/// returns `None`.
110///
111/// Shutdown binding is optional: `new` leaves the context unbound
112/// (compile-time scans, tests), while
113/// [`with_shutdown_slot`](Self::with_shutdown_slot) binds it to the
114/// owning context's shutdown-token slot so resolved producers observe
115/// Runtime shutdown.
116pub struct RegistryComponentContext {
117    registry: std::sync::Weak<std::sync::Mutex<Registry>>,
118    metrics: Arc<dyn camel_api::MetricsCollector>,
119    components_enabled: bool,
120    shutdown_slot: Option<Arc<std::sync::Mutex<CancellationToken>>>,
121}
122
123impl RegistryComponentContext {
124    /// Builds the context, resolving the collector once: `metrics` when
125    /// wired, `NoOpMetrics` otherwise.
126    pub fn new(
127        registry: Arc<std::sync::Mutex<Registry>>,
128        metrics: Option<Arc<dyn camel_api::MetricsCollector>>,
129        components_enabled: bool,
130    ) -> Self {
131        Self {
132            registry: Arc::downgrade(&registry),
133            metrics: metrics.unwrap_or_else(|| Arc::new(camel_api::NoOpMetrics)),
134            components_enabled,
135            shutdown_slot: None,
136        }
137    }
138
139    /// Binds this context to a shutdown-token slot that mirrors the
140    /// owning `CamelContext`'s CURRENT shutdown token.
141    ///
142    /// The slot is written at context build time and re-written by every
143    /// `start()`, so [`shutdown_token`](camel_component_api::ComponentContext::shutdown_token)
144    /// resolves the token per call: after a stop/start cycle the resolved
145    /// token is the new boot's token, never a stale cancelled one. The
146    /// `Arc` shares only the slot — it does NOT extend the context's
147    /// lifetime, mirroring this struct's `Weak`-registry ownership: once
148    /// the owning context drops, its token is never cancelled (the same
149    /// observable outcome as `resolve_component` returning `None`).
150    pub fn with_shutdown_slot(mut self, slot: Arc<std::sync::Mutex<CancellationToken>>) -> Self {
151        self.shutdown_slot = Some(slot);
152        self
153    }
154}
155
156impl camel_component_api::ComponentContext for RegistryComponentContext {
157    fn resolve_component(&self, scheme: &str) -> Option<Arc<dyn camel_component_api::Component>> {
158        let registry = self.registry.upgrade()?;
159        let component = registry.lock().ok()?.get(scheme);
160        // Release the lock guard (statement end) and the upgraded anchor
161        // so the caller's return path retains neither.
162        drop(registry);
163        component
164    }
165
166    fn resolve_language(&self, _name: &str) -> Option<Arc<dyn camel_language_api::Language>> {
167        None
168    }
169
170    fn metrics(&self) -> Arc<dyn camel_api::MetricsCollector> {
171        self.metrics.clone()
172    }
173
174    fn component_metrics_enabled(&self) -> bool {
175        self.components_enabled
176    }
177
178    fn platform_service(&self) -> Arc<dyn camel_api::PlatformService> {
179        Arc::new(camel_api::NoopPlatformService::default())
180    }
181
182    fn register_route_health_check(
183        &self,
184        _route_id: &str,
185        _check: Arc<dyn camel_api::AsyncHealthCheck>,
186    ) {
187    }
188
189    fn unregister_route_health_check(&self, _route_id: &str) {}
190
191    /// The owning context's CURRENT shutdown token, read through the slot
192    /// mirror at call time (see [`with_shutdown_slot`](Self::with_shutdown_slot)).
193    /// Poison-tolerant: a poisoned slot resolves `None` rather than
194    /// panicking, degrading to the slot-less default (callers mint an
195    /// uncancelled local root — fail-safe, never fail-boot).
196    fn shutdown_token(&self) -> Option<CancellationToken> {
197        let slot = self.shutdown_slot.as_ref()?;
198        let guard = slot.lock().ok()?;
199        Some(guard.clone())
200    }
201}
202
203#[cfg(test)]
204mod tests {
205    use super::*;
206    use std::time::Duration;
207
208    use camel_api::MetricsCollector;
209    use camel_api::component_metadata::ComponentMetadata;
210    use camel_component_api::{ComponentContext, RuntimeObservability};
211    use camel_component_log::LogComponent;
212    use camel_component_timer::TimerComponent;
213
214    /// Recording state owned by the [`RecordingMetrics`] double. Owned
215    /// `String`s throughout — the facade passes formatted labels.
216    struct RecordingState {
217        errors: Vec<(String, String)>,
218        component_ops: Vec<(String, String, String)>,
219        counters: Vec<(String, f64)>,
220    }
221
222    /// Local recording double capturing the families the registry context
223    /// can emit: error pairs, component-op triples, generic counters. All
224    /// other trait methods are empty.
225    struct RecordingMetrics {
226        state: Arc<std::sync::Mutex<RecordingState>>,
227    }
228
229    impl RecordingMetrics {
230        fn new() -> Self {
231            Self {
232                state: Arc::new(std::sync::Mutex::new(RecordingState {
233                    errors: Vec::new(),
234                    component_ops: Vec::new(),
235                    counters: Vec::new(),
236                })),
237            }
238        }
239
240        fn recorded_errors(&self) -> Vec<(String, String)> {
241            self.state
242                .lock()
243                .expect("recording state lock")
244                .errors
245                .clone()
246        }
247
248        fn recorded_component_operations(&self) -> Vec<(String, String, String)> {
249            self.state
250                .lock()
251                .expect("recording state lock")
252                .component_ops
253                .clone()
254        }
255
256        fn recorded_counters(&self) -> Vec<(String, f64)> {
257            self.state
258                .lock()
259                .expect("recording state lock")
260                .counters
261                .clone()
262        }
263    }
264
265    impl MetricsCollector for RecordingMetrics {
266        fn record_exchange_duration(&self, _route_id: &str, _duration: Duration) {}
267
268        fn increment_errors(&self, route_id: &str, error_type: &str) {
269            self.state
270                .lock()
271                .expect("recording state lock")
272                .errors
273                .push((route_id.to_string(), error_type.to_string()));
274        }
275
276        fn increment_exchanges(&self, _route_id: &str) {}
277
278        fn set_queue_depth(&self, _queue: &str, _depth: usize) {}
279
280        fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {}
281
282        fn record_counter(&self, name: &str, value: f64, _labels: &[(&str, &str)]) {
283            self.state
284                .lock()
285                .expect("recording state lock")
286                .counters
287                .push((name.to_string(), value));
288        }
289
290        fn record_component_operation(&self, component: &str, operation: &str, outcome: &str) {
291            self.state
292                .lock()
293                .expect("recording state lock")
294                .component_ops
295                .push((
296                    component.to_string(),
297                    operation.to_string(),
298                    outcome.to_string(),
299                ));
300        }
301    }
302
303    #[test]
304    fn registry_starts_empty() {
305        let registry = Registry::new();
306        assert!(registry.is_empty());
307        assert_eq!(registry.len(), 0);
308        assert!(registry.get("timer").is_none());
309    }
310
311    #[test]
312    fn registry_registers_and_gets_components() {
313        let mut registry = Registry::new();
314        registry.register(Arc::new(TimerComponent::new()));
315        registry.register(Arc::new(LogComponent::new()));
316
317        assert_eq!(registry.len(), 2);
318        assert!(registry.get("timer").is_some());
319        assert!(registry.get("log").is_some());
320        assert!(!registry.is_empty());
321    }
322
323    #[test]
324    fn registry_get_or_err_reports_missing_component() {
325        let mut registry = Registry::new();
326        registry.register(Arc::new(TimerComponent::new()));
327
328        let err = match registry.get_or_err("missing") {
329            Ok(_) => panic!("must fail"),
330            Err(err) => err,
331        };
332        assert!(matches!(err, CamelError::ComponentNotFound(_)));
333    }
334
335    #[test]
336    fn registry_replaces_component_with_same_scheme() {
337        let mut registry = Registry::new();
338        registry.register(Arc::new(TimerComponent::new()));
339        registry.register(Arc::new(TimerComponent::new()));
340
341        assert_eq!(registry.len(), 1);
342        assert!(registry.get("timer").is_some());
343        assert_eq!(registry.all_metadata().len(), 1);
344    }
345
346    #[test]
347    fn registry_harvests_metadata_on_register() {
348        let mut registry = Registry::new();
349        registry.register(Arc::new(TimerComponent::new()));
350
351        let meta = registry.get_metadata("timer");
352        assert!(meta.is_some());
353        let meta = meta.unwrap(); // allow-unwrap
354        assert_eq!(meta.scheme, "timer");
355        assert_eq!(meta.schema_version, ComponentMetadata::SCHEMA_VERSION);
356    }
357
358    #[test]
359    fn registry_all_metadata_returns_all_schemes() {
360        let mut registry = Registry::new();
361        registry.register(Arc::new(TimerComponent::new()));
362        registry.register(Arc::new(LogComponent::new()));
363
364        let all = registry.all_metadata();
365        assert_eq!(all.len(), 2);
366    }
367
368    #[test]
369    fn registry_metadata_schemes_lists_all_keys() {
370        let mut registry = Registry::new();
371        registry.register(Arc::new(TimerComponent::new()));
372        registry.register(Arc::new(LogComponent::new()));
373
374        let mut schemes = registry.metadata_schemes();
375        schemes.sort();
376        assert_eq!(schemes, vec!["log".to_string(), "timer".to_string()]);
377    }
378
379    #[test]
380    fn metrics_returns_wired_collector_and_is_stable() {
381        let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
382        let collector = Arc::new(RecordingMetrics::new());
383        let wired_dyn: Arc<dyn MetricsCollector> = collector.clone();
384        let ctx = RegistryComponentContext::new(registry, Some(collector), false);
385
386        let first = ComponentContext::metrics(&ctx);
387        let second = ComponentContext::metrics(&ctx);
388        assert!(Arc::ptr_eq(&first, &wired_dyn));
389        assert!(Arc::ptr_eq(&second, &wired_dyn));
390    }
391
392    #[test]
393    fn component_metrics_enabled_reflects_constructor_lever() {
394        let on = RegistryComponentContext::new(
395            Arc::new(std::sync::Mutex::new(Registry::new())),
396            None,
397            true,
398        );
399        let off = RegistryComponentContext::new(
400            Arc::new(std::sync::Mutex::new(Registry::new())),
401            None,
402            false,
403        );
404
405        assert!(ComponentContext::component_metrics_enabled(&on));
406        assert!(!ComponentContext::component_metrics_enabled(&off));
407    }
408
409    #[test]
410    fn facade_error_family_reaches_wired_collector_with_lever_off() {
411        let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
412        let collector = Arc::new(RecordingMetrics::new());
413        let ctx = RegistryComponentContext::new(registry, Some(collector.clone()), false);
414
415        let facade = RuntimeObservability::component_metrics(&ctx);
416        facade.observe("wasm", "invoke", true);
417
418        assert_eq!(
419            collector.recorded_errors(),
420            vec![("wasm".to_string(), "e:wasm:invoke".to_string())]
421        );
422    }
423
424    #[test]
425    fn facade_component_family_gated_by_lever() {
426        let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
427        let collector = Arc::new(RecordingMetrics::new());
428        let ctx_on = RegistryComponentContext::new(registry.clone(), Some(collector.clone()), true);
429        let ctx_off = RegistryComponentContext::new(registry, Some(collector.clone()), false);
430
431        RuntimeObservability::component_metrics(&ctx_on).observe("wasm", "invoke", false);
432        RuntimeObservability::component_metrics(&ctx_off).observe("wasm", "invoke", false);
433
434        assert_eq!(
435            collector.recorded_component_operations(),
436            vec![(
437                "wasm".to_string(),
438                "invoke".to_string(),
439                "success".to_string()
440            )]
441        );
442        assert!(collector.recorded_errors().is_empty());
443        assert!(collector.recorded_counters().is_empty());
444    }
445
446    #[test]
447    fn late_registered_collector_reaches_registry_context() {
448        let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
449        let handle = Arc::new(camel_api::MetricsHandle::new());
450        let handle_dyn: Arc<dyn MetricsCollector> = handle.clone();
451        let recording = Arc::new(RecordingMetrics::new());
452        let ctx = RegistryComponentContext::new(registry, Some(handle_dyn), false);
453
454        handle.register(recording.clone());
455        ComponentContext::metrics(&ctx).increment_errors("wasm", "e:wasm:invoke");
456
457        assert_eq!(
458            recording.recorded_errors(),
459            vec![("wasm".to_string(), "e:wasm:invoke".to_string())]
460        );
461    }
462
463    #[test]
464    fn none_falls_back_to_noop_semantics() {
465        let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
466        let ctx = RegistryComponentContext::new(registry, None, false);
467
468        // None of these may panic: metrics() resolves the fallback
469        // collector, the facade builds over it, and the error flows into
470        // NoOp silently.
471        ComponentContext::metrics(&ctx);
472        RuntimeObservability::component_metrics(&ctx).observe("wasm", "invoke", true);
473
474        assert!(!ComponentContext::component_metrics_enabled(&ctx));
475    }
476
477    #[test]
478    fn resolve_component_unaffected_by_observability_params() {
479        let mut registry = Registry::new();
480        registry.register(Arc::new(TimerComponent::new()));
481        let registry = Arc::new(std::sync::Mutex::new(registry));
482        // Anchor: the context holds the registry weakly, so a strong
483        // reference must stay alive across the assertion.
484        let _registry_anchor = Arc::clone(&registry);
485        let ctx = RegistryComponentContext::new(registry, None, false);
486
487        assert!(ctx.resolve_component("timer").is_some());
488    }
489
490    #[test]
491    fn resolve_component_returns_component_while_anchored() {
492        let mut registry = Registry::new();
493        registry.register(Arc::new(TimerComponent::new()));
494        // The local binding is the strong anchor keeping the weak
495        // reference inside the context live.
496        let registry = Arc::new(std::sync::Mutex::new(registry));
497        let ctx = RegistryComponentContext::new(Arc::clone(&registry), None, false);
498
499        let Some(component) = ctx.resolve_component("timer") else {
500            panic!("anchored resolution must return the component");
501        };
502        assert_eq!(component.scheme(), "timer");
503    }
504
505    #[test]
506    fn resolve_component_returns_none_after_last_strong_ref_drops() {
507        let mut registry = Registry::new();
508        registry.register(Arc::new(TimerComponent::new()));
509        let registry = Arc::new(std::sync::Mutex::new(registry));
510        let ctx = RegistryComponentContext::new(Arc::clone(&registry), None, false);
511        // This binding held the only strong reference.
512        drop(registry);
513
514        assert!(ctx.resolve_component("timer").is_none());
515    }
516
517    /// Slot-less construction (security_boot compile-time scan shape) must
518    /// keep resolving `None` so callers mint an uncancelled local root.
519    #[test]
520    fn shutdown_token_is_none_without_shutdown_slot() {
521        let ctx = RegistryComponentContext::new(
522            Arc::new(std::sync::Mutex::new(Registry::new())),
523            None,
524            false,
525        );
526
527        assert!(ComponentContext::shutdown_token(&ctx).is_none());
528    }
529
530    /// Per-call resolution: replacing the slot's token (what `start()`
531    /// does on every boot) is observed by the next resolution, and the
532    /// replaced lineage does not cancel the fresh token.
533    #[test]
534    fn shutdown_token_resolves_the_current_slot_token() {
535        use tokio_util::sync::CancellationToken;
536
537        let slot = Arc::new(std::sync::Mutex::new(CancellationToken::new()));
538        let ctx = RegistryComponentContext::new(
539            Arc::new(std::sync::Mutex::new(Registry::new())),
540            None,
541            false,
542        )
543        .with_shutdown_slot(Arc::clone(&slot));
544
545        let stale = ComponentContext::shutdown_token(&ctx).expect("slot-bound context");
546        *slot.lock().expect("test slot lock") = CancellationToken::new();
547        let fresh = ComponentContext::shutdown_token(&ctx).expect("slot-bound context");
548
549        stale.cancel();
550        assert!(
551            !fresh.is_cancelled(),
552            "resolution is per call: the replacement token must not share the replaced lineage"
553        );
554    }
555}