Skip to main content

camel_core/
context.rs

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
34/// The CamelContext is the runtime engine that manages components, routes, and their lifecycle.
35///
36/// # Lifecycle
37///
38/// Call [`start()`](Self::start) to launch routes, then [`stop()`](Self::stop)
39/// or [`abort()`](Self::abort) to shut down. A stopped context can be restarted
40/// by calling `start()` again — the controller actor stays alive across stop/start
41/// cycles; only [`abort()`](Self::abort) is destructive.
42pub 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    /// Shared late-bound metrics cell (rc-hrm1.3): the SAME handle instance
50    /// seeds the route controller's `tracer_metrics` and the RuntimeBus
51    /// collector, so a collector registered at any time — including after
52    /// routes are added — is observed by all subsequent emission calls.
53    metrics: Arc<MetricsHandle>,
54    /// Snapshot of the metric-family levers from the last
55    /// `set_tracer_config` call (default: components off). Read
56    /// synchronously by `ComponentContext::component_metrics_enabled`
57    /// — the controller actor holds its own copy for pipeline gating,
58    /// which is not reachable from a sync context method.
59    metrics_levers: MetricsLeversConfig,
60    // Platform ports
61    platform_service: Arc<dyn PlatformService>,
62    languages: SharedLanguageRegistry,
63    shutdown_timeout: std::time::Duration,
64    services: Vec<Box<dyn Lifecycle>>,
65    health_registry: Arc<HealthCheckRegistry>,
66    component_configs: HashMap<TypeId, Box<dyn Any + Send + Sync>>,
67    function_invoker: Option<Arc<dyn FunctionInvoker>>,
68    template_registry: Arc<TemplateRegistry>,
69    idempotent_repositories: crate::registry::SharedIdempotentRegistry,
70    claim_check_repositories: crate::registry::SharedClaimCheckRegistry,
71    cache_repositories: crate::registry::SharedCacheRegistry,
72    /// Fail-closed startup validation registry (ADR-0033). Checks are drained
73    /// and executed synchronously at the head of [`start()`](Self::start) before
74    /// any route consumer is started.
75    startup_checks: Vec<Box<dyn ConfigCheck>>,
76    /// Build identity (dashboard-observability T3.2): re-published when a
77    /// metrics collector registers late, because the composite does not
78    /// replay past observations to new members.
79    build_version: &'static str,
80    build_git_sha: &'static str,
81    /// Anchor for `camel_uptime_seconds` (context build time).
82    build_started_at: std::time::Instant,
83    /// Context-global accepted-not-completed gauge (drainclaim): every
84    /// live `InFlightClaim` on this context increments it, and the last
85    /// release fires the gauge's idle notification. The SAME `Arc`
86    /// is installed into the route controller, so consumers, producers,
87    /// and the inline dispatcher all mint claims against one gauge.
88    in_flight_total: Arc<InFlightGauge>,
89}
90
91/// Parts bag used by [`CamelContextBuilder::build`] to construct a [`CamelContext`]
92/// without accessing private fields from a sibling module.
93pub(crate) struct FromParts {
94    pub(crate) registry: Arc<std::sync::Mutex<Registry>>,
95    pub(crate) route_controller: RouteControllerHandle,
96    pub(crate) _actor_join: tokio::task::JoinHandle<()>,
97    pub(crate) supervision_join: Option<tokio::task::JoinHandle<()>>,
98    pub(crate) runtime: Arc<RuntimeBus>,
99    pub(crate) cancel_token: CancellationToken,
100    pub(crate) metrics: Arc<MetricsHandle>,
101    pub(crate) platform_service: Arc<dyn PlatformService>,
102    pub(crate) languages: SharedLanguageRegistry,
103    pub(crate) shutdown_timeout: std::time::Duration,
104    pub(crate) services: Vec<Box<dyn Lifecycle>>,
105    pub(crate) health_registry: Arc<HealthCheckRegistry>,
106    pub(crate) component_configs: HashMap<TypeId, Box<dyn Any + Send + Sync>>,
107    pub(crate) function_invoker: Option<Arc<dyn FunctionInvoker>>,
108    pub(crate) template_registry: Arc<TemplateRegistry>,
109    pub(crate) idempotent_repositories: crate::registry::SharedIdempotentRegistry,
110    pub(crate) claim_check_repositories: crate::registry::SharedClaimCheckRegistry,
111    pub(crate) cache_repositories: crate::registry::SharedCacheRegistry,
112    pub(crate) startup_checks: Vec<Box<dyn ConfigCheck>>,
113    pub(crate) build_version: &'static str,
114    pub(crate) build_git_sha: &'static str,
115    pub(crate) build_started_at: std::time::Instant,
116    pub(crate) in_flight_total: Arc<InFlightGauge>,
117}
118
119impl CamelContext {
120    pub(crate) fn from_parts(parts: FromParts) -> Self {
121        Self {
122            registry: parts.registry,
123            route_controller: parts.route_controller,
124            actor_join: Some(parts._actor_join),
125            supervision_join: parts.supervision_join,
126            runtime: parts.runtime,
127            cancel_token: parts.cancel_token,
128            metrics: parts.metrics,
129            metrics_levers: MetricsLeversConfig::default(),
130            platform_service: parts.platform_service,
131            languages: parts.languages,
132            shutdown_timeout: parts.shutdown_timeout,
133            services: parts.services,
134            health_registry: parts.health_registry,
135            component_configs: parts.component_configs,
136            function_invoker: parts.function_invoker,
137            template_registry: parts.template_registry,
138            idempotent_repositories: parts.idempotent_repositories,
139            claim_check_repositories: parts.claim_check_repositories,
140            cache_repositories: parts.cache_repositories,
141            startup_checks: parts.startup_checks,
142            build_version: parts.build_version,
143            build_git_sha: parts.build_git_sha,
144            build_started_at: parts.build_started_at,
145            in_flight_total: parts.in_flight_total,
146        }
147    }
148}
149
150/// Opaque handle for runtime side-effect execution operations.
151///
152/// This intentionally does not expose direct lifecycle mutation APIs to callers.
153#[derive(Clone)]
154pub struct RuntimeExecutionHandle {
155    pub(crate) controller: RouteControllerHandle,
156    pub(crate) runtime: Arc<RuntimeBus>,
157    pub(crate) function_invoker: Option<Arc<dyn FunctionInvoker>>,
158    /// Lifecycle handles to inject into the compiled pipeline during
159    /// `apply_swap`.  Used in tests to simulate lifecycle-bearing routes
160    /// (e.g. resequencer).  Always `None` in production.
161    #[cfg(test)]
162    #[allow(clippy::type_complexity)]
163    pub(crate) test_lifecycle_inject: Arc<std::sync::Mutex<Option<Vec<Arc<dyn StepLifecycle>>>>>,
164}
165
166impl RuntimeExecutionHandle {
167    pub(crate) async fn add_route_definition(
168        &self,
169        definition: RouteDefinition,
170    ) -> Result<(), CamelError> {
171        use crate::lifecycle::application::ports::RouteRegistrationPort;
172        self.runtime
173            .register_route(definition)
174            .await
175            .map_err(Into::into)
176    }
177
178    /// Compile a route definition into a bare BoxProcessor (no lifecycle).
179    /// Kept for tests; hot-reload uses the lifecycle-preserving variant instead.
180    #[allow(dead_code)]
181    pub(crate) async fn compile_route_definition(
182        &self,
183        definition: RouteDefinition,
184    ) -> Result<camel_api::BoxProcessor, CamelError> {
185        self.controller.compile_route_definition(definition).await
186    }
187
188    #[allow(dead_code)] // kept for potential future hot-reload paths
189    pub(crate) async fn compile_route_definition_with_generation(
190        &self,
191        definition: RouteDefinition,
192        generation: u64,
193    ) -> Result<camel_api::BoxProcessor, CamelError> {
194        self.controller
195            .compile_route_definition_with_generation(definition, generation)
196            .await
197    }
198
199    pub(crate) async fn compile_route_definition_pipeline(
200        &self,
201        definition: RouteDefinition,
202        generation: u64,
203    ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
204        self.controller
205            .compile_route_definition_pipeline(definition, generation)
206            .await
207    }
208
209    /// Compile without function generation, returning full CompiledPipeline.
210    /// Oracle Fix 1: stateless hot-reload path preserves lifecycle handles.
211    pub(crate) async fn compile_route_definition_dry_pipeline(
212        &self,
213        definition: RouteDefinition,
214    ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
215        self.controller
216            .compile_route_definition_dry_pipeline(definition)
217            .await
218    }
219
220    pub(crate) async fn prepare_route_definition_with_generation(
221        &self,
222        definition: RouteDefinition,
223        generation: u64,
224    ) -> Result<crate::lifecycle::domain::route_compilation::PreparedRoute, CamelError> {
225        self.controller
226            .prepare_route_definition_with_generation(definition, generation)
227            .await
228    }
229
230    pub(crate) async fn insert_prepared_route(
231        &self,
232        prepared: crate::lifecycle::domain::route_compilation::PreparedRoute,
233    ) -> Result<(), CamelError> {
234        self.controller.insert_prepared_route(prepared).await
235    }
236
237    pub(crate) async fn discard_prepared_staging(&self, route_id: &str) -> Result<(), CamelError> {
238        self.controller.discard_prepared_staging(route_id).await
239    }
240
241    pub(crate) async fn remove_route_preserving_functions(
242        &self,
243        route_id: String,
244    ) -> Result<(), CamelError> {
245        self.controller
246            .remove_route_preserving_functions(route_id)
247            .await
248    }
249
250    pub(crate) async fn register_route_aggregate(
251        &self,
252        route_id: String,
253    ) -> Result<(), CamelError> {
254        self.runtime.register_aggregate_only(route_id).await
255    }
256
257    pub(crate) async fn swap_route_pipeline(
258        &self,
259        route_id: &str,
260        pipeline: camel_api::BoxProcessor,
261    ) -> Result<(), CamelError> {
262        self.controller.swap_pipeline(route_id, pipeline).await
263    }
264
265    /// Stop the route via the reload path (graceful lifecycle drain).
266    pub(crate) async fn stop_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
267        self.controller.stop_route_reload(route_id).await
268    }
269
270    /// Start the route via the reload path (re-create consumer).
271    pub(crate) async fn start_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
272        self.controller.start_route_reload(route_id).await
273    }
274
275    /// Raw pipeline swap — bypasses the lifecycle/aggregate rejection check.
276    /// Only safe after the route has been stopped (Restart path).
277    pub(crate) async fn swap_route_pipeline_raw(
278        &self,
279        route_id: &str,
280        pipeline: camel_api::BoxProcessor,
281        lifecycle: Vec<Arc<dyn camel_api::StepLifecycle>>,
282    ) -> Result<(), CamelError> {
283        self.controller
284            .swap_pipeline_raw(route_id, pipeline, lifecycle)
285            .await
286    }
287
288    pub(crate) async fn execute_runtime_command(
289        &self,
290        cmd: camel_api::RuntimeCommand,
291    ) -> Result<camel_api::RuntimeCommandResult, CamelError> {
292        self.runtime.execute(cmd).await
293    }
294
295    pub(crate) async fn runtime_route_status(
296        &self,
297        route_id: &str,
298    ) -> Result<Option<String>, CamelError> {
299        match self
300            .runtime
301            .ask(camel_api::RuntimeQuery::GetRouteStatus {
302                route_id: route_id.to_string(),
303            })
304            .await
305        {
306            Ok(camel_api::RuntimeQueryResult::RouteStatus { status, .. }) => Ok(Some(status)),
307            Ok(_) => Err(CamelError::RouteError(
308                "unexpected runtime query response for route status".to_string(),
309            )),
310            Err(CamelError::RouteError(msg)) if msg.contains("not found") => Ok(None),
311            Err(err) => Err(err),
312        }
313    }
314
315    pub(crate) async fn runtime_route_ids(&self) -> Result<Vec<String>, CamelError> {
316        match self.runtime.ask(camel_api::RuntimeQuery::ListRoutes).await {
317            Ok(camel_api::RuntimeQueryResult::Routes { route_ids }) => Ok(route_ids),
318            Ok(_) => Err(CamelError::RouteError(
319                "unexpected runtime query response for route listing".to_string(),
320            )),
321            Err(err) => Err(err),
322        }
323    }
324
325    pub(crate) async fn route_source_hash(&self, route_id: &str) -> Option<u64> {
326        self.controller.route_source_hash(route_id).await
327    }
328
329    pub(crate) async fn in_flight_count(&self, route_id: &str) -> Result<u64, CamelError> {
330        if !self.controller.route_exists(route_id).await? {
331            return Err(CamelError::RouteError(format!(
332                "Route '{}' not found",
333                route_id
334            )));
335        }
336        Ok(self
337            .controller
338            .in_flight_count(route_id)
339            .await?
340            .unwrap_or(0))
341    }
342
343    /// Check whether the running route has lifecycle-bearing steps.
344    pub(crate) async fn route_has_lifecycle(&self, route_id: &str) -> bool {
345        self.controller
346            .route_has_lifecycle(route_id)
347            .await
348            .unwrap_or(false)
349    }
350
351    pub(crate) fn function_invoker(&self) -> Option<Arc<dyn FunctionInvoker>> {
352        self.function_invoker.clone()
353    }
354
355    #[cfg(test)]
356    pub(crate) async fn force_start_route_for_test(
357        &self,
358        route_id: &str,
359    ) -> Result<(), CamelError> {
360        self.controller.start_route(route_id).await
361    }
362
363    pub async fn controller_route_count_for_test(&self) -> usize {
364        self.controller.route_count().await.unwrap_or(0)
365    }
366}
367
368#[async_trait::async_trait]
369impl crate::hot_reload::ports::ReloadExecutorPort for RuntimeExecutionHandle {
370    async fn add_route_definition(&self, definition: RouteDefinition) -> Result<(), CamelError> {
371        RuntimeExecutionHandle::add_route_definition(self, definition).await
372    }
373
374    async fn compile_route_definition_pipeline(
375        &self,
376        definition: RouteDefinition,
377        generation: u64,
378    ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
379        RuntimeExecutionHandle::compile_route_definition_pipeline(self, definition, generation)
380            .await
381    }
382
383    async fn compile_route_definition_dry_pipeline(
384        &self,
385        definition: RouteDefinition,
386    ) -> Result<crate::lifecycle::domain::CompiledPipeline, CamelError> {
387        RuntimeExecutionHandle::compile_route_definition_dry_pipeline(self, definition).await
388    }
389
390    async fn prepare_route_definition_with_generation(
391        &self,
392        definition: RouteDefinition,
393        generation: u64,
394    ) -> Result<crate::lifecycle::domain::route_compilation::PreparedRoute, CamelError> {
395        RuntimeExecutionHandle::prepare_route_definition_with_generation(
396            self, definition, generation,
397        )
398        .await
399    }
400
401    async fn insert_prepared_route(
402        &self,
403        prepared: crate::lifecycle::domain::route_compilation::PreparedRoute,
404    ) -> Result<(), CamelError> {
405        RuntimeExecutionHandle::insert_prepared_route(self, prepared).await
406    }
407
408    async fn discard_prepared_staging(&self, route_id: &str) -> Result<(), CamelError> {
409        RuntimeExecutionHandle::discard_prepared_staging(self, route_id).await
410    }
411
412    async fn remove_route_preserving_functions(&self, route_id: String) -> Result<(), CamelError> {
413        RuntimeExecutionHandle::remove_route_preserving_functions(self, route_id).await
414    }
415
416    async fn register_route_aggregate(&self, route_id: String) -> Result<(), CamelError> {
417        RuntimeExecutionHandle::register_route_aggregate(self, route_id).await
418    }
419
420    async fn swap_route_pipeline(
421        &self,
422        route_id: &str,
423        pipeline: camel_api::BoxProcessor,
424    ) -> Result<(), CamelError> {
425        RuntimeExecutionHandle::swap_route_pipeline(self, route_id, pipeline).await
426    }
427
428    async fn stop_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
429        RuntimeExecutionHandle::stop_route_reload(self, route_id).await
430    }
431
432    async fn start_route_reload(&self, route_id: &str) -> Result<(), CamelError> {
433        RuntimeExecutionHandle::start_route_reload(self, route_id).await
434    }
435
436    async fn swap_route_pipeline_raw(
437        &self,
438        route_id: &str,
439        pipeline: camel_api::BoxProcessor,
440        lifecycle: Vec<std::sync::Arc<dyn camel_api::StepLifecycle>>,
441    ) -> Result<(), CamelError> {
442        RuntimeExecutionHandle::swap_route_pipeline_raw(self, route_id, pipeline, lifecycle).await
443    }
444
445    async fn execute_runtime_command(
446        &self,
447        cmd: camel_api::RuntimeCommand,
448    ) -> Result<camel_api::RuntimeCommandResult, CamelError> {
449        RuntimeExecutionHandle::execute_runtime_command(self, cmd).await
450    }
451
452    async fn runtime_route_status(&self, route_id: &str) -> Result<Option<String>, CamelError> {
453        RuntimeExecutionHandle::runtime_route_status(self, route_id).await
454    }
455
456    async fn in_flight_count(&self, route_id: &str) -> Result<u64, CamelError> {
457        RuntimeExecutionHandle::in_flight_count(self, route_id).await
458    }
459
460    async fn route_has_lifecycle(&self, route_id: &str) -> bool {
461        RuntimeExecutionHandle::route_has_lifecycle(self, route_id).await
462    }
463
464    #[cfg(test)]
465    fn take_test_lifecycle_inject(
466        &self,
467    ) -> Option<Vec<std::sync::Arc<dyn camel_api::StepLifecycle>>> {
468        self.test_lifecycle_inject.lock().unwrap().take()
469    }
470}
471
472impl CamelContext {
473    pub fn builder() -> CamelContextBuilder {
474        CamelContextBuilder::new()
475    }
476
477    /// Set a global error handler applied to all routes without a per-route handler.
478    pub async fn set_error_handler(&mut self, config: ErrorHandlerConfig) {
479        let _ = self.route_controller.set_error_handler(config).await;
480    }
481
482    /// Install per-bind public-exposure acknowledgements (ADR-0061).
483    /// Built by the CLI from `CamelConfig.binds`; the per-bind gate fails
484    /// closed on non-loopback binds until acknowledged.
485    pub async fn set_bind_exposure_acks(
486        &mut self,
487        acks: crate::lifecycle::adapters::route_controller_trait::BindExposureAcks,
488    ) {
489        let _ = self.route_controller.set_bind_exposure_acks(acks).await;
490    }
491
492    /// Enable or disable tracing globally.
493    pub async fn set_tracing(&mut self, enabled: bool) {
494        let config = TracerConfig {
495            enabled,
496            ..Default::default()
497        };
498        // Keep the lever snapshot in lockstep with what is forwarded
499        // (this path resets the levers to their defaults).
500        self.metrics_levers = config.metrics_levers.clone();
501        let _ = self.route_controller.set_tracer_config(config).await;
502    }
503
504    /// Configure tracing with full config.
505    pub async fn set_tracer_config(&mut self, config: TracerConfig) {
506        // Snapshot the levers for the sync `component_metrics_enabled`
507        // surface before the config moves into the controller actor.
508        self.metrics_levers = config.metrics_levers.clone();
509        let _ = self.route_controller.set_tracer_config(config).await;
510    }
511
512    /// Builder-style: enable tracing with default config.
513    pub async fn with_tracing(mut self) -> Self {
514        self.set_tracing(true).await;
515        self
516    }
517
518    /// Builder-style: configure tracing with custom config.
519    /// Note: tracing subscriber initialization (stdout/file output) is handled
520    /// separately via init_tracing_subscriber (called in camel-config bridge).
521    pub async fn with_tracer_config(mut self, config: TracerConfig) -> Self {
522        self.set_tracer_config(config).await;
523        self
524    }
525
526    /// Register a lifecycle service (Apache Camel: addService pattern)
527    ///
528    /// For services exposing `as_function_invoker()`, the invoker is propagated
529    /// to the route controller so that subsequent route definitions with function
530    /// steps work correctly.
531    ///
532    /// Prefer [`CamelContextBuilder::with_lifecycle`] when possible, which wires
533    /// the invoker at build time before any routes are added.
534    pub fn with_lifecycle<L: Lifecycle + 'static>(mut self, service: L) -> Self {
535        self.add_lifecycle(service);
536        self
537    }
538
539    /// `&mut self` sibling of [`Self::with_lifecycle`]: register a lifecycle
540    /// service on a context the caller still owns. Required by the
541    /// `camel_bundles::boot` path (ADR-0069 section 10), which receives a
542    /// `&mut CamelContext` and cannot run the consuming builder.
543    pub fn add_lifecycle<L: Lifecycle + 'static>(&mut self, service: L) {
544        if let Some(collector) = service.as_metrics_collector() {
545            // Late-bound registration (compose, never replace): the shared
546            // handle fans the collector into every emission path seeded at
547            // build time — context slot, controller tracer path, RuntimeBus.
548            self.metrics.register(collector);
549            // The composite does not replay past observations to a newly
550            // registered collector, so re-publish the identity gauges: a
551            // collector wired post-build (e.g. configure_context's
552            // PrometheusService) still reports build info and uptime.
553            self.metrics
554                .record_build_info(self.build_version, self.build_git_sha);
555            self.metrics
556                .record_uptime(self.build_started_at.elapsed().as_secs_f64());
557        }
558        if let Some(invoker) = service.as_function_invoker() {
559            self.function_invoker = Some(invoker.clone());
560            if let Err(e) = self.route_controller.try_set_function_invoker(invoker) {
561                tracing::debug!("Failed to propagate function invoker to route controller: {e}");
562            }
563        }
564
565        self.services.push(Box::new(service));
566    }
567
568    /// Register a component with this context.
569    ///
570    /// Delegates to [`register_component_dyn`](Self::register_component_dyn)
571    /// so metadata harvesting happens regardless of which entry point
572    /// is used.
573    pub fn register_component<C: Component + 'static>(&mut self, component: C) {
574        self.register_component_dyn(Arc::new(component));
575    }
576
577    /// Install route send-point interception rules (pre-first-use only).
578    ///
579    /// Returns `CamelError::Config` once frozen: after the first route is
580    /// registered or the context is started, compiled pipelines have
581    /// captured the rule set and it cannot change. Builder-time rules via
582    /// [`CamelContextBuilder::with_intercept_rules`] bypass this gate.
583    pub async fn set_intercept_rules(&self, rules: InterceptRules) -> Result<(), CamelError> {
584        self.route_controller.set_intercept_rules(rules).await
585    }
586
587    /// Register a startup `ConfigCheck` to be evaluated at the head of
588    /// [`start()`](Self::start). Established by ADR-0033.
589    ///
590    /// Checks are drained and executed synchronously before any route consumer
591    /// is started. If any check returns `Err`, `start()` fails closed with
592    /// `CamelError::Config(_)` and no route is started. The check list is
593    /// consumed (moved) during `start()` so this method may be called multiple
594    /// times to register an arbitrary number of checks.
595    pub fn add_startup_check(&mut self, check: Box<dyn ConfigCheck>) {
596        self.startup_checks.push(check);
597    }
598
599    /// Register a language with this context, keyed by name.
600    ///
601    /// Returns `Err(LanguageRegistryError::AlreadyRegistered)` if a language
602    /// with the same name is already registered. Use
603    /// [`resolve_language`](Self::resolve_language) to check before
604    /// registering, or choose a distinct name.
605    pub fn register_language(
606        &mut self,
607        name: impl Into<String>,
608        lang: Box<dyn Language>,
609    ) -> Result<(), LanguageRegistryError> {
610        let name = name.into();
611        let mut languages = self
612            .languages
613            .lock()
614            .expect("mutex poisoned: another thread panicked while holding this lock"); // allow-unwrap
615        if languages.contains_key(&name) {
616            return Err(LanguageRegistryError::AlreadyRegistered { name });
617        }
618        languages.insert(name, Arc::from(lang));
619        Ok(())
620    }
621
622    /// Resolve a language by name. Returns `None` if not registered.
623    pub fn resolve_language(&self, name: &str) -> Option<Arc<dyn Language>> {
624        let languages = self
625            .languages
626            .lock()
627            .expect("mutex poisoned: another thread panicked while holding this lock"); // allow-unwrap
628        languages.get(name).cloned()
629    }
630
631    /// Add a route definition to this context.
632    ///
633    /// The route must have an ID. Steps are resolved immediately using registered components.
634    pub async fn add_route_definition(
635        &self,
636        definition: RouteDefinition,
637    ) -> Result<(), CamelError> {
638        use crate::lifecycle::application::ports::RouteRegistrationPort;
639        debug!(
640            from = definition.from_uri(),
641            route_id = %definition.route_id(),
642            "Adding route definition"
643        );
644        self.runtime
645            .register_route(definition)
646            .await
647            .map_err(Into::into)
648    }
649
650    /// Access the component registry.
651    pub fn registry(&self) -> std::sync::MutexGuard<'_, Registry> {
652        self.registry
653            .lock()
654            .expect("mutex poisoned: another thread panicked while holding this lock") // allow-unwrap
655    }
656
657    /// Access the shared component registry Arc.
658    pub fn registry_arc(&self) -> Arc<std::sync::Mutex<Registry>> {
659        Arc::clone(&self.registry)
660    }
661
662    /// Get runtime execution handle for file-watcher integrations.
663    pub fn runtime_execution_handle(&self) -> RuntimeExecutionHandle {
664        RuntimeExecutionHandle {
665            controller: self.route_controller.clone(),
666            runtime: Arc::clone(&self.runtime),
667            function_invoker: self.function_invoker.clone(),
668            #[cfg(test)]
669            test_lifecycle_inject: Arc::new(std::sync::Mutex::new(None)),
670        }
671    }
672
673    /// Get the metrics collector (the shared late-bound handle).
674    pub fn metrics(&self) -> Arc<dyn MetricsCollector> {
675        Arc::clone(&self.metrics) as Arc<dyn MetricsCollector>
676    }
677
678    /// Canonical accepted-not-completed count for the whole context
679    /// (drainclaim): the number of exchanges accepted through a counted
680    /// path (seda enqueue, `ConsumerContext::send`, inline dispatch)
681    /// whose pipelines have not completed yet.
682    ///
683    /// A single atomic load — the drain verdict is linearizable by
684    /// construction; a read of zero states that no counted exchange is
685    /// awaiting completion. Distinct from the per-route `drain_in_flight`
686    /// stop-bookkeeping (ADR-0043), which is untouched.
687    pub fn total_in_flight(&self) -> u64 {
688        self.in_flight_total.total()
689    }
690
691    /// Clone of the shared in-flight gauge (drainclaim): the runner seam
692    /// for notification-based settle — await `idle()` for the zero
693    /// transition instead of polling [`Self::total_in_flight`].
694    pub fn in_flight_gauge(&self) -> Arc<InFlightGauge> {
695        Arc::clone(&self.in_flight_total)
696    }
697
698    /// Get the platform service.
699    pub fn platform_service(&self) -> Arc<dyn PlatformService> {
700        Arc::clone(&self.platform_service)
701    }
702
703    /// Get the readiness gate port.
704    pub fn readiness_gate(&self) -> Arc<dyn ReadinessGate> {
705        self.platform_service.readiness_gate()
706    }
707
708    /// Get the platform identity.
709    pub fn platform_identity(&self) -> PlatformIdentity {
710        self.platform_service.identity()
711    }
712
713    /// Get the leadership service port.
714    pub fn leadership(&self) -> Arc<dyn camel_api::LeadershipService> {
715        self.platform_service.leadership()
716    }
717
718    /// Get runtime command/query bus handle.
719    pub fn runtime(&self) -> Arc<dyn camel_api::RuntimeHandle> {
720        self.runtime.clone()
721    }
722
723    /// Build a producer context wired to this runtime.
724    pub fn producer_context(&self) -> camel_api::ProducerContext {
725        camel_api::ProducerContext::new().with_runtime(self.runtime())
726    }
727
728    /// Query route status via runtime read-model.
729    pub async fn runtime_route_status(&self, route_id: &str) -> Result<Option<String>, CamelError> {
730        match self
731            .runtime()
732            .ask(camel_api::RuntimeQuery::GetRouteStatus {
733                route_id: route_id.to_string(),
734            })
735            .await
736        {
737            Ok(camel_api::RuntimeQueryResult::RouteStatus { status, .. }) => Ok(Some(status)),
738            Ok(_) => Err(CamelError::RouteError(
739                "unexpected runtime query response for route status".to_string(),
740            )),
741            Err(CamelError::RouteError(msg)) if msg.contains("not found") => Ok(None),
742            Err(err) => Err(err),
743        }
744    }
745
746    /// Start all routes. Each route's consumer will begin producing exchanges.
747    ///
748    /// Only routes with `auto_startup == true` will be started, in order of their
749    /// `startup_order` (lower values start first).
750    ///
751    /// Algorithm lives in `lifecycle::application::context_lifecycle::start_context`
752    /// (Tier C C2). Public signature is unchanged.
753    pub async fn start(&mut self) -> Result<(), CamelError> {
754        crate::lifecycle::application::context_lifecycle::start_context(
755            &mut self.services,
756            &mut self.startup_checks,
757            &self.runtime,
758            &self.route_controller,
759            &mut self.cancel_token,
760        )
761        .await?;
762        // Trip the intercept-rules freeze so it applies even with zero
763        // routes. A failed `start_context` above returns early — a failed
764        // start does not freeze.
765        self.route_controller.mark_started().await
766    }
767
768    /// Graceful shutdown with default 30-second timeout.
769    pub async fn stop(&mut self) -> Result<(), CamelError> {
770        self.stop_timeout(self.shutdown_timeout).await
771    }
772
773    /// Graceful shutdown with custom timeout.
774    ///
775    /// Note: The timeout parameter is currently not propagated to the
776    /// RouteController's per-route shutdown timeout. The RouteController
777    /// uses a hardcoded 5-second default (`DEFAULT_SHUTDOWN_TIMEOUT`).
778    /// Full propagation is planned for a future version.
779    ///
780    /// Algorithm lives in `lifecycle::application::context_lifecycle::stop_context`
781    /// (Tier C C2). Public signature is unchanged.
782    pub async fn stop_timeout(&mut self, _timeout: std::time::Duration) -> Result<(), CamelError> {
783        crate::lifecycle::application::context_lifecycle::stop_context(
784            &self.cancel_token,
785            &mut self.supervision_join,
786            &self.runtime,
787            &self.route_controller,
788            &mut self.services,
789        )
790        .await
791    }
792
793    /// Get the graceful shutdown timeout used by [`stop()`](Self::stop).
794    pub fn shutdown_timeout(&self) -> std::time::Duration {
795        self.shutdown_timeout
796    }
797
798    /// Set the graceful shutdown timeout used by [`stop()`](Self::stop).
799    pub fn set_shutdown_timeout(&mut self, timeout: std::time::Duration) {
800        self.shutdown_timeout = timeout;
801    }
802
803    /// Test-only: take the actor join handle out of the context.
804    /// Used to verify the actor exits gracefully after stop().
805    #[cfg(test)]
806    pub(crate) fn take_actor_join(&mut self) -> Option<tokio::task::JoinHandle<()>> {
807        self.actor_join.take()
808    }
809
810    /// Immediate abort — kills all tasks without draining.
811    ///
812    /// Algorithm lives in `lifecycle::application::context_lifecycle::abort_context`
813    /// (Tier C C2). Public signature is unchanged. The use-case routes
814    /// through `RouteOrderingPort` + `RouteDestructiveTeardownPort`; the
815    /// same underlying `RouteControllerHandle` is passed twice as both
816    /// trait objects.
817    pub async fn abort(&mut self) {
818        crate::lifecycle::application::context_lifecycle::abort_context(
819            &self.cancel_token,
820            &mut self.supervision_join,
821            &self.runtime,
822            &self.route_controller as &dyn crate::lifecycle::application::ports::RouteOrderingPort,
823            &self.route_controller
824                as &dyn crate::lifecycle::application::ports::RouteDestructiveTeardownPort,
825            &mut self.services,
826            self.health_registry.cancel_token(),
827            &mut self.actor_join,
828        )
829        .await
830    }
831
832    /// Check health status of all registered services and lifecycle services.
833    pub async fn health_check(&self) -> HealthReport {
834        use camel_api::HealthSource;
835        self.health_report().await
836    }
837
838    pub fn health_registry(&self) -> Arc<HealthCheckRegistry> {
839        Arc::clone(&self.health_registry)
840    }
841
842    /// Store a component config. Overwrites any previously stored config of the same type.
843    pub fn set_component_config<T: 'static + Send + Sync>(&mut self, config: T) {
844        self.component_configs
845            .insert(TypeId::of::<T>(), Box::new(config));
846    }
847
848    /// Retrieve a stored component config by type. Returns None if not stored.
849    pub fn get_component_config<T: 'static + Send + Sync>(&self) -> Option<&T> {
850        self.component_configs
851            .get(&TypeId::of::<T>())
852            .and_then(|b| b.downcast_ref::<T>())
853    }
854
855    // --- Component Metadata ---
856
857    /// Get a component's harvested metadata by URI scheme.
858    pub fn component_metadata(&self, scheme: &str) -> Option<ComponentMetadata> {
859        self.registry.lock().ok()?.get_metadata(scheme)
860    }
861
862    /// Get metadata for every registered component.
863    pub fn all_component_metadata(&self) -> Vec<ComponentMetadata> {
864        self.registry
865            .lock()
866            .expect("mutex poisoned: another thread panicked while holding this lock") // allow-unwrap
867            .all_metadata()
868    }
869
870    /// Get a metadata catalog handle implementing
871    /// [`ComponentMetadataCatalog`](camel_api::component_metadata::ComponentMetadataCatalog).
872    ///
873    /// The returned handle shares the same `Arc<Mutex<Registry>>` as the
874    /// context, so registrations made through one are visible through the
875    /// other.
876    pub fn metadata_catalog(
877        &self,
878    ) -> crate::component_metadata_catalog::RuntimeComponentMetadataCatalog {
879        crate::component_metadata_catalog::RuntimeComponentMetadataCatalog::new(Arc::clone(
880            &self.registry,
881        ))
882    }
883
884    // --- Route Template Registry (data-only) ---
885
886    /// Register a route template specification.
887    ///
888    /// Returns `Err(CamelError)` if a template with the same ID is already registered.
889    pub fn add_route_template(&self, spec: RouteTemplateSpec) -> Result<(), CamelError> {
890        self.template_registry.register(spec)
891    }
892
893    /// Retrieve a route template specification by its ID.
894    pub fn get_route_template(&self, id: &str) -> Option<RouteTemplateSpec> {
895        self.template_registry.get(id)
896    }
897
898    /// Return all registered template IDs.
899    pub fn template_ids(&self) -> Vec<String> {
900        self.template_registry.template_ids()
901    }
902
903    /// Record a newly instantiated template instance.
904    pub fn record_template_instance(&self, record: TemplateInstanceRecord) {
905        self.template_registry.record_instance(record)
906    }
907
908    /// Return all instance records for a given template ID.
909    pub fn template_instances(&self, template_id: &str) -> Vec<TemplateInstanceRecord> {
910        self.template_registry.instances(template_id)
911    }
912
913    // --- Idempotent Repository Registry ---
914
915    /// Register an idempotent repository.
916    ///
917    /// Returns `Err(RegistryError::AlreadyRegistered)` if a repository with
918    /// the same name is already registered.
919    pub fn register_idempotent_repository(
920        &mut self,
921        name: impl Into<String>,
922        repo: Arc<dyn camel_api::IdempotentRepository>,
923    ) -> Result<(), RegistryError> {
924        self.idempotent_repositories.register(name, repo)
925    }
926
927    /// Retrieve an idempotent repository by name.
928    pub fn idempotent_repository(
929        &self,
930        name: &str,
931    ) -> Option<Arc<dyn camel_api::IdempotentRepository>> {
932        self.idempotent_repositories.get(name)
933    }
934
935    // --- Claim Check Repository Registry ---
936
937    /// Register a claim check repository.
938    ///
939    /// Returns `Err(RegistryError::AlreadyRegistered)` if a repository with
940    /// the same name is already registered.
941    pub fn register_claim_check_repository(
942        &mut self,
943        name: impl Into<String>,
944        repo: Arc<dyn camel_api::ClaimCheckRepository>,
945    ) -> Result<(), RegistryError> {
946        self.claim_check_repositories.register(name, repo)
947    }
948
949    /// Retrieve a claim check repository by name.
950    pub fn claim_check_repository(
951        &self,
952        name: &str,
953    ) -> Option<Arc<dyn camel_api::ClaimCheckRepository>> {
954        self.claim_check_repositories.get(name)
955    }
956
957    // --- Cache Repository Registry ---
958
959    /// Register a cache repository.
960    ///
961    /// Returns `Err(RegistryError::AlreadyRegistered)` if a repository with
962    /// the same name is already registered.
963    pub fn register_cache_repository(
964        &mut self,
965        name: impl Into<String>,
966        repo: Arc<dyn camel_api::CacheRepository>,
967    ) -> Result<(), RegistryError> {
968        self.cache_repositories.register(name, repo)
969    }
970
971    /// Replace an existing cache repository, returning the evicted value.
972    ///
973    /// Returns `None` if no repository was registered under `name`.
974    pub fn replace_cache_repository(
975        &mut self,
976        name: impl Into<String>,
977        repo: Arc<dyn camel_api::CacheRepository>,
978    ) -> Option<Arc<dyn camel_api::CacheRepository>> {
979        self.cache_repositories.register_or_replace(name, repo)
980    }
981
982    /// Retrieve a cache repository by name.
983    pub fn cache_repository(&self, name: &str) -> Option<Arc<dyn camel_api::CacheRepository>> {
984        self.cache_repositories.get(name)
985    }
986
987    /// Access the shutdown cancellation token.
988    ///
989    /// Repositories that need to bind background sweep tasks to context
990    /// shutdown (e.g. `RedbCacheRepository`) can `child_token()` from this.
991    pub fn shutdown_token(&self) -> CancellationToken {
992        self.cancel_token.clone()
993    }
994}
995
996/// Dropping the context is a NON-graceful termination of the controller
997/// actor and supervision tasks: both background tasks are aborted, so an
998/// outstanding controller command is interrupted at its await point (the
999/// same cancellation class as [`abort()`](Self::abort); ADR-0018
1000/// sequencing is not extended to the drop path). Route and service
1001/// teardown remains the caller's [`stop()`](Self::stop) responsibility.
1002/// [`stop()`](Self::stop) → [`start()`](Self::start) restart semantics are
1003/// unchanged, because Drop fires only when the context value itself is
1004/// discarded — a stopped-but-alive context keeps its actor for restart.
1005/// If a join handle was already taken (e.g. via the test-only
1006/// `take_actor_join`), the corresponding
1007/// `Option` is `None` and Drop skips it.
1008impl Drop for CamelContext {
1009    fn drop(&mut self) {
1010        if let Some(handle) = self.actor_join.take() {
1011            handle.abort();
1012        }
1013        if let Some(handle) = self.supervision_join.take() {
1014            handle.abort();
1015        }
1016    }
1017}
1018
1019impl ComponentRegistrar for CamelContext {
1020    fn register_component_dyn(&mut self, component: Arc<dyn Component>) {
1021        let scheme = component.scheme().to_string();
1022        self.registry
1023            .lock()
1024            .expect("mutex poisoned: another thread panicked while holding this lock") // allow-unwrap
1025            .register(component);
1026        trace!(scheme, "Registered component");
1027    }
1028}
1029
1030impl ComponentContext for CamelContext {
1031    fn resolve_component(&self, scheme: &str) -> Option<Arc<dyn Component>> {
1032        self.registry.lock().ok()?.get(scheme)
1033    }
1034
1035    fn resolve_language(&self, name: &str) -> Option<Arc<dyn Language>> {
1036        self.languages.lock().ok()?.get(name).cloned()
1037    }
1038
1039    fn metrics(&self) -> Arc<dyn MetricsCollector> {
1040        Arc::clone(&self.metrics) as Arc<dyn MetricsCollector>
1041    }
1042
1043    /// Snapshot of the `[observability.metrics].components` lever
1044    /// (dashboard-observability Task 4.1): gates only the uniform
1045    /// component-operations family offered through
1046    /// `RuntimeObservability::component_metrics()`. Error-family
1047    /// emission is never lever-gated, so it is not part of this flag.
1048    fn component_metrics_enabled(&self) -> bool {
1049        self.metrics_levers.components_enabled()
1050    }
1051
1052    fn health(&self) -> Arc<dyn camel_component_api::HealthCheckRegistry> {
1053        // The concrete HealthCheckRegistry struct implements the trait via
1054        // the impl added in health_registry.rs.
1055        Arc::clone(&self.health_registry) as Arc<dyn camel_component_api::HealthCheckRegistry>
1056    }
1057
1058    fn platform_service(&self) -> Arc<dyn PlatformService> {
1059        Arc::clone(&self.platform_service)
1060    }
1061
1062    fn register_route_health_check(
1063        &self,
1064        route_id: &str,
1065        check: Arc<dyn camel_api::AsyncHealthCheck>,
1066    ) {
1067        self.health_registry.register_for_route(route_id, check);
1068    }
1069
1070    fn unregister_route_health_check(&self, route_id: &str) {
1071        self.health_registry.unregister_for_route(route_id);
1072    }
1073
1074    /// The context-global accepted-not-completed gauge (drainclaim):
1075    /// producers created through this context capture it once at
1076    /// `create_producer` and mint `InFlightClaim`s against it.
1077    fn in_flight_counter(&self) -> Option<Arc<InFlightGauge>> {
1078        Some(Arc::clone(&self.in_flight_total))
1079    }
1080}
1081
1082#[async_trait::async_trait]
1083impl camel_api::HealthSource for CamelContext {
1084    async fn liveness(&self) -> camel_api::HealthStatus {
1085        let has_failed = self
1086            .services
1087            .iter()
1088            .any(|s| s.status() == camel_api::ServiceStatus::Failed);
1089        if has_failed {
1090            camel_api::HealthStatus::Unhealthy
1091        } else {
1092            camel_api::HealthStatus::Healthy
1093        }
1094    }
1095
1096    async fn readiness(&self) -> camel_api::HealthStatus {
1097        let has_failed = self
1098            .services
1099            .iter()
1100            .any(|s| s.status() == camel_api::ServiceStatus::Failed);
1101        if has_failed {
1102            return camel_api::HealthStatus::Unhealthy;
1103        }
1104        let has_stopped = self
1105            .services
1106            .iter()
1107            .any(|s| s.status() == camel_api::ServiceStatus::Stopped);
1108        if has_stopped {
1109            return camel_api::HealthStatus::Degraded;
1110        }
1111        self.health_registry.check_all().await.status
1112    }
1113
1114    async fn health_report(&self) -> camel_api::HealthReport {
1115        let mut report = self.health_registry.check_all().await;
1116        let mut worst = report.status;
1117        for service in &self.services {
1118            let svc_status = service.status();
1119            let health = match svc_status {
1120                camel_api::ServiceStatus::Started => camel_api::HealthStatus::Healthy,
1121                camel_api::ServiceStatus::Stopped => camel_api::HealthStatus::Degraded,
1122                camel_api::ServiceStatus::Failed => camel_api::HealthStatus::Unhealthy,
1123                // Forward-safe fail-closed: an unknown future ServiceStatus keeps
1124                // the pod NotReady (no traffic) rather than guessing Degraded.
1125                _ => camel_api::HealthStatus::Unhealthy,
1126            };
1127            if matches!(worst, camel_api::HealthStatus::Healthy)
1128                && matches!(
1129                    health,
1130                    camel_api::HealthStatus::Degraded | camel_api::HealthStatus::Unhealthy
1131                )
1132            {
1133                worst = health;
1134            }
1135            if matches!(worst, camel_api::HealthStatus::Degraded)
1136                && matches!(health, camel_api::HealthStatus::Unhealthy)
1137            {
1138                worst = health;
1139            }
1140            report.services.push(camel_api::ServiceHealth {
1141                name: service.name().to_string(),
1142                status: svc_status,
1143                message: None,
1144            });
1145        }
1146        report.status = worst;
1147        report
1148    }
1149
1150    async fn startup(&self) -> camel_api::HealthStatus {
1151        camel_api::HealthStatus::Healthy
1152    }
1153}
1154
1155#[cfg(test)]
1156#[path = "context_tests.rs"]
1157mod context_tests;