1use std::any::{Any, TypeId};
2use std::collections::HashMap;
3use std::sync::Arc;
4use tokio_util::sync::CancellationToken;
5use tracing::{debug, trace};
6
7#[cfg(test)]
8use camel_api::StepLifecycle;
9use camel_api::component_metadata::ComponentMetadata;
10use camel_api::error_handler::ErrorHandlerConfig;
11use camel_api::{
12 CamelError, FunctionInvoker, HealthReport, InFlightGauge, Lifecycle, MetricsCollector,
13 MetricsHandle, PlatformIdentity, PlatformService, ReadinessGate, RouteTemplateSpec,
14 RuntimeCommandBus, RuntimeQueryBus, TemplateInstanceRecord,
15};
16use camel_component_api::{Component, ComponentContext, ComponentRegistrar};
17use camel_language_api::Language;
18
19use crate::health_registry::HealthCheckRegistry;
20use crate::intercept::InterceptRules;
21use crate::language_registry::LanguageRegistryError;
22use crate::lifecycle::adapters::controller_actor::RouteControllerHandle;
23use crate::lifecycle::adapters::route_controller::SharedLanguageRegistry;
24use crate::lifecycle::application::route_definition::RouteDefinition;
25use crate::lifecycle::application::runtime_bus::RuntimeBus;
26use crate::registry::RegistryError;
27use crate::shared::components::domain::Registry;
28use crate::shared::observability::domain::{MetricsLeversConfig, TracerConfig};
29use crate::startup_validation::ConfigCheck;
30use crate::template::TemplateRegistry;
31
32pub use crate::context_builder::CamelContextBuilder;
33
34pub struct CamelContext {
43 registry: Arc<std::sync::Mutex<Registry>>,
44 route_controller: RouteControllerHandle,
45 actor_join: Option<tokio::task::JoinHandle<()>>,
46 supervision_join: Option<tokio::task::JoinHandle<()>>,
47 runtime: Arc<RuntimeBus>,
48 cancel_token: CancellationToken,
49 shutdown_token_slot: Arc<std::sync::Mutex<CancellationToken>>,
57 metrics: Arc<MetricsHandle>,
62 metrics_levers: MetricsLeversConfig,
68 platform_service: Arc<dyn PlatformService>,
70 languages: SharedLanguageRegistry,
71 shutdown_timeout: std::time::Duration,
72 services: Vec<Box<dyn Lifecycle>>,
73 health_registry: Arc<HealthCheckRegistry>,
74 component_configs: HashMap<TypeId, Box<dyn Any + Send + Sync>>,
75 function_invoker: Option<Arc<dyn FunctionInvoker>>,
76 template_registry: Arc<TemplateRegistry>,
77 idempotent_repositories: crate::registry::SharedIdempotentRegistry,
78 claim_check_repositories: crate::registry::SharedClaimCheckRegistry,
79 cache_repositories: crate::registry::SharedCacheRegistry,
80 startup_checks: Vec<Box<dyn ConfigCheck>>,
84 build_version: &'static str,
88 build_git_sha: &'static str,
89 build_started_at: std::time::Instant,
91 in_flight_total: Arc<InFlightGauge>,
97}
98
99pub(crate) struct FromParts {
102 pub(crate) registry: Arc<std::sync::Mutex<Registry>>,
103 pub(crate) route_controller: RouteControllerHandle,
104 pub(crate) _actor_join: tokio::task::JoinHandle<()>,
105 pub(crate) supervision_join: Option<tokio::task::JoinHandle<()>>,
106 pub(crate) runtime: Arc<RuntimeBus>,
107 pub(crate) cancel_token: CancellationToken,
108 pub(crate) shutdown_token_slot: Arc<std::sync::Mutex<CancellationToken>>,
109 pub(crate) metrics: Arc<MetricsHandle>,
110 pub(crate) platform_service: Arc<dyn PlatformService>,
111 pub(crate) languages: SharedLanguageRegistry,
112 pub(crate) shutdown_timeout: std::time::Duration,
113 pub(crate) services: Vec<Box<dyn Lifecycle>>,
114 pub(crate) health_registry: Arc<HealthCheckRegistry>,
115 pub(crate) component_configs: HashMap<TypeId, Box<dyn Any + Send + Sync>>,
116 pub(crate) function_invoker: Option<Arc<dyn FunctionInvoker>>,
117 pub(crate) template_registry: Arc<TemplateRegistry>,
118 pub(crate) idempotent_repositories: crate::registry::SharedIdempotentRegistry,
119 pub(crate) claim_check_repositories: crate::registry::SharedClaimCheckRegistry,
120 pub(crate) cache_repositories: crate::registry::SharedCacheRegistry,
121 pub(crate) startup_checks: Vec<Box<dyn ConfigCheck>>,
122 pub(crate) build_version: &'static str,
123 pub(crate) build_git_sha: &'static str,
124 pub(crate) build_started_at: std::time::Instant,
125 pub(crate) in_flight_total: Arc<InFlightGauge>,
126}
127
128impl CamelContext {
129 pub(crate) fn from_parts(parts: FromParts) -> Self {
130 Self {
131 registry: parts.registry,
132 route_controller: parts.route_controller,
133 actor_join: Some(parts._actor_join),
134 supervision_join: parts.supervision_join,
135 runtime: parts.runtime,
136 cancel_token: parts.cancel_token,
137 shutdown_token_slot: parts.shutdown_token_slot,
138 metrics: parts.metrics,
139 metrics_levers: MetricsLeversConfig::default(),
140 platform_service: parts.platform_service,
141 languages: parts.languages,
142 shutdown_timeout: parts.shutdown_timeout,
143 services: parts.services,
144 health_registry: parts.health_registry,
145 component_configs: parts.component_configs,
146 function_invoker: parts.function_invoker,
147 template_registry: parts.template_registry,
148 idempotent_repositories: parts.idempotent_repositories,
149 claim_check_repositories: parts.claim_check_repositories,
150 cache_repositories: parts.cache_repositories,
151 startup_checks: parts.startup_checks,
152 build_version: parts.build_version,
153 build_git_sha: parts.build_git_sha,
154 build_started_at: parts.build_started_at,
155 in_flight_total: parts.in_flight_total,
156 }
157 }
158}
159
160#[derive(Clone)]
164pub struct RuntimeExecutionHandle {
165 pub(crate) controller: RouteControllerHandle,
166 pub(crate) runtime: Arc<RuntimeBus>,
167 pub(crate) function_invoker: Option<Arc<dyn FunctionInvoker>>,
168 #[cfg(test)]
172 #[allow(clippy::type_complexity)]
173 pub(crate) test_lifecycle_inject: Arc<std::sync::Mutex<Option<Vec<Arc<dyn StepLifecycle>>>>>,
174}
175
176impl RuntimeExecutionHandle {
177 pub(crate) async fn add_route_definition(
178 &self,
179 definition: RouteDefinition,
180 ) -> Result<(), CamelError> {
181 use crate::lifecycle::application::ports::RouteRegistrationPort;
182 self.runtime
183 .register_route(definition)
184 .await
185 .map_err(Into::into)
186 }
187
188 #[allow(dead_code)]
191 pub(crate) async fn compile_route_definition(
192 &self,
193 definition: RouteDefinition,
194 ) -> Result<camel_api::BoxProcessor, CamelError> {
195 self.controller.compile_route_definition(definition).await
196 }
197
198 #[allow(dead_code)] pub(crate) async fn compile_route_definition_with_generation(
200 &self,
201 definition: RouteDefinition,
202 generation: u64,
203 ) -> Result<camel_api::BoxProcessor, CamelError> {
204 self.controller
205 .compile_route_definition_with_generation(definition, generation)
206 .await
207 }
208
209 pub(crate) async fn compile_route_definition_pipeline(
210 &self,
211 definition: RouteDefinition,
212 generation: u64,
213 ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
214 self.controller
215 .compile_route_definition_pipeline(definition, generation)
216 .await
217 }
218
219 pub(crate) async fn compile_route_definition_dry_pipeline(
222 &self,
223 definition: RouteDefinition,
224 ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
225 self.controller
226 .compile_route_definition_dry_pipeline(definition)
227 .await
228 }
229
230 pub(crate) async fn prepare_route_definition_with_generation(
231 &self,
232 definition: RouteDefinition,
233 generation: u64,
234 ) -> Result<crate::lifecycle::domain::route_compilation::PreparedRoute, CamelError> {
235 self.controller
236 .prepare_route_definition_with_generation(definition, generation)
237 .await
238 }
239
240 pub(crate) async fn insert_prepared_route(
241 &self,
242 prepared: crate::lifecycle::domain::route_compilation::PreparedRoute,
243 ) -> Result<(), CamelError> {
244 self.controller.insert_prepared_route(prepared).await
245 }
246
247 pub(crate) async fn discard_prepared_staging(&self, route_id: &str) -> Result<(), CamelError> {
248 self.controller.discard_prepared_staging(route_id).await
249 }
250
251 pub(crate) async fn remove_route_preserving_functions(
252 &self,
253 route_id: String,
254 ) -> Result<(), CamelError> {
255 self.controller
256 .remove_route_preserving_functions(route_id)
257 .await
258 }
259
260 pub(crate) async fn register_route_aggregate(
261 &self,
262 route_id: String,
263 ) -> Result<(), CamelError> {
264 self.runtime.register_aggregate_only(route_id).await
265 }
266
267 pub(crate) async fn swap_route_pipeline(
268 &self,
269 route_id: &str,
270 pipeline: camel_api::BoxProcessor,
271 ) -> Result<(), CamelError> {
272 self.controller.swap_pipeline(route_id, pipeline).await
273 }
274
275 pub(crate) async fn stop_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
277 self.controller.stop_route_reload(route_id).await
278 }
279
280 pub(crate) async fn start_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
282 self.controller.start_route_reload(route_id).await
283 }
284
285 pub(crate) async fn swap_route_pipeline_raw(
288 &self,
289 route_id: &str,
290 pipeline: camel_api::BoxProcessor,
291 lifecycle: Vec<Arc<dyn camel_api::StepLifecycle>>,
292 ) -> Result<(), CamelError> {
293 self.controller
294 .swap_pipeline_raw(route_id, pipeline, lifecycle)
295 .await
296 }
297
298 pub(crate) async fn execute_runtime_command(
299 &self,
300 cmd: camel_api::RuntimeCommand,
301 ) -> Result<camel_api::RuntimeCommandResult, CamelError> {
302 self.runtime.execute(cmd).await
303 }
304
305 pub(crate) async fn runtime_route_status(
306 &self,
307 route_id: &str,
308 ) -> Result<Option<String>, CamelError> {
309 match self
310 .runtime
311 .ask(camel_api::RuntimeQuery::GetRouteStatus {
312 route_id: route_id.to_string(),
313 })
314 .await
315 {
316 Ok(camel_api::RuntimeQueryResult::RouteStatus { status, .. }) => Ok(Some(status)),
317 Ok(_) => Err(CamelError::RouteError(
318 "unexpected runtime query response for route status".to_string(),
319 )),
320 Err(CamelError::RouteError(msg)) if msg.contains("not found") => Ok(None),
321 Err(err) => Err(err),
322 }
323 }
324
325 pub(crate) async fn runtime_route_ids(&self) -> Result<Vec<String>, CamelError> {
326 match self.runtime.ask(camel_api::RuntimeQuery::ListRoutes).await {
327 Ok(camel_api::RuntimeQueryResult::Routes { route_ids }) => Ok(route_ids),
328 Ok(_) => Err(CamelError::RouteError(
329 "unexpected runtime query response for route listing".to_string(),
330 )),
331 Err(err) => Err(err),
332 }
333 }
334
335 pub(crate) async fn route_source_hash(&self, route_id: &str) -> Option<u64> {
336 self.controller.route_source_hash(route_id).await
337 }
338
339 pub(crate) async fn in_flight_count(&self, route_id: &str) -> Result<u64, CamelError> {
340 if !self.controller.route_exists(route_id).await? {
341 return Err(CamelError::RouteError(format!(
342 "Route '{}' not found",
343 route_id
344 )));
345 }
346 Ok(self
347 .controller
348 .in_flight_count(route_id)
349 .await?
350 .unwrap_or(0))
351 }
352
353 pub(crate) async fn route_has_lifecycle(&self, route_id: &str) -> bool {
355 self.controller
356 .route_has_lifecycle(route_id)
357 .await
358 .unwrap_or(false)
359 }
360
361 pub(crate) fn function_invoker(&self) -> Option<Arc<dyn FunctionInvoker>> {
362 self.function_invoker.clone()
363 }
364
365 #[cfg(test)]
366 pub(crate) async fn force_start_route_for_test(
367 &self,
368 route_id: &str,
369 ) -> Result<(), CamelError> {
370 self.controller.start_route(route_id).await
371 }
372
373 pub async fn controller_route_count_for_test(&self) -> usize {
374 self.controller.route_count().await.unwrap_or(0)
375 }
376}
377
378#[async_trait::async_trait]
379impl crate::hot_reload::ports::ReloadExecutorPort for RuntimeExecutionHandle {
380 async fn add_route_definition(&self, definition: RouteDefinition) -> Result<(), CamelError> {
381 RuntimeExecutionHandle::add_route_definition(self, definition).await
382 }
383
384 async fn compile_route_definition_pipeline(
385 &self,
386 definition: RouteDefinition,
387 generation: u64,
388 ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
389 RuntimeExecutionHandle::compile_route_definition_pipeline(self, definition, generation)
390 .await
391 }
392
393 async fn compile_route_definition_dry_pipeline(
394 &self,
395 definition: RouteDefinition,
396 ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
397 RuntimeExecutionHandle::compile_route_definition_dry_pipeline(self, definition).await
398 }
399
400 async fn prepare_route_definition_with_generation(
401 &self,
402 definition: RouteDefinition,
403 generation: u64,
404 ) -> Result<crate::lifecycle::domain::route_compilation::PreparedRoute, CamelError> {
405 RuntimeExecutionHandle::prepare_route_definition_with_generation(
406 self, definition, generation,
407 )
408 .await
409 }
410
411 async fn insert_prepared_route(
412 &self,
413 prepared: crate::lifecycle::domain::route_compilation::PreparedRoute,
414 ) -> Result<(), CamelError> {
415 RuntimeExecutionHandle::insert_prepared_route(self, prepared).await
416 }
417
418 async fn discard_prepared_staging(&self, route_id: &str) -> Result<(), CamelError> {
419 RuntimeExecutionHandle::discard_prepared_staging(self, route_id).await
420 }
421
422 async fn remove_route_preserving_functions(&self, route_id: String) -> Result<(), CamelError> {
423 RuntimeExecutionHandle::remove_route_preserving_functions(self, route_id).await
424 }
425
426 async fn register_route_aggregate(&self, route_id: String) -> Result<(), CamelError> {
427 RuntimeExecutionHandle::register_route_aggregate(self, route_id).await
428 }
429
430 async fn swap_route_pipeline(
431 &self,
432 route_id: &str,
433 pipeline: camel_api::BoxProcessor,
434 ) -> Result<(), CamelError> {
435 RuntimeExecutionHandle::swap_route_pipeline(self, route_id, pipeline).await
436 }
437
438 async fn stop_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
439 RuntimeExecutionHandle::stop_route_reload(self, route_id).await
440 }
441
442 async fn start_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
443 RuntimeExecutionHandle::start_route_reload(self, route_id).await
444 }
445
446 async fn swap_route_pipeline_raw(
447 &self,
448 route_id: &str,
449 pipeline: camel_api::BoxProcessor,
450 lifecycle: Vec<std::sync::Arc<dyn camel_api::StepLifecycle>>,
451 ) -> Result<(), CamelError> {
452 RuntimeExecutionHandle::swap_route_pipeline_raw(self, route_id, pipeline, lifecycle).await
453 }
454
455 async fn execute_runtime_command(
456 &self,
457 cmd: camel_api::RuntimeCommand,
458 ) -> Result<camel_api::RuntimeCommandResult, CamelError> {
459 RuntimeExecutionHandle::execute_runtime_command(self, cmd).await
460 }
461
462 async fn runtime_route_status(&self, route_id: &str) -> Result<Option<String>, CamelError> {
463 RuntimeExecutionHandle::runtime_route_status(self, route_id).await
464 }
465
466 async fn in_flight_count(&self, route_id: &str) -> Result<u64, CamelError> {
467 RuntimeExecutionHandle::in_flight_count(self, route_id).await
468 }
469
470 async fn route_has_lifecycle(&self, route_id: &str) -> bool {
471 RuntimeExecutionHandle::route_has_lifecycle(self, route_id).await
472 }
473
474 #[cfg(test)]
475 fn take_test_lifecycle_inject(
476 &self,
477 ) -> Option<Vec<std::sync::Arc<dyn camel_api::StepLifecycle>>> {
478 self.test_lifecycle_inject.lock().unwrap().take()
479 }
480}
481
482impl CamelContext {
483 pub fn builder() -> CamelContextBuilder {
484 CamelContextBuilder::new()
485 }
486
487 pub async fn set_error_handler(&mut self, config: ErrorHandlerConfig) {
489 let _ = self.route_controller.set_error_handler(config).await;
490 }
491
492 pub async fn set_bind_exposure_acks(
496 &mut self,
497 acks: crate::lifecycle::adapters::route_controller_trait::BindExposureAcks,
498 ) {
499 let _ = self.route_controller.set_bind_exposure_acks(acks).await;
500 }
501
502 pub async fn set_tracing(&mut self, enabled: bool) {
504 let config = TracerConfig {
505 enabled,
506 ..Default::default()
507 };
508 self.metrics_levers = config.metrics_levers.clone();
511 let _ = self.route_controller.set_tracer_config(config).await;
512 }
513
514 pub async fn set_tracer_config(&mut self, config: TracerConfig) {
516 self.metrics_levers = config.metrics_levers.clone();
519 let _ = self.route_controller.set_tracer_config(config).await;
520 }
521
522 pub async fn with_tracing(mut self) -> Self {
524 self.set_tracing(true).await;
525 self
526 }
527
528 pub async fn with_tracer_config(mut self, config: TracerConfig) -> Self {
532 self.set_tracer_config(config).await;
533 self
534 }
535
536 pub fn with_lifecycle<L: Lifecycle + 'static>(mut self, service: L) -> Self {
545 self.add_lifecycle(service);
546 self
547 }
548
549 pub fn add_lifecycle<L: Lifecycle + 'static>(&mut self, service: L) {
554 if let Some(collector) = service.as_metrics_collector() {
555 self.metrics.register(collector);
559 self.metrics
564 .record_build_info(self.build_version, self.build_git_sha);
565 self.metrics
566 .record_uptime(self.build_started_at.elapsed().as_secs_f64());
567 }
568 if let Some(invoker) = service.as_function_invoker() {
569 self.function_invoker = Some(invoker.clone());
570 if let Err(e) = self.route_controller.try_set_function_invoker(invoker) {
571 tracing::debug!("Failed to propagate function invoker to route controller: {e}");
572 }
573 }
574
575 self.services.push(Box::new(service));
576 }
577
578 pub fn register_component<C: Component + 'static>(&mut self, component: C) {
584 self.register_component_dyn(Arc::new(component));
585 }
586
587 pub async fn set_intercept_rules(&self, rules: InterceptRules) -> Result<(), CamelError> {
594 self.route_controller.set_intercept_rules(rules).await
595 }
596
597 pub fn add_startup_check(&mut self, check: Box<dyn ConfigCheck>) {
606 self.startup_checks.push(check);
607 }
608
609 pub fn register_language(
616 &mut self,
617 name: impl Into<String>,
618 lang: Box<dyn Language>,
619 ) -> Result<(), LanguageRegistryError> {
620 let name = name.into();
621 let mut languages = self
622 .languages
623 .lock()
624 .expect("mutex poisoned: another thread panicked while holding this lock"); if languages.contains_key(&name) {
626 return Err(LanguageRegistryError::AlreadyRegistered { name });
627 }
628 languages.insert(name, Arc::from(lang));
629 Ok(())
630 }
631
632 pub fn resolve_language(&self, name: &str) -> Option<Arc<dyn Language>> {
634 let languages = self
635 .languages
636 .lock()
637 .expect("mutex poisoned: another thread panicked while holding this lock"); languages.get(name).cloned()
639 }
640
641 pub async fn add_route_definition(
645 &self,
646 definition: RouteDefinition,
647 ) -> Result<(), CamelError> {
648 use crate::lifecycle::application::ports::RouteRegistrationPort;
649 debug!(
650 from = definition.from_uri(),
651 route_id = %definition.route_id(),
652 "Adding route definition"
653 );
654 self.runtime
655 .register_route(definition)
656 .await
657 .map_err(Into::into)
658 }
659
660 pub fn registry(&self) -> std::sync::MutexGuard<'_, Registry> {
662 self.registry
663 .lock()
664 .expect("mutex poisoned: another thread panicked while holding this lock") }
666
667 pub fn registry_arc(&self) -> Arc<std::sync::Mutex<Registry>> {
669 Arc::clone(&self.registry)
670 }
671
672 pub fn runtime_execution_handle(&self) -> RuntimeExecutionHandle {
674 RuntimeExecutionHandle {
675 controller: self.route_controller.clone(),
676 runtime: Arc::clone(&self.runtime),
677 function_invoker: self.function_invoker.clone(),
678 #[cfg(test)]
679 test_lifecycle_inject: Arc::new(std::sync::Mutex::new(None)),
680 }
681 }
682
683 pub fn metrics(&self) -> Arc<dyn MetricsCollector> {
685 Arc::clone(&self.metrics) as Arc<dyn MetricsCollector>
686 }
687
688 pub fn total_in_flight(&self) -> u64 {
698 self.in_flight_total.total()
699 }
700
701 pub fn in_flight_gauge(&self) -> Arc<InFlightGauge> {
705 Arc::clone(&self.in_flight_total)
706 }
707
708 pub fn platform_service(&self) -> Arc<dyn PlatformService> {
710 Arc::clone(&self.platform_service)
711 }
712
713 pub fn readiness_gate(&self) -> Arc<dyn ReadinessGate> {
715 self.platform_service.readiness_gate()
716 }
717
718 pub fn platform_identity(&self) -> PlatformIdentity {
720 self.platform_service.identity()
721 }
722
723 pub fn leadership(&self) -> Arc<dyn camel_api::LeadershipService> {
725 self.platform_service.leadership()
726 }
727
728 pub fn runtime(&self) -> Arc<dyn camel_api::RuntimeHandle> {
730 self.runtime.clone()
731 }
732
733 pub fn producer_context(&self) -> camel_api::ProducerContext {
735 camel_api::ProducerContext::new().with_runtime(self.runtime())
736 }
737
738 pub async fn runtime_route_status(&self, route_id: &str) -> Result<Option<String>, CamelError> {
740 match self
741 .runtime()
742 .ask(camel_api::RuntimeQuery::GetRouteStatus {
743 route_id: route_id.to_string(),
744 })
745 .await
746 {
747 Ok(camel_api::RuntimeQueryResult::RouteStatus { status, .. }) => Ok(Some(status)),
748 Ok(_) => Err(CamelError::RouteError(
749 "unexpected runtime query response for route status".to_string(),
750 )),
751 Err(CamelError::RouteError(msg)) if msg.contains("not found") => Ok(None),
752 Err(err) => Err(err),
753 }
754 }
755
756 pub async fn start(&mut self) -> Result<(), CamelError> {
764 crate::lifecycle::application::context_lifecycle::start_context(
765 &mut self.services,
766 &mut self.startup_checks,
767 &self.runtime,
768 &self.route_controller,
769 &mut self.cancel_token,
770 &self.shutdown_token_slot,
771 )
772 .await?;
773 self.route_controller.mark_started().await
777 }
778
779 pub async fn stop(&mut self) -> Result<(), CamelError> {
781 self.stop_timeout(self.shutdown_timeout).await
782 }
783
784 pub async fn stop_timeout(&mut self, _timeout: std::time::Duration) -> Result<(), CamelError> {
794 crate::lifecycle::application::context_lifecycle::stop_context(
795 &self.cancel_token,
796 &mut self.supervision_join,
797 &self.runtime,
798 &self.route_controller,
799 &mut self.services,
800 )
801 .await
802 }
803
804 pub fn shutdown_timeout(&self) -> std::time::Duration {
806 self.shutdown_timeout
807 }
808
809 pub fn set_shutdown_timeout(&mut self, timeout: std::time::Duration) {
811 self.shutdown_timeout = timeout;
812 }
813
814 #[cfg(test)]
817 pub(crate) fn take_actor_join(&mut self) -> Option<tokio::task::JoinHandle<()>> {
818 self.actor_join.take()
819 }
820
821 pub async fn abort(&mut self) {
829 crate::lifecycle::application::context_lifecycle::abort_context(
830 &self.cancel_token,
831 &mut self.supervision_join,
832 &self.runtime,
833 &self.route_controller as &dyn crate::lifecycle::application::ports::RouteOrderingPort,
834 &self.route_controller
835 as &dyn crate::lifecycle::application::ports::RouteDestructiveTeardownPort,
836 &mut self.services,
837 self.health_registry.cancel_token(),
838 &mut self.actor_join,
839 )
840 .await
841 }
842
843 pub async fn health_check(&self) -> HealthReport {
845 use camel_api::HealthSource;
846 self.health_report().await
847 }
848
849 pub fn health_registry(&self) -> Arc<HealthCheckRegistry> {
850 Arc::clone(&self.health_registry)
851 }
852
853 pub fn set_component_config<T: 'static + Send + Sync>(&mut self, config: T) {
855 self.component_configs
856 .insert(TypeId::of::<T>(), Box::new(config));
857 }
858
859 pub fn get_component_config<T: 'static + Send + Sync>(&self) -> Option<&T> {
861 self.component_configs
862 .get(&TypeId::of::<T>())
863 .and_then(|b| b.downcast_ref::<T>())
864 }
865
866 pub fn component_metadata(&self, scheme: &str) -> Option<ComponentMetadata> {
870 self.registry.lock().ok()?.get_metadata(scheme)
871 }
872
873 pub fn all_component_metadata(&self) -> Vec<ComponentMetadata> {
875 self.registry
876 .lock()
877 .expect("mutex poisoned: another thread panicked while holding this lock") .all_metadata()
879 }
880
881 pub fn metadata_catalog(
888 &self,
889 ) -> crate::component_metadata_catalog::RuntimeComponentMetadataCatalog {
890 crate::component_metadata_catalog::RuntimeComponentMetadataCatalog::new(Arc::clone(
891 &self.registry,
892 ))
893 }
894
895 pub fn add_route_template(&self, spec: RouteTemplateSpec) -> Result<(), CamelError> {
901 self.template_registry.register(spec)
902 }
903
904 pub fn get_route_template(&self, id: &str) -> Option<RouteTemplateSpec> {
906 self.template_registry.get(id)
907 }
908
909 pub fn template_ids(&self) -> Vec<String> {
911 self.template_registry.template_ids()
912 }
913
914 pub fn record_template_instance(&self, record: TemplateInstanceRecord) {
916 self.template_registry.record_instance(record)
917 }
918
919 pub fn template_instances(&self, template_id: &str) -> Vec<TemplateInstanceRecord> {
921 self.template_registry.instances(template_id)
922 }
923
924 pub fn register_idempotent_repository(
931 &mut self,
932 name: impl Into<String>,
933 repo: Arc<dyn camel_api::IdempotentRepository>,
934 ) -> Result<(), RegistryError> {
935 self.idempotent_repositories.register(name, repo)
936 }
937
938 pub fn idempotent_repository(
940 &self,
941 name: &str,
942 ) -> Option<Arc<dyn camel_api::IdempotentRepository>> {
943 self.idempotent_repositories.get(name)
944 }
945
946 pub fn register_claim_check_repository(
953 &mut self,
954 name: impl Into<String>,
955 repo: Arc<dyn camel_api::ClaimCheckRepository>,
956 ) -> Result<(), RegistryError> {
957 self.claim_check_repositories.register(name, repo)
958 }
959
960 pub fn claim_check_repository(
962 &self,
963 name: &str,
964 ) -> Option<Arc<dyn camel_api::ClaimCheckRepository>> {
965 self.claim_check_repositories.get(name)
966 }
967
968 pub fn register_cache_repository(
975 &mut self,
976 name: impl Into<String>,
977 repo: Arc<dyn camel_api::CacheRepository>,
978 ) -> Result<(), RegistryError> {
979 self.cache_repositories.register(name, repo)
980 }
981
982 pub fn replace_cache_repository(
986 &mut self,
987 name: impl Into<String>,
988 repo: Arc<dyn camel_api::CacheRepository>,
989 ) -> Option<Arc<dyn camel_api::CacheRepository>> {
990 self.cache_repositories.register_or_replace(name, repo)
991 }
992
993 pub fn cache_repository(&self, name: &str) -> Option<Arc<dyn camel_api::CacheRepository>> {
995 self.cache_repositories.get(name)
996 }
997
998 pub fn shutdown_token(&self) -> CancellationToken {
1003 self.cancel_token.clone()
1004 }
1005
1006 pub fn shutdown_token_slot(&self) -> Arc<std::sync::Mutex<CancellationToken>> {
1017 Arc::clone(&self.shutdown_token_slot)
1018 }
1019}
1020
1021impl Drop for CamelContext {
1034 fn drop(&mut self) {
1035 if let Some(handle) = self.actor_join.take() {
1036 handle.abort();
1037 }
1038 if let Some(handle) = self.supervision_join.take() {
1039 handle.abort();
1040 }
1041 }
1042}
1043
1044impl ComponentRegistrar for CamelContext {
1045 fn register_component_dyn(&mut self, component: Arc<dyn Component>) {
1046 let scheme = component.scheme().to_string();
1047 self.registry
1048 .lock()
1049 .expect("mutex poisoned: another thread panicked while holding this lock") .register(component);
1051 trace!(scheme, "Registered component");
1052 }
1053}
1054
1055impl ComponentContext for CamelContext {
1056 fn resolve_component(&self, scheme: &str) -> Option<Arc<dyn Component>> {
1057 self.registry.lock().ok()?.get(scheme)
1058 }
1059
1060 fn resolve_language(&self, name: &str) -> Option<Arc<dyn Language>> {
1061 self.languages.lock().ok()?.get(name).cloned()
1062 }
1063
1064 fn metrics(&self) -> Arc<dyn MetricsCollector> {
1065 Arc::clone(&self.metrics) as Arc<dyn MetricsCollector>
1066 }
1067
1068 fn component_metrics_enabled(&self) -> bool {
1074 self.metrics_levers.components_enabled()
1075 }
1076
1077 fn health(&self) -> Arc<dyn camel_component_api::HealthCheckRegistry> {
1078 Arc::clone(&self.health_registry) as Arc<dyn camel_component_api::HealthCheckRegistry>
1081 }
1082
1083 fn platform_service(&self) -> Arc<dyn PlatformService> {
1084 Arc::clone(&self.platform_service)
1085 }
1086
1087 fn register_route_health_check(
1088 &self,
1089 route_id: &str,
1090 check: Arc<dyn camel_api::AsyncHealthCheck>,
1091 ) {
1092 self.health_registry.register_for_route(route_id, check);
1093 }
1094
1095 fn unregister_route_health_check(&self, route_id: &str) {
1096 self.health_registry.unregister_for_route(route_id);
1097 }
1098
1099 fn in_flight_counter(&self) -> Option<Arc<InFlightGauge>> {
1103 Some(Arc::clone(&self.in_flight_total))
1104 }
1105
1106 fn shutdown_token(&self) -> Option<CancellationToken> {
1110 Some(CamelContext::shutdown_token(self))
1111 }
1112}
1113
1114#[async_trait::async_trait]
1115impl camel_api::HealthSource for CamelContext {
1116 async fn liveness(&self) -> camel_api::HealthStatus {
1117 let has_failed = self
1118 .services
1119 .iter()
1120 .any(|s| s.status() == camel_api::ServiceStatus::Failed);
1121 if has_failed {
1122 camel_api::HealthStatus::Unhealthy
1123 } else {
1124 camel_api::HealthStatus::Healthy
1125 }
1126 }
1127
1128 async fn readiness(&self) -> camel_api::HealthStatus {
1129 let has_failed = self
1130 .services
1131 .iter()
1132 .any(|s| s.status() == camel_api::ServiceStatus::Failed);
1133 if has_failed {
1134 return camel_api::HealthStatus::Unhealthy;
1135 }
1136 let has_stopped = self
1137 .services
1138 .iter()
1139 .any(|s| s.status() == camel_api::ServiceStatus::Stopped);
1140 if has_stopped {
1141 return camel_api::HealthStatus::Degraded;
1142 }
1143 self.health_registry.check_all().await.status
1144 }
1145
1146 async fn health_report(&self) -> camel_api::HealthReport {
1147 let mut report = self.health_registry.check_all().await;
1148 let mut worst = report.status;
1149 for service in &self.services {
1150 let svc_status = service.status();
1151 let health = match svc_status {
1152 camel_api::ServiceStatus::Started => camel_api::HealthStatus::Healthy,
1153 camel_api::ServiceStatus::Stopped => camel_api::HealthStatus::Degraded,
1154 camel_api::ServiceStatus::Failed => camel_api::HealthStatus::Unhealthy,
1155 _ => camel_api::HealthStatus::Unhealthy,
1158 };
1159 if matches!(worst, camel_api::HealthStatus::Healthy)
1160 && matches!(
1161 health,
1162 camel_api::HealthStatus::Degraded | camel_api::HealthStatus::Unhealthy
1163 )
1164 {
1165 worst = health;
1166 }
1167 if matches!(worst, camel_api::HealthStatus::Degraded)
1168 && matches!(health, camel_api::HealthStatus::Unhealthy)
1169 {
1170 worst = health;
1171 }
1172 report.services.push(camel_api::ServiceHealth {
1173 name: service.name().to_string(),
1174 status: svc_status,
1175 message: None,
1176 });
1177 }
1178 report.status = worst;
1179 report
1180 }
1181
1182 async fn startup(&self) -> camel_api::HealthStatus {
1183 camel_api::HealthStatus::Healthy
1184 }
1185}
1186
1187#[cfg(test)]
1188#[path = "context_tests.rs"]
1189mod context_tests;