camel_core/shared/components/domain/
registry.rs1use std::collections::HashMap;
2use std::sync::Arc;
3
4use camel_api::CamelError;
5use camel_api::component_metadata::ComponentMetadata;
6use camel_component_api::Component;
7
8pub struct Registry {
15 components: HashMap<String, Arc<dyn Component>>,
16 metadata: HashMap<String, ComponentMetadata>,
17}
18
19impl Registry {
20 pub fn new() -> Self {
22 Self {
23 components: HashMap::new(),
24 metadata: HashMap::new(),
25 }
26 }
27
28 pub fn register(&mut self, component: Arc<dyn Component>) {
35 let scheme = component.scheme().to_string();
36 let mut metadata = component.metadata();
37 if let Err(e) = metadata.validate_scheme(&scheme) {
38 tracing::warn!(scheme = %scheme, error = %e, "metadata scheme mismatch, normalizing");
39 metadata.scheme = scheme.clone();
40 }
41 self.metadata.insert(scheme.clone(), metadata);
42 self.components.insert(scheme, component);
43 }
44
45 pub fn get(&self, scheme: &str) -> Option<Arc<dyn Component>> {
47 self.components.get(scheme).cloned()
48 }
49
50 pub fn get_or_err(&self, scheme: &str) -> Result<Arc<dyn Component>, CamelError> {
52 self.get(scheme)
53 .ok_or_else(|| CamelError::ComponentNotFound(scheme.to_string()))
54 }
55
56 pub fn get_metadata(&self, scheme: &str) -> Option<ComponentMetadata> {
58 self.metadata.get(scheme).cloned()
59 }
60
61 pub fn all_metadata(&self) -> Vec<ComponentMetadata> {
63 self.metadata.values().cloned().collect()
64 }
65
66 pub fn metadata_schemes(&self) -> Vec<String> {
68 self.metadata.keys().cloned().collect()
69 }
70
71 pub fn len(&self) -> usize {
73 self.components.len()
74 }
75
76 pub fn is_empty(&self) -> bool {
78 self.components.is_empty()
79 }
80}
81
82impl Default for Registry {
83 fn default() -> Self {
84 Self::new()
85 }
86}
87
88pub struct RegistryComponentContext {
110 registry: std::sync::Weak<std::sync::Mutex<Registry>>,
111 metrics: Arc<dyn camel_api::MetricsCollector>,
112 components_enabled: bool,
113}
114
115impl RegistryComponentContext {
116 pub fn new(
119 registry: Arc<std::sync::Mutex<Registry>>,
120 metrics: Option<Arc<dyn camel_api::MetricsCollector>>,
121 components_enabled: bool,
122 ) -> Self {
123 Self {
124 registry: Arc::downgrade(®istry),
125 metrics: metrics.unwrap_or_else(|| Arc::new(camel_api::NoOpMetrics)),
126 components_enabled,
127 }
128 }
129}
130
131impl camel_component_api::ComponentContext for RegistryComponentContext {
132 fn resolve_component(&self, scheme: &str) -> Option<Arc<dyn camel_component_api::Component>> {
133 let registry = self.registry.upgrade()?;
134 let component = registry.lock().ok()?.get(scheme);
135 drop(registry);
138 component
139 }
140
141 fn resolve_language(&self, _name: &str) -> Option<Arc<dyn camel_language_api::Language>> {
142 None
143 }
144
145 fn metrics(&self) -> Arc<dyn camel_api::MetricsCollector> {
146 self.metrics.clone()
147 }
148
149 fn component_metrics_enabled(&self) -> bool {
150 self.components_enabled
151 }
152
153 fn platform_service(&self) -> Arc<dyn camel_api::PlatformService> {
154 Arc::new(camel_api::NoopPlatformService::default())
155 }
156
157 fn register_route_health_check(
158 &self,
159 _route_id: &str,
160 _check: Arc<dyn camel_api::AsyncHealthCheck>,
161 ) {
162 }
163
164 fn unregister_route_health_check(&self, _route_id: &str) {}
165}
166
167#[cfg(test)]
168mod tests {
169 use super::*;
170 use std::time::Duration;
171
172 use camel_api::MetricsCollector;
173 use camel_api::component_metadata::ComponentMetadata;
174 use camel_component_api::{ComponentContext, RuntimeObservability};
175 use camel_component_log::LogComponent;
176 use camel_component_timer::TimerComponent;
177
178 struct RecordingState {
181 errors: Vec<(String, String)>,
182 component_ops: Vec<(String, String, String)>,
183 counters: Vec<(String, f64)>,
184 }
185
186 struct RecordingMetrics {
190 state: Arc<std::sync::Mutex<RecordingState>>,
191 }
192
193 impl RecordingMetrics {
194 fn new() -> Self {
195 Self {
196 state: Arc::new(std::sync::Mutex::new(RecordingState {
197 errors: Vec::new(),
198 component_ops: Vec::new(),
199 counters: Vec::new(),
200 })),
201 }
202 }
203
204 fn recorded_errors(&self) -> Vec<(String, String)> {
205 self.state
206 .lock()
207 .expect("recording state lock")
208 .errors
209 .clone()
210 }
211
212 fn recorded_component_operations(&self) -> Vec<(String, String, String)> {
213 self.state
214 .lock()
215 .expect("recording state lock")
216 .component_ops
217 .clone()
218 }
219
220 fn recorded_counters(&self) -> Vec<(String, f64)> {
221 self.state
222 .lock()
223 .expect("recording state lock")
224 .counters
225 .clone()
226 }
227 }
228
229 impl MetricsCollector for RecordingMetrics {
230 fn record_exchange_duration(&self, _route_id: &str, _duration: Duration) {}
231
232 fn increment_errors(&self, route_id: &str, error_type: &str) {
233 self.state
234 .lock()
235 .expect("recording state lock")
236 .errors
237 .push((route_id.to_string(), error_type.to_string()));
238 }
239
240 fn increment_exchanges(&self, _route_id: &str) {}
241
242 fn set_queue_depth(&self, _queue: &str, _depth: usize) {}
243
244 fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {}
245
246 fn record_counter(&self, name: &str, value: f64, _labels: &[(&str, &str)]) {
247 self.state
248 .lock()
249 .expect("recording state lock")
250 .counters
251 .push((name.to_string(), value));
252 }
253
254 fn record_component_operation(&self, component: &str, operation: &str, outcome: &str) {
255 self.state
256 .lock()
257 .expect("recording state lock")
258 .component_ops
259 .push((
260 component.to_string(),
261 operation.to_string(),
262 outcome.to_string(),
263 ));
264 }
265 }
266
267 #[test]
268 fn registry_starts_empty() {
269 let registry = Registry::new();
270 assert!(registry.is_empty());
271 assert_eq!(registry.len(), 0);
272 assert!(registry.get("timer").is_none());
273 }
274
275 #[test]
276 fn registry_registers_and_gets_components() {
277 let mut registry = Registry::new();
278 registry.register(Arc::new(TimerComponent::new()));
279 registry.register(Arc::new(LogComponent::new()));
280
281 assert_eq!(registry.len(), 2);
282 assert!(registry.get("timer").is_some());
283 assert!(registry.get("log").is_some());
284 assert!(!registry.is_empty());
285 }
286
287 #[test]
288 fn registry_get_or_err_reports_missing_component() {
289 let mut registry = Registry::new();
290 registry.register(Arc::new(TimerComponent::new()));
291
292 let err = match registry.get_or_err("missing") {
293 Ok(_) => panic!("must fail"),
294 Err(err) => err,
295 };
296 assert!(matches!(err, CamelError::ComponentNotFound(_)));
297 }
298
299 #[test]
300 fn registry_replaces_component_with_same_scheme() {
301 let mut registry = Registry::new();
302 registry.register(Arc::new(TimerComponent::new()));
303 registry.register(Arc::new(TimerComponent::new()));
304
305 assert_eq!(registry.len(), 1);
306 assert!(registry.get("timer").is_some());
307 assert_eq!(registry.all_metadata().len(), 1);
308 }
309
310 #[test]
311 fn registry_harvests_metadata_on_register() {
312 let mut registry = Registry::new();
313 registry.register(Arc::new(TimerComponent::new()));
314
315 let meta = registry.get_metadata("timer");
316 assert!(meta.is_some());
317 let meta = meta.unwrap(); assert_eq!(meta.scheme, "timer");
319 assert_eq!(meta.schema_version, ComponentMetadata::SCHEMA_VERSION);
320 }
321
322 #[test]
323 fn registry_all_metadata_returns_all_schemes() {
324 let mut registry = Registry::new();
325 registry.register(Arc::new(TimerComponent::new()));
326 registry.register(Arc::new(LogComponent::new()));
327
328 let all = registry.all_metadata();
329 assert_eq!(all.len(), 2);
330 }
331
332 #[test]
333 fn registry_metadata_schemes_lists_all_keys() {
334 let mut registry = Registry::new();
335 registry.register(Arc::new(TimerComponent::new()));
336 registry.register(Arc::new(LogComponent::new()));
337
338 let mut schemes = registry.metadata_schemes();
339 schemes.sort();
340 assert_eq!(schemes, vec!["log".to_string(), "timer".to_string()]);
341 }
342
343 #[test]
344 fn metrics_returns_wired_collector_and_is_stable() {
345 let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
346 let collector = Arc::new(RecordingMetrics::new());
347 let wired_dyn: Arc<dyn MetricsCollector> = collector.clone();
348 let ctx = RegistryComponentContext::new(registry, Some(collector), false);
349
350 let first = ComponentContext::metrics(&ctx);
351 let second = ComponentContext::metrics(&ctx);
352 assert!(Arc::ptr_eq(&first, &wired_dyn));
353 assert!(Arc::ptr_eq(&second, &wired_dyn));
354 }
355
356 #[test]
357 fn component_metrics_enabled_reflects_constructor_lever() {
358 let on = RegistryComponentContext::new(
359 Arc::new(std::sync::Mutex::new(Registry::new())),
360 None,
361 true,
362 );
363 let off = RegistryComponentContext::new(
364 Arc::new(std::sync::Mutex::new(Registry::new())),
365 None,
366 false,
367 );
368
369 assert!(ComponentContext::component_metrics_enabled(&on));
370 assert!(!ComponentContext::component_metrics_enabled(&off));
371 }
372
373 #[test]
374 fn facade_error_family_reaches_wired_collector_with_lever_off() {
375 let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
376 let collector = Arc::new(RecordingMetrics::new());
377 let ctx = RegistryComponentContext::new(registry, Some(collector.clone()), false);
378
379 let facade = RuntimeObservability::component_metrics(&ctx);
380 facade.observe("wasm", "invoke", true);
381
382 assert_eq!(
383 collector.recorded_errors(),
384 vec![("wasm".to_string(), "e:wasm:invoke".to_string())]
385 );
386 }
387
388 #[test]
389 fn facade_component_family_gated_by_lever() {
390 let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
391 let collector = Arc::new(RecordingMetrics::new());
392 let ctx_on = RegistryComponentContext::new(registry.clone(), Some(collector.clone()), true);
393 let ctx_off = RegistryComponentContext::new(registry, Some(collector.clone()), false);
394
395 RuntimeObservability::component_metrics(&ctx_on).observe("wasm", "invoke", false);
396 RuntimeObservability::component_metrics(&ctx_off).observe("wasm", "invoke", false);
397
398 assert_eq!(
399 collector.recorded_component_operations(),
400 vec![(
401 "wasm".to_string(),
402 "invoke".to_string(),
403 "success".to_string()
404 )]
405 );
406 assert!(collector.recorded_errors().is_empty());
407 assert!(collector.recorded_counters().is_empty());
408 }
409
410 #[test]
411 fn late_registered_collector_reaches_registry_context() {
412 let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
413 let handle = Arc::new(camel_api::MetricsHandle::new());
414 let handle_dyn: Arc<dyn MetricsCollector> = handle.clone();
415 let recording = Arc::new(RecordingMetrics::new());
416 let ctx = RegistryComponentContext::new(registry, Some(handle_dyn), false);
417
418 handle.register(recording.clone());
419 ComponentContext::metrics(&ctx).increment_errors("wasm", "e:wasm:invoke");
420
421 assert_eq!(
422 recording.recorded_errors(),
423 vec![("wasm".to_string(), "e:wasm:invoke".to_string())]
424 );
425 }
426
427 #[test]
428 fn none_falls_back_to_noop_semantics() {
429 let registry = Arc::new(std::sync::Mutex::new(Registry::new()));
430 let ctx = RegistryComponentContext::new(registry, None, false);
431
432 ComponentContext::metrics(&ctx);
436 RuntimeObservability::component_metrics(&ctx).observe("wasm", "invoke", true);
437
438 assert!(!ComponentContext::component_metrics_enabled(&ctx));
439 }
440
441 #[test]
442 fn resolve_component_unaffected_by_observability_params() {
443 let mut registry = Registry::new();
444 registry.register(Arc::new(TimerComponent::new()));
445 let registry = Arc::new(std::sync::Mutex::new(registry));
446 let _registry_anchor = Arc::clone(®istry);
449 let ctx = RegistryComponentContext::new(registry, None, false);
450
451 assert!(ctx.resolve_component("timer").is_some());
452 }
453
454 #[test]
455 fn resolve_component_returns_component_while_anchored() {
456 let mut registry = Registry::new();
457 registry.register(Arc::new(TimerComponent::new()));
458 let registry = Arc::new(std::sync::Mutex::new(registry));
461 let ctx = RegistryComponentContext::new(Arc::clone(®istry), None, false);
462
463 let Some(component) = ctx.resolve_component("timer") else {
464 panic!("anchored resolution must return the component");
465 };
466 assert_eq!(component.scheme(), "timer");
467 }
468
469 #[test]
470 fn resolve_component_returns_none_after_last_strong_ref_drops() {
471 let mut registry = Registry::new();
472 registry.register(Arc::new(TimerComponent::new()));
473 let registry = Arc::new(std::sync::Mutex::new(registry));
474 let ctx = RegistryComponentContext::new(Arc::clone(®istry), None, false);
475 drop(registry);
477
478 assert!(ctx.resolve_component("timer").is_none());
479 }
480}