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;
7use tokio_util::sync::CancellationToken;
8
9pub struct Registry {
16 components: HashMap<String, Arc<dyn Component>>,
17 metadata: HashMap<String, ComponentMetadata>,
18}
19
20impl Registry {
21 pub fn new() -> Self {
23 Self {
24 components: HashMap::new(),
25 metadata: HashMap::new(),
26 }
27 }
28
29 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 pub fn get(&self, scheme: &str) -> Option<Arc<dyn Component>> {
48 self.components.get(scheme).cloned()
49 }
50
51 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 pub fn get_metadata(&self, scheme: &str) -> Option<ComponentMetadata> {
59 self.metadata.get(scheme).cloned()
60 }
61
62 pub fn all_metadata(&self) -> Vec<ComponentMetadata> {
64 self.metadata.values().cloned().collect()
65 }
66
67 pub fn metadata_schemes(&self) -> Vec<String> {
69 self.metadata.keys().cloned().collect()
70 }
71
72 pub fn len(&self) -> usize {
74 self.components.len()
75 }
76
77 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
89pub 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 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(®istry),
133 metrics: metrics.unwrap_or_else(|| Arc::new(camel_api::NoOpMetrics)),
134 components_enabled,
135 shutdown_slot: None,
136 }
137 }
138
139 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 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 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 struct RecordingState {
217 errors: Vec<(String, String)>,
218 component_ops: Vec<(String, String, String)>,
219 counters: Vec<(String, f64)>,
220 }
221
222 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(); 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 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 let _registry_anchor = Arc::clone(®istry);
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 let registry = Arc::new(std::sync::Mutex::new(registry));
497 let ctx = RegistryComponentContext::new(Arc::clone(®istry), 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(®istry), None, false);
511 drop(registry);
513
514 assert!(ctx.resolve_component("timer").is_none());
515 }
516
517 #[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 #[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}