Skip to main content

camel_core/lifecycle/adapters/
route_controller.rs

1//! Default implementation of RouteController.
2//!
3//! This module provides [`DefaultRouteController`], which manages route lifecycle
4//! including starting, stopping, suspending, and resuming routes.
5
6use std::collections::HashMap;
7use std::sync::{Arc, Weak};
8use std::time::Duration;
9
10use tokio::sync::mpsc;
11use tokio_util::sync::CancellationToken;
12use tower::{Layer, ServiceExt};
13use tracing::{debug, info, warn};
14
15use super::route_controller_trait::{
16    BindExposureAcks, bind_key_from_uri, enforce_bind_exposure_gate,
17};
18use camel_api::error_handler::ErrorHandlerConfig;
19use camel_api::metrics::MetricsCollector;
20#[allow(unused_imports)]
21use camel_api::{
22    BoxProcessor, CamelError, Exchange, FunctionInvoker, IdentityProcessor, InFlightClaim,
23    InFlightGauge, NoOpMetrics, NoopPlatformService, PlatformService, ProducerContext,
24    RouteController, RuntimeHandle, StepLifecycle,
25};
26use camel_component_api::{Consumer, ConsumerContext, consumer::ExchangeEnvelope};
27use camel_processor::aggregator::AggregatorService;
28pub use camel_processor::aggregator::SharedLanguageRegistry;
29use camel_processor::aggregator::{AggregateEmission, AggregationReceipt};
30
31use crate::health_registry::HealthCheckRegistry;
32use crate::intercept::InterceptRules;
33use crate::lifecycle::CohortActivationGate;
34use crate::lifecycle::adapters::controller_component_context::ControllerComponentContext;
35use crate::lifecycle::adapters::route_compiler::TracerPipelineGating;
36use crate::lifecycle::adapters::route_compiler_ext::{
37    RouteCompilerExt, build_eh_config_pipeline, transport_from_uri,
38};
39use crate::lifecycle::adapters::route_helpers::{
40    AggregateSplitInfo, CrashNotification, ManagedRoute, assert_no_mixed_top_level_splits,
41    handle_is_running, inferred_lifecycle_label, is_pending, send_reply_or_b_prime,
42};
43#[cfg(test)]
44pub(super) use crate::lifecycle::adapters::route_helpers::{
45    emit_start_route_event, set_start_route_event_hook,
46};
47use crate::lifecycle::adapters::route_registry::RouteRegistry;
48use crate::lifecycle::adapters::route_runtime_state;
49use crate::lifecycle::adapters::step_compilers::CompiledStep;
50use crate::lifecycle::application::route_definition::{BuilderStep, RouteDefinition};
51pub(crate) use crate::lifecycle::domain::CompiledPipeline;
52use crate::shared::components::domain::Registry;
53use crate::shared::observability::domain::{DetailLevel, TracerConfig};
54use camel_bean::BeanRegistry;
55
56/// Default implementation of [`RouteController`].
57///
58/// Manages route lifecycle with support for:
59/// - Starting/stopping individual routes
60/// - Suspending and resuming routes
61/// - Auto-startup with startup ordering
62/// - Graceful shutdown
63pub struct DefaultRouteController {
64    /// Routes indexed by route ID.
65    pub(super) routes: RouteRegistry,
66    /// Reference to the component registry for resolving endpoints.
67    pub(super) registry: Arc<std::sync::Mutex<Registry>>,
68    /// Shared language registry for resolving declarative language expressions.
69    pub(super) languages: SharedLanguageRegistry,
70    /// Bean registry for bean method invocation.
71    pub(super) beans: Arc<std::sync::Mutex<BeanRegistry>>,
72    /// Runtime handle injected into ProducerContext for command/query operations.
73    pub(super) runtime: Option<Weak<dyn RuntimeHandle>>,
74    /// Optional global error handler applied to all routes without a per-route handler.
75    pub(super) global_error_handler: Option<ErrorHandlerConfig>,
76    /// Optional crash notifier for supervision.
77    pub(super) crash_notifier: Option<mpsc::Sender<CrashNotification>>,
78    /// Whether tracing is enabled for route pipelines.
79    pub(super) tracer_gating: TracerPipelineGating,
80    /// Detail level for tracing when enabled.
81    pub(super) tracer_detail_level: DetailLevel,
82    /// Metrics collector for tracing processor.
83    pub(super) tracer_metrics: Option<Arc<dyn MetricsCollector>>,
84    pub(super) platform_service: Arc<dyn PlatformService>,
85    pub(super) function_invoker: Option<Arc<dyn FunctionInvoker>>,
86    pub(super) health_registry: Option<Arc<HealthCheckRegistry>>,
87    /// Shared idempotent repository registry. Defaults to an empty registry;
88    /// the CamelContext builder installs a populated handle that includes the
89    /// built-in `"memory"` repository.
90    pub(super) idempotent_repositories: crate::SharedIdempotentRegistry,
91    pub(super) claim_check_repositories: crate::SharedClaimCheckRegistry,
92    pub(super) cache_repositories: crate::SharedCacheRegistry,
93    /// F2 staging: prepared-but-not-inserted ManagedRoutes keyed by route_id.
94    /// `prepare_*` writes here; `insert_prepared_route` drains via `remove()`.
95    /// On insert-failure error paths, the caller (`reload_actions.rs`) must
96    /// explicitly drain to avoid orphan CancellationToken/SharedPipeline leaks.
97    pub(super) prepared_staging: HashMap<String, ManagedRoute>,
98    /// Source endpoint URI to route_id index (one-to-many).
99    pub(super) endpoint_index: super::endpoint_index::EndpointIndex,
100    /// Operator acknowledgements for per-bind public exposure (ADR-0061).
101    /// Empty by default → the gate fails closed on non-loopback binds.
102    pub(super) bind_acks: super::route_controller_trait::BindExposureAcks,
103    /// Route send-point interception rules captured by step compilation.
104    pub(super) intercept: InterceptRules,
105    /// Intercept-rules freeze. Trips on `add_route` success and on the
106    /// `MarkStarted` actor command; never reset (stop/restart included),
107    /// because compiled pipelines capture the rules at compile time.
108    pub(super) frozen: bool,
109    /// Startup-cohort activation barrier (rc-jxkj). Cloned into the
110    /// `RouteControllerHandle` at spawn; reset/activate act on this shared
111    /// gate directly, never through the actor.
112    pub(super) cohort: Arc<CohortActivationGate>,
113    /// Context-global accepted-not-completed gauge (drainclaim): the
114    /// SAME `Arc` the owning `CamelContext` exposes through
115    /// `total_in_flight()` (installed by the builder). Standalone
116    /// controllers keep an isolated zero gauge — claims flow, but only
117    /// the constructing scope can read them.
118    pub(super) in_flight_total: Arc<InFlightGauge>,
119}
120
121impl DefaultRouteController {
122    /// Open the startup-cohort activation barrier (rc-jxkj).
123    ///
124    /// The CamelContext lifecycle opens this automatically once the startup
125    /// cohort completes. Consumers that drive a bare `DefaultRouteController`
126    /// (outside a full context) must call this before dispatching
127    /// (typically after starting routes), or pipeline dispatch parks
128    /// every envelope until the caller's call timeout surfaces as a
129    /// failure.
130    pub fn activate_cohort(&self) {
131        self.cohort.open();
132    }
133
134    pub(super) fn health_registry(&self) -> Arc<HealthCheckRegistry> {
135        self.health_registry.clone().unwrap_or_else(|| {
136            debug!("health_registry not configured — creating isolated fallback");
137            Arc::new(HealthCheckRegistry::new(Duration::from_secs(5)))
138        })
139    }
140
141    /// Create a new `DefaultRouteController` with the given registry.
142    pub fn new(
143        registry: Arc<std::sync::Mutex<Registry>>,
144        platform_service: Arc<dyn PlatformService>,
145    ) -> Self {
146        Self::with_beans_and_platform_service(
147            registry,
148            Arc::new(std::sync::Mutex::new(BeanRegistry::new())),
149            platform_service,
150        )
151    }
152
153    /// Create a new `DefaultRouteController` with shared bean registry.
154    pub fn with_beans(
155        registry: Arc<std::sync::Mutex<Registry>>,
156        beans: Arc<std::sync::Mutex<BeanRegistry>>,
157    ) -> Self {
158        Self::with_beans_and_platform_service(
159            registry,
160            beans,
161            Arc::new(NoopPlatformService::default()),
162        )
163    }
164
165    fn with_beans_and_platform_service(
166        registry: Arc<std::sync::Mutex<Registry>>,
167        beans: Arc<std::sync::Mutex<BeanRegistry>>,
168        platform_service: Arc<dyn PlatformService>,
169    ) -> Self {
170        Self {
171            routes: RouteRegistry::new(),
172            registry,
173            languages: Arc::new(std::sync::Mutex::new(HashMap::new())),
174            beans,
175            runtime: None,
176            global_error_handler: None,
177            crash_notifier: None,
178            tracer_gating: TracerPipelineGating::off(),
179            tracer_detail_level: DetailLevel::Minimal,
180            tracer_metrics: None,
181            platform_service,
182            function_invoker: None,
183            health_registry: None,
184            idempotent_repositories: Arc::new(crate::IdempotentRegistry::new()),
185            claim_check_repositories: Arc::new(crate::ClaimCheckRegistry::new()),
186            cache_repositories: Arc::new(crate::CacheRegistry::new()),
187            prepared_staging: HashMap::new(),
188            endpoint_index: super::endpoint_index::EndpointIndex::new(),
189            bind_acks: Default::default(),
190            intercept: InterceptRules::default(),
191            frozen: false,
192            cohort: Arc::new(CohortActivationGate::new_closed()),
193            in_flight_total: Arc::new(InFlightGauge::new()),
194        }
195    }
196
197    /// Create a new `DefaultRouteController` with shared language registry.
198    pub fn with_languages(
199        registry: Arc<std::sync::Mutex<Registry>>,
200        languages: SharedLanguageRegistry,
201        platform_service: Arc<dyn PlatformService>,
202    ) -> Self {
203        Self {
204            routes: RouteRegistry::new(),
205            registry,
206            languages,
207            beans: Arc::new(std::sync::Mutex::new(BeanRegistry::new())),
208            runtime: None,
209            global_error_handler: None,
210            crash_notifier: None,
211            tracer_gating: TracerPipelineGating::off(),
212            tracer_detail_level: DetailLevel::Minimal,
213            tracer_metrics: None,
214            platform_service,
215            function_invoker: None,
216            health_registry: None,
217            idempotent_repositories: Arc::new(crate::IdempotentRegistry::new()),
218            claim_check_repositories: Arc::new(crate::ClaimCheckRegistry::new()),
219            cache_repositories: Arc::new(crate::CacheRegistry::new()),
220            prepared_staging: HashMap::new(),
221            endpoint_index: super::endpoint_index::EndpointIndex::new(),
222            bind_acks: Default::default(),
223            intercept: InterceptRules::default(),
224            frozen: false,
225            cohort: Arc::new(CohortActivationGate::new_closed()),
226            in_flight_total: Arc::new(InFlightGauge::new()),
227        }
228    }
229
230    pub fn with_languages_and_beans(
231        registry: Arc<std::sync::Mutex<Registry>>,
232        languages: SharedLanguageRegistry,
233        platform_service: Arc<dyn PlatformService>,
234        beans: Arc<std::sync::Mutex<BeanRegistry>>,
235    ) -> Self {
236        Self {
237            routes: RouteRegistry::new(),
238            registry,
239            languages,
240            beans,
241            runtime: None,
242            global_error_handler: None,
243            crash_notifier: None,
244            tracer_gating: TracerPipelineGating::off(),
245            tracer_detail_level: DetailLevel::Minimal,
246            tracer_metrics: None,
247            platform_service,
248            function_invoker: None,
249            health_registry: None,
250            idempotent_repositories: Arc::new(crate::IdempotentRegistry::new()),
251            claim_check_repositories: Arc::new(crate::ClaimCheckRegistry::new()),
252            cache_repositories: Arc::new(crate::CacheRegistry::new()),
253            prepared_staging: HashMap::new(),
254            endpoint_index: super::endpoint_index::EndpointIndex::new(),
255            bind_acks: Default::default(),
256            intercept: InterceptRules::default(),
257            frozen: false,
258            cohort: Arc::new(CohortActivationGate::new_closed()),
259            in_flight_total: Arc::new(InFlightGauge::new()),
260        }
261    }
262
263    pub fn with_function_invoker(mut self, function_invoker: Arc<dyn FunctionInvoker>) -> Self {
264        self.function_invoker = Some(function_invoker);
265        self
266    }
267
268    pub(crate) fn set_idempotent_repositories(
269        &mut self,
270        repositories: crate::SharedIdempotentRegistry,
271    ) {
272        self.idempotent_repositories = repositories;
273    }
274
275    pub(crate) fn set_claim_check_repositories(
276        &mut self,
277        repositories: crate::SharedClaimCheckRegistry,
278    ) {
279        self.claim_check_repositories = repositories;
280    }
281
282    pub(crate) fn set_cache_repositories(&mut self, repositories: crate::SharedCacheRegistry) {
283        self.cache_repositories = repositories;
284    }
285
286    pub fn set_health_registry(&mut self, registry: Arc<HealthCheckRegistry>) {
287        self.health_registry = Some(registry);
288    }
289
290    pub fn set_function_invoker(&mut self, invoker: Arc<dyn FunctionInvoker>) {
291        self.function_invoker = Some(invoker);
292    }
293
294    /// Set runtime handle for ProducerContext creation.
295    pub fn set_runtime_handle(&mut self, runtime: Arc<dyn RuntimeHandle>) {
296        self.runtime = Some(Arc::downgrade(&runtime));
297    }
298
299    /// Set the crash notifier for supervision.
300    ///
301    /// When set, the controller will send a `CrashNotification` whenever
302    /// a consumer crashes.
303    pub fn set_crash_notifier(&mut self, tx: mpsc::Sender<CrashNotification>) {
304        self.crash_notifier = Some(tx);
305    }
306
307    /// Set a global error handler applied to all routes without a per-route handler.
308    /// Install operator acknowledgements for per-bind public exposure
309    /// (ADR-0061). Built by the CLI from `CamelConfig.binds`.
310    pub fn set_bind_exposure_acks(&mut self, acks: BindExposureAcks) {
311        self.bind_acks = acks;
312    }
313
314    /// Install route send-point interception rules (pre-first-use only).
315    ///
316    /// Fails with `CamelError::Config` once the freeze has tripped: compiled
317    /// pipelines capture the rules at compile time, so the rule set must not
318    /// change after a route is added or the context is started.
319    pub fn set_intercept_rules(&mut self, rules: InterceptRules) -> Result<(), CamelError> {
320        if self.frozen {
321            return Err(CamelError::Config(
322                "intercept rules are frozen: a route was added or the context was started; \
323                 rules cannot be changed after first use"
324                    .into(),
325            ));
326        }
327        self.intercept = rules;
328        Ok(())
329    }
330
331    /// Builder-style build-time configuration on a fresh controller (which
332    /// is never frozen); mirrors `with_function_invoker`.
333    pub fn with_intercept_rules(mut self, rules: InterceptRules) -> Self {
334        self.intercept = rules;
335        self
336    }
337
338    /// Trip the intercept-rules freeze. Dispatched by the `MarkStarted`
339    /// actor command so the freeze applies even with zero routes. Never
340    /// unset.
341    pub fn mark_started(&mut self) {
342        self.frozen = true;
343    }
344
345    /// All compiled security plans whose routes bind the same listener
346    /// address (bind key), for the per-bind exposure gate. Routes still
347    /// staging (no plan yet) are skipped — classification failures already
348    /// aborted their own staging (Task 1.8).
349    pub(super) fn plans_for_bind(
350        &self,
351        bind_key: &str,
352    ) -> Vec<(String, camel_api::security_policy::RouteSecurityPlan)> {
353        self.routes
354            .iter()
355            .filter_map(|(route_id, managed)| {
356                let plan = managed.compiled.security_plan.as_ref()?;
357                let bind = bind_key_from_uri(&managed.from_uri)?;
358                if bind.key == bind_key {
359                    Some((route_id.clone(), plan.clone()))
360                } else {
361                    None
362                }
363            })
364            .collect()
365    }
366
367    /// Return whether a listener for `bind_key` is already serving a route.
368    ///
369    /// This is deliberately based on the consumer handle, rather than on the
370    /// presence of a compiled route: the exposure gate applies only to late
371    /// registration against an already-running bind (Task 2.2). Startup and
372    /// resume perform their own gate checks in the lifecycle implementation.
373    fn bind_is_running(&self, bind_key: &str) -> bool {
374        self.routes.iter().any(|(_, managed)| {
375            bind_key_from_uri(&managed.from_uri).is_some_and(|bind| {
376                bind.key == bind_key && handle_is_running(&managed.consumer_handle)
377            })
378        })
379    }
380
381    /// Gate a route before it is inserted into a listener that is already
382    /// running. The candidate is included with all sibling plans so the gate
383    /// has the same aggregation semantics as the start path.
384    fn enforce_late_registration_gate(
385        &self,
386        from_uri: &str,
387        route_id: &str,
388        plan: Option<&camel_api::security_policy::RouteSecurityPlan>,
389    ) -> Result<(), CamelError> {
390        let Some(bind) = bind_key_from_uri(from_uri) else {
391            return Ok(());
392        };
393        if !self.bind_is_running(&bind.key) {
394            return Ok(());
395        }
396
397        let mut owned = self.plans_for_bind(&bind.key);
398        if let Some(plan) = plan {
399            owned.push((route_id.to_string(), plan.clone()));
400        }
401        let plans: Vec<(&str, &camel_api::security_policy::RouteSecurityPlan)> =
402            owned.iter().map(|(id, plan)| (id.as_str(), plan)).collect();
403        enforce_bind_exposure_gate(
404            &bind.key,
405            bind.loopback,
406            &plans,
407            self.bind_acks.acknowledged(&bind.key),
408        )
409        .map_err(|err| {
410            CamelError::RouteError(format!(
411                "late registration of route '{route_id}' rejected for bind '{}': {err}",
412                bind.key
413            ))
414        })
415    }
416
417    pub fn set_error_handler(&mut self, config: ErrorHandlerConfig) {
418        self.global_error_handler = Some(config);
419    }
420
421    /// Configure tracing for this route controller.
422    pub fn set_tracer_config(&mut self, config: &TracerConfig) {
423        // Pipeline wrapping follows tracing unless the effective assembly
424        // raised it (exporter active); spans follow `enabled` alone.
425        self.tracer_gating = TracerPipelineGating {
426            pipeline_enabled: config.enabled || config.pipeline_enabled,
427            spans_enabled: config.enabled,
428            levers: config.metrics_levers.clone(),
429        };
430        self.tracer_detail_level = config.detail_level.clone();
431    }
432
433    /// Seed the tracer metrics collector — the shared late-bound
434    /// `MetricsHandle` built once by `CamelContextBuilder::build()`.
435    ///
436    /// Replaces the deleted `TracerConfig.metrics_collector` snapshot
437    /// injection: the collector is wired here, at construction, and late
438    /// registrations flow through the handle without re-snapshotting.
439    pub fn set_tracer_metrics(&mut self, metrics: Arc<dyn MetricsCollector>) {
440        self.tracer_metrics = Some(metrics);
441    }
442
443    /// Install the context-global accepted-not-completed gauge
444    /// (drainclaim) — the SAME `Arc<InFlightGauge>` the owning
445    /// `CamelContext` reads through `total_in_flight()`. Called once by
446    /// `CamelContextBuilder::build()`; contexts built without a builder
447    /// (or standalone controllers) keep the isolated zero gauge from
448    /// construction.
449    pub fn set_in_flight_total(&mut self, counter: Arc<InFlightGauge>) {
450        self.in_flight_total = counter;
451    }
452
453    fn build_producer_context(&self, route_id: &str) -> Result<ProducerContext, CamelError> {
454        let mut producer_ctx = ProducerContext::new().with_route_id(route_id);
455        if let Some(runtime) = self.runtime.as_ref().and_then(Weak::upgrade) {
456            producer_ctx = producer_ctx.with_runtime(runtime);
457        }
458        Ok(producer_ctx)
459    }
460
461    /// Create a transient [`RouteCompilerExt`] from this controller's fields.
462    fn route_compiler_ext(&self) -> RouteCompilerExt<'_> {
463        RouteCompilerExt {
464            registry: &self.registry,
465            languages: &self.languages,
466            beans: &self.beans,
467            function_invoker: &self.function_invoker,
468            tracer_gating: self.tracer_gating.clone(),
469            tracer_detail_level: &self.tracer_detail_level,
470            tracer_metrics: &self.tracer_metrics,
471            platform_service: &self.platform_service,
472            runtime: &self.runtime,
473            global_error_handler: &self.global_error_handler,
474            health_registry: &self.health_registry,
475            route_registry: &self.routes,
476            idempotent_repositories: Arc::clone(&self.idempotent_repositories),
477            claim_check_repositories: Arc::clone(&self.claim_check_repositories),
478            cache_repositories: Arc::clone(&self.cache_repositories),
479            intercept: &self.intercept,
480            in_flight_total: Arc::clone(&self.in_flight_total),
481        }
482    }
483
484    /// Resolve BuilderSteps into BoxProcessors.
485    #[allow(dead_code)] // used by tests and may be needed for future split paths
486    pub(crate) fn resolve_steps(
487        &self,
488        steps: Vec<BuilderStep>,
489        producer_ctx: &ProducerContext,
490        registry: &Arc<std::sync::Mutex<Registry>>,
491        route_id: Option<&str>,
492        staging_mode: &super::step_resolution::FunctionStagingMode,
493    ) -> Result<Vec<CompiledStep>, CamelError> {
494        let component_ctx = Arc::new(ControllerComponentContext::new(
495            Arc::clone(registry),
496            Arc::clone(&self.languages),
497            self.tracer_metrics
498                .clone()
499                .unwrap_or_else(|| Arc::new(NoOpMetrics)),
500            Arc::clone(&self.platform_service),
501            self.health_registry(),
502            route_id.map(|s| s.to_string()),
503            self.tracer_gating.levers.components_enabled(),
504        ));
505        let rt: Arc<dyn camel_component_api::RuntimeObservability> =
506            Arc::clone(&component_ctx) as Arc<_>;
507
508        super::step_resolution::resolve_steps(
509            steps,
510            producer_ctx,
511            rt,
512            registry,
513            &self.languages,
514            &self.beans,
515            self.function_invoker.clone(),
516            component_ctx,
517            route_id,
518            staging_mode,
519            &self.idempotent_repositories,
520            &self.claim_check_repositories,
521            &self.cache_repositories,
522            self.intercept.clone(),
523        )
524    }
525
526    /// Add a route definition to the controller.
527    ///
528    /// Steps are resolved immediately using the registry.
529    ///
530    /// # Errors
531    ///
532    /// Returns an error if:
533    /// - A route with the same ID already exists
534    /// - Step resolution fails
535    pub async fn add_route(&mut self, definition: RouteDefinition) -> Result<(), CamelError> {
536        let route_id = definition.route_id().to_string();
537        let from_uri = definition.from_uri().to_string();
538
539        if self.routes.contains_key(&route_id) {
540            return Err(CamelError::RouteError(format!(
541                "duplicate route ID '{route_id}'"
542            )));
543        }
544
545        debug!(route_id = %route_id, "Adding route to controller");
546
547        let managed = match self.build_managed_route(
548            definition,
549            &super::step_resolution::FunctionStagingMode::DirectAdd,
550        ) {
551            Ok(managed) => managed,
552            Err(err) => {
553                self.discard_function_staging();
554                return Err(err);
555            }
556        };
557
558        // A running listener can observe a newly inserted route immediately.
559        // Gate the complete sibling set before committing function staging or
560        // inserting the route, so a rejected late registration is unreachable.
561        if let Err(err) = self.enforce_late_registration_gate(
562            &from_uri,
563            &route_id,
564            managed.compiled.security_plan.as_ref(),
565        ) {
566            self.discard_function_staging();
567            return Err(err);
568        }
569
570        if let Some(invoker) = &self.function_invoker
571            && let Err(err) = invoker.commit_staged().await
572        {
573            invoker.discard_staging(0);
574            return Err(CamelError::Config(err.to_string()));
575        }
576
577        self.routes
578            .insert(managed.definition.route_id().to_string(), managed);
579
580        self.endpoint_index.insert(&from_uri, &route_id);
581        // First successful route registration freezes the intercept rules:
582        // this route's pipeline compiled against the current rule set.
583        self.frozen = true;
584        Ok(())
585    }
586
587    pub(super) fn build_managed_route(
588        &self,
589        definition: RouteDefinition,
590        staging_mode: &super::step_resolution::FunctionStagingMode,
591    ) -> Result<ManagedRoute, CamelError> {
592        let route_id = definition.route_id().to_string();
593
594        let definition_info = definition.to_info();
595
596        // Security plan compilation (Task 1.8): consumer-backed routes get a
597        // plan BEFORE any consumer starts; a declared route that fails
598        // classification aborts staging (never a Public downgrade). Routes
599        // without a provider registry compile against an empty view.
600        let empty_providers;
601        let providers = match &definition.provider_registry {
602            Some(registry) => registry.as_ref(),
603            None => {
604                empty_providers = camel_auth::ProviderRegistry::new();
605                &empty_providers
606            }
607        };
608        let security_plan =
609            super::route_compiler_ext::compile_route_security_plan(&definition, providers)?;
610
611        let RouteDefinition {
612            from_uri,
613            steps,
614            error_handler,
615            circuit_breaker,
616            circuit_breaker_fallback,
617            security_policy,
618            security_authenticator,
619            provider_registry,
620            unit_of_work,
621            concurrency,
622            ..
623        } = definition;
624
625        let producer_ctx = self.build_producer_context(&route_id)?;
626
627        // N2: reject mixed Aggregate + Resequence top-level splits
628        assert_no_mixed_top_level_splits(&steps)?;
629
630        let (aggregate_split, processors_with_contracts) = self
631            .route_compiler_ext()
632            .detect_and_validate_route_split(steps, &producer_ctx, &route_id, staging_mode)?;
633        let mut lifecycle = super::route_helpers::collect_lifecycle(&processors_with_contracts);
634
635        // CB fallback (mirrors on_miss lifecycle packing): attach via the
636        // shared helper, then merge the fallback lifecycle handles into the
637        // route vec.
638        let (circuit_breaker, fallback_lifecycle) = self.route_compiler_ext().attach_cb_fallback(
639            circuit_breaker,
640            circuit_breaker_fallback,
641            &producer_ctx,
642            &route_id,
643            staging_mode,
644        )?;
645        lifecycle.extend(fallback_lifecycle);
646        let route_id_for_tracing = route_id.clone();
647        let eh_config = error_handler.or_else(|| self.global_error_handler.clone());
648        let transport = transport_from_uri(&from_uri);
649
650        let mut pipeline = build_eh_config_pipeline(
651            eh_config.as_ref(),
652            Arc::clone(&self.registry),
653            Arc::clone(&self.languages),
654            self.tracer_metrics.clone(),
655            Arc::clone(&self.platform_service),
656            self.health_registry(),
657            &route_id_for_tracing,
658            &producer_ctx,
659            processors_with_contracts,
660            self.tracer_gating.clone(),
661            self.tracer_detail_level.clone(),
662            security_policy.clone(),
663            transport,
664            circuit_breaker,
665            Arc::clone(&self.in_flight_total),
666        )?;
667
668        let uow_counter = if let Some(uow_config) = &unit_of_work {
669            let component_ctx = Arc::new(
670                ControllerComponentContext::new(
671                    Arc::clone(&self.registry),
672                    Arc::clone(&self.languages),
673                    self.tracer_metrics
674                        .clone()
675                        .unwrap_or_else(|| Arc::new(NoOpMetrics)),
676                    Arc::clone(&self.platform_service),
677                    self.health_registry(),
678                    Some(route_id.clone()),
679                    self.tracer_gating.levers.components_enabled(),
680                )
681                .with_in_flight(Arc::clone(&self.in_flight_total)),
682            );
683            let rt: Arc<dyn camel_component_api::RuntimeObservability> =
684                Arc::clone(&component_ctx) as Arc<_>;
685            let (uow_layer, counter) = super::route_compiler_ext::resolve_uow_layer(
686                uow_config,
687                &producer_ctx,
688                rt,
689                component_ctx.as_ref(),
690                None,
691            )?;
692            pipeline = BoxProcessor::new(uow_layer.layer(pipeline));
693            Some(counter)
694        } else {
695            None
696        };
697
698        Ok(ManagedRoute {
699            definition: definition_info,
700            from_uri,
701            pipeline: super::pipeline_runtime::new_shared_pipeline_with_lifecycle(
702                pipeline, lifecycle,
703            ),
704            concurrency,
705            consumer_handle: None,
706            pipeline_handle: None,
707            consumer_cancel_token: CancellationToken::new(),
708            pipeline_cancel_token: CancellationToken::new(),
709            channel_sender: None,
710            in_flight: uow_counter,
711            drain_in_flight: Arc::new(std::sync::atomic::AtomicU64::new(0)),
712            aggregate_split,
713            agg_service: None,
714            compiled: route_runtime_state::CompiledRoute {
715                security_policy,
716                security_authenticator,
717                provider_registry,
718                security_plan,
719            },
720        })
721    }
722
723    pub async fn add_route_with_generation(
724        &mut self,
725        definition: RouteDefinition,
726        generation: u64,
727    ) -> Result<(), CamelError> {
728        let route_id = definition.route_id().to_string();
729        let from_uri = definition.from_uri().to_string();
730
731        if self.routes.contains_key(&route_id) {
732            return Err(CamelError::RouteError(format!(
733                "duplicate route ID '{route_id}'"
734            )));
735        }
736
737        debug!(route_id = %route_id, generation, "Adding route to controller with generation");
738
739        let managed = self.build_managed_route(
740            definition,
741            &super::step_resolution::FunctionStagingMode::HotReload { generation },
742        )?;
743
744        // Symmetric with `add_route`: an already-running bind must never
745        // observe an ungated insertion through the hot-reload path.
746        if let Err(err) = self.enforce_late_registration_gate(
747            &from_uri,
748            &route_id,
749            managed.compiled.security_plan.as_ref(),
750        ) {
751            self.discard_function_staging();
752            return Err(err);
753        }
754
755        self.routes.insert(route_id.clone(), managed);
756
757        self.endpoint_index.insert(&from_uri, &route_id);
758        Ok(())
759    }
760
761    pub async fn remove_route_preserving_functions(
762        &mut self,
763        route_id: &str,
764    ) -> Result<(), CamelError> {
765        let managed = self.routes.get(route_id).ok_or_else(|| {
766            CamelError::RouteError(format!("Route '{}' not found for removal", route_id))
767        })?;
768        if handle_is_running(&managed.consumer_handle)
769            || handle_is_running(&managed.pipeline_handle)
770        {
771            return Err(CamelError::RouteError(format!(
772                "Route '{}' must be stopped before removal (current execution lifecycle: {})",
773                route_id,
774                inferred_lifecycle_label(managed)
775            )));
776        }
777        self.routes.remove(route_id);
778        if let Some(reg) = &self.health_registry {
779            reg.unregister_for_route(route_id);
780        }
781        self.endpoint_index.remove(route_id);
782        debug!(route_id = %route_id, "Route removed from controller (functions preserved for reload finalize)");
783        Ok(())
784    }
785
786    /// Compile a route definition into a processor pipeline, without adding it
787    /// to the controller. Used for validation and testing.
788    pub fn compile_route_definition(
789        &self,
790        def: RouteDefinition,
791    ) -> Result<BoxProcessor, CamelError> {
792        self.route_compiler_ext().compile_route_definition(def)
793    }
794
795    /// Compile a route definition with a specific generation (for hot-reload).
796    pub fn compile_route_definition_with_generation(
797        &self,
798        def: RouteDefinition,
799        generation: u64,
800    ) -> Result<BoxProcessor, CamelError> {
801        self.route_compiler_ext()
802            .compile_route_definition_with_generation(def, generation)
803    }
804
805    /// Compile a route definition into a [`CompiledPipeline`] (processor +
806    /// lifecycle handles). Used by the hot-reload Restart path so that
807    /// lifecycle handles are threaded through
808    /// [`swap_pipeline_raw`](Self::swap_pipeline_raw).
809    pub(crate) fn compile_route_definition_pipeline(
810        &self,
811        def: RouteDefinition,
812        generation: u64,
813    ) -> Result<CompiledPipeline, CamelError> {
814        self.route_compiler_ext()
815            .compile_route_definition_pipeline(def, generation)
816    }
817
818    /// Compile without function generation, returning full [`CompiledPipeline`].
819    ///
820    /// Oracle Fix 1: used by the stateless hot-reload path so that
821    /// lifecycle-bearing routes have their handles preserved.
822    pub(crate) fn compile_route_definition_dry_pipeline(
823        &self,
824        def: RouteDefinition,
825    ) -> Result<CompiledPipeline, CamelError> {
826        self.route_compiler_ext()
827            .compile_route_definition_dry_pipeline(def)
828    }
829
830    /// Remove a route from the controller map.
831    ///
832    /// The route **must** be stopped before removal (status `Stopped` or `Failed`).
833    /// Returns an error if the route is still running or does not exist.
834    /// Does not cancel any running tasks — call `stop_route` first.
835    pub async fn remove_route(&mut self, route_id: &str) -> Result<(), CamelError> {
836        let managed = self.routes.get(route_id).ok_or_else(|| {
837            CamelError::RouteError(format!("Route '{}' not found for removal", route_id))
838        })?;
839        if handle_is_running(&managed.consumer_handle)
840            || handle_is_running(&managed.pipeline_handle)
841        {
842            return Err(CamelError::RouteError(format!(
843                "Route '{}' must be stopped before removal (current execution lifecycle: {})",
844                route_id,
845                inferred_lifecycle_label(managed)
846            )));
847        }
848        if let Some(invoker) = &self.function_invoker {
849            for (id, rid) in self.collect_function_refs(route_id) {
850                if let Err(e) = invoker.unregister(&id, rid.as_deref()).await {
851                    warn!(route_id = %route_id, error = %e, "Failed to unregister function during route removal");
852                }
853            }
854        }
855        self.routes.remove(route_id);
856        if let Some(reg) = &self.health_registry {
857            reg.unregister_for_route(route_id);
858        }
859        self.endpoint_index.remove(route_id);
860        info!(route_id = %route_id, "Route removed from controller");
861        Ok(())
862    }
863
864    fn collect_function_refs(
865        &self,
866        route_id: &str,
867    ) -> Vec<(camel_api::FunctionId, Option<String>)> {
868        self.function_invoker
869            .as_ref()
870            .map(|invoker| invoker.function_refs_for_route(route_id))
871            .unwrap_or_default()
872    }
873
874    fn discard_function_staging(&self) {
875        if let Some(invoker) = &self.function_invoker {
876            invoker.discard_staging(0);
877        }
878    }
879
880    /// Returns the number of routes in the controller.
881    pub fn route_count(&self) -> usize {
882        self.routes.route_count()
883    }
884
885    pub fn in_flight_count(&self, route_id: &str) -> Option<u64> {
886        self.routes.in_flight_count(route_id)
887    }
888
889    /// Returns `true` if a route with the given ID exists.
890    pub fn route_exists(&self, route_id: &str) -> bool {
891        self.routes.route_exists(route_id)
892    }
893
894    /// Returns all route IDs.
895    pub fn route_ids(&self) -> Vec<String> {
896        self.routes.route_ids()
897    }
898
899    pub fn route_source_hash(&self, route_id: &str) -> Option<u64> {
900        self.routes.route_source_hash(route_id)
901    }
902
903    /// Returns route IDs that should auto-start, sorted by startup order (ascending).
904    pub fn auto_startup_route_ids(&self) -> Vec<String> {
905        self.routes.auto_startup_route_ids()
906    }
907
908    /// Returns route IDs sorted by shutdown order (startup order descending).
909    pub fn shutdown_route_ids(&self) -> Vec<String> {
910        self.routes.shutdown_route_ids()
911    }
912
913    /// Atomically swap the pipeline of a route (zero-downtime).
914    ///
915    /// In-flight requests finish with the old pipeline (kept alive by Arc).
916    /// New requests immediately use the new pipeline.
917    ///
918    /// ## Rejection policy
919    ///
920    /// Returns an error if the route has lifecycle-bearing steps or an active
921    /// aggregate — these require the **Restart path** (stop → swap → start).
922    ///
923    /// The caller (e.g. `reload_actions::apply_swap`) MUST catch this rejection
924    /// and fall back to:
925    /// 1. `stop_route_reload` — drain lifecycle, stop consumer
926    /// 2. `swap_pipeline_raw` — bypass the lifecycle check (route is stopped)
927    /// 3. `start_route_reload` — re-create consumer with the new pipeline
928    ///
929    /// This is the "reject, don't defer" policy (oracle Fix 3): the swap is
930    /// refused upfront rather than silently deferring or partially swapping.
931    pub fn swap_pipeline(
932        &self,
933        route_id: &str,
934        new_pipeline: BoxProcessor,
935    ) -> Result<(), CamelError> {
936        let managed = self
937            .routes
938            .get(route_id)
939            .ok_or_else(|| CamelError::RouteError(format!("Route '{}' not found", route_id)))?;
940
941        let assembly = managed.pipeline.load();
942        let has_lifecycle = !assembly.lifecycle.is_empty();
943
944        if has_lifecycle || managed.agg_service.is_some() {
945            warn!(
946                route_id = %route_id,
947                "Hot-swap rejected — route has lifecycle/agg steps; use Restart path"
948            );
949            return Err(CamelError::RouteError(format!(
950                "Route '{}' contains stateful steps (lifecycle-bearing). Hot-swap not supported — use restart.",
951                route_id
952            )));
953        }
954
955        drop(assembly);
956
957        if managed.aggregate_split.is_some() {
958            warn!(
959                route_id = %route_id,
960                "swap_pipeline: aggregate routes with timeout do not support hot-reload of pre/post segments"
961            );
962        }
963
964        super::pipeline_runtime::swap_pipeline_raw(&managed.pipeline, new_pipeline, vec![]);
965        debug!(route_id = %route_id, "Pipeline swapped atomically");
966        Ok(())
967    }
968
969    /// Non-checking raw pipeline swap — bypasses lifecycle/aggregate rejection.
970    ///
971    /// Only for use after the route has been stopped (Restart path).
972    /// Does NOT check for lifecycle handles or aggregate service — the caller
973    /// is responsible for ensuring the route is safe to swap.
974    ///
975    /// Accepts `lifecycle` so that the new pipeline assembly records the
976    /// lifecycle handles from the compiled steps.  When the route is
977    /// subsequently stopped, these handles are drained.
978    pub(crate) fn swap_pipeline_raw(
979        &self,
980        route_id: &str,
981        new_pipeline: BoxProcessor,
982        lifecycle: Vec<Arc<dyn StepLifecycle>>,
983    ) -> Result<(), CamelError> {
984        let managed = self
985            .routes
986            .get(route_id)
987            .ok_or_else(|| CamelError::RouteError(format!("Route '{}' not found", route_id)))?;
988        super::pipeline_runtime::swap_pipeline_raw(&managed.pipeline, new_pipeline, lifecycle);
989        debug!(route_id = %route_id, "Pipeline swapped (raw — lifecycle bypass)");
990        Ok(())
991    }
992
993    /// Returns the from_uri of a route, if it exists.
994    pub fn route_from_uri(&self, route_id: &str) -> Option<String> {
995        self.routes.route_from_uri(route_id)
996    }
997
998    /// Return all route_ids that consume from the given source endpoint URI.
999    pub fn routes_for_endpoint(&self, uri: &str) -> Vec<String> {
1000        self.endpoint_index.routes_for(uri)
1001    }
1002
1003    /// Return all registered source endpoint URIs.
1004    pub fn list_endpoint_uris(&self) -> Vec<String> {
1005        self.endpoint_index.list_uris()
1006    }
1007
1008    /// Get a clone of the current pipeline for a route.
1009    ///
1010    /// This is useful for testing and introspection.
1011    /// Returns `None` if the route doesn't exist.
1012    pub fn get_pipeline(&self, route_id: &str) -> Option<BoxProcessor> {
1013        self.routes.get_pipeline(route_id)
1014    }
1015
1016    /// Check whether the running route has lifecycle-bearing steps.
1017    ///
1018    /// Returns `false` when the route is missing.
1019    pub(crate) fn route_has_lifecycle(&self, route_id: &str) -> bool {
1020        self.routes
1021            .get(route_id)
1022            .map(|managed| !managed.pipeline.load().lifecycle.is_empty())
1023            .unwrap_or(false)
1024    }
1025
1026    /// Internal stop implementation that can set custom status.
1027    pub(super) async fn stop_route_internal(&mut self, route_id: &str) -> Result<(), CamelError> {
1028        self.routes.stop_route(route_id).await
1029    }
1030
1031    pub async fn start_route_reload(&mut self, route_id: &str) -> Result<(), CamelError> {
1032        self.start_route(route_id).await
1033    }
1034
1035    pub async fn stop_route_reload(&mut self, route_id: &str) -> Result<(), CamelError> {
1036        self.stop_route(route_id).await
1037    }
1038}
1039
1040// ── Aggregator route helpers ──
1041
1042impl DefaultRouteController {
1043    /// Start a route with an aggregate split (pre-pipeline → aggregator → post-pipeline).
1044    ///
1045    /// Spawns a biased-select forward loop that routes exchanges through the
1046    /// pre-pipeline, aggregator, and post-pipeline in sequence, with late-exchange
1047    /// handling and force-completion on stop.
1048    ///
1049    /// drainclaim claim propagation: the forward loop submits each
1050    /// envelope's [`InFlightClaim`] into the aggregator WITH its exchange
1051    /// ([`AggregatorService::submit_with_claim`]) — a pending stash parks
1052    /// the claim inside the bucket so the stashed exchange stays counted,
1053    /// and every emission path (sync completion, timeout `late_tx` fire,
1054    /// `force_complete_all` at stop) hands the bucket's claims back to
1055    /// this loop, which holds them across the post-pipeline
1056    /// continuation. Paths that drop a bucket without emitting (TTL
1057    /// eviction, unarmed-bucket release, discard-on-timeout, saturated
1058    /// late channel) release by dropping. Exactly one release per claim
1059    /// on every path. claimfamily (rc-e1a4f): every oneshot dispatch
1060    /// through the pre/post pipelines splits a sibling claim onto the
1061    /// exchange (in-band results take it back, stash emissions escape
1062    /// with theirs), so stash sites embedded in those pipelines stay
1063    /// counted.
1064    #[allow(clippy::too_many_arguments)]
1065    pub(super) async fn start_aggregate_route(
1066        &mut self,
1067        route_id: &str,
1068        split: AggregateSplitInfo,
1069        consumer: Box<dyn Consumer>,
1070        consumer_ctx: ConsumerContext,
1071        mut rx: mpsc::Receiver<ExchangeEnvelope>,
1072        crash_notifier: Option<mpsc::Sender<CrashNotification>>,
1073        runtime_for_consumer: Option<Weak<dyn RuntimeHandle>>,
1074        tx_for_storage: mpsc::Sender<ExchangeEnvelope>,
1075        // Pipeline cancellation — a child of the managed route's pipeline_cancel_token.
1076        pipeline_cancel: CancellationToken,
1077        drain_in_flight: Arc<std::sync::atomic::AtomicU64>,
1078    ) -> Result<(), CamelError> {
1079        // drainclaim claim propagation: the late channel carries each
1080        // emission's stashed claims alongside the aggregated exchange —
1081        // the loop holds them across the post-pipeline continuation.
1082        let (late_tx, late_rx) =
1083            mpsc::channel::<camel_processor::aggregator::AggregateEmission>(256);
1084
1085        let route_cancel_clone = pipeline_cancel.clone();
1086        let mut svc = AggregatorService::new(
1087            split.agg_config.clone(),
1088            late_tx,
1089            Arc::clone(&self.languages),
1090            route_cancel_clone,
1091        );
1092        // Queue-depth visibility: the TTL-sweep pass reports the buffered
1093        // group count as camel_queue_depth{queue="aggregator:<route>"}.
1094        if let Some(metrics) = self.tracer_metrics.clone() {
1095            svc = svc.with_queue_metrics(metrics, format!("aggregator:{route_id}"));
1096        }
1097        let agg = Arc::new(svc);
1098
1099        let pipeline_cancel_for_monitor = pipeline_cancel.clone();
1100        // rc-e2r9: capture route_id and metrics for b′ emission at reply-drop
1101        // sites inside the spawned task.
1102        let route_id_for_metrics = route_id.to_string();
1103        let metrics_for_reply_drop = self.tracer_metrics.clone();
1104        // rc-jxkj cohort gate: owned by the forward loop — the envelope arm
1105        // parks dispatch until the startup cohort opens the gate.
1106        let mut cohort_rx = self.cohort.subscribe();
1107        let agg_for_monitor = Arc::clone(&agg);
1108
1109        {
1110            let managed = self
1111                .routes
1112                .get_mut(route_id)
1113                .expect("invariant: route must exist"); // allow-unwrap
1114            managed.agg_service = Some(Arc::clone(&agg));
1115        }
1116
1117        let late_rx = Arc::new(tokio::sync::Mutex::new(late_rx));
1118        let pre_pipeline = split.pre_pipeline;
1119        let post_pipeline = split.post_pipeline;
1120
1121        // Spawn biased select forward loop
1122        let pipeline_handle = tokio::spawn(async move {
1123            loop {
1124                tokio::select! {
1125                    biased;
1126
1127                    // Ungated by design (D3, rc-jxkj): late exchanges exist
1128                    // only after a dispatched envelope traversed the
1129                    // aggregator — transitively post-activation. Gating here
1130                    // would be dead code and a self-deadlock risk.
1131                    late_ex = async {
1132                        let mut rx = late_rx.lock().await;
1133                        rx.recv().await
1134                    } => {
1135                        match late_ex {
1136                            Some(emission) => {
1137                                // drainclaim: hold the emission's stashed
1138                                // claims across the post-pipeline — drop
1139                                // at iteration end is the release.
1140                                let AggregateEmission {
1141                                    mut exchange,
1142                                    claims: _in_flight_claims,
1143                                } = emission;
1144                                // claimfamily (rc-e1a4f): split a sibling onto the emission exchange so
1145                                // residency inside stash sites embedded in the post-pipeline stays
1146                                // counted after this iteration. An in-band result drops at iteration end
1147                                // (its sibling releases with it); a stash emission escapes with its
1148                                // sibling; an Err drops it with the exchange.
1149                                exchange.in_flight_claim = _in_flight_claims
1150                                    .iter()
1151                                    .flatten()
1152                                    .next()
1153                                    .map(InFlightClaim::split);
1154                                let pipe = post_pipeline.load();
1155                                if let Err(e) =
1156                                    pipe.processor.clone_inner().oneshot(exchange).await
1157                                {
1158                                    tracing::warn!(error = %e, "late exchange post-pipeline failed");
1159                                }
1160                            }
1161                            None => return,
1162                        }
1163                    }
1164
1165                    envelope_opt = rx.recv() => {
1166                        match envelope_opt {
1167                            Some(envelope) => {
1168                                // rc-jxkj cohort gate: park dispatch until
1169                                // the startup cohort completes (same guard
1170                                // as the non-aggregate drain loops).
1171                                // rc-z5qz: `biased` with the gate polled
1172                                // FIRST — with the cohort open AND the
1173                                // pipeline token cancelled, both branches
1174                                // are ready and an unbiased select picks
1175                                // randomly, letting the cancel arm drop a
1176                                // deliverable envelope (force_complete_all
1177                                // then sees no buckets). Gate-open must win
1178                                // deterministically; a genuinely closed
1179                                // gate still drops on cancel (rc-jxkj
1180                                // semantics preserved).
1181                                tokio::select! {
1182                                    biased;
1183                                    _ = cohort_rx.wait_for(|open| *open) => {}
1184                                    _ = pipeline_cancel.cancelled() => {
1185                                        // Drop the envelope; reply_tx (if
1186                                        // any) resolves to ChannelClosed for
1187                                        // the send_and_wait waiter.
1188                                        // `continue`, not `return`: the token
1189                                        // is already cancelled, so the next
1190                                        // loop iteration lands in the biased
1191                                        // outer select's cancel arm below,
1192                                        // which runs the force_complete_all +
1193                                        // late_rx cleanup. A `return` would
1194                                        // skip that cleanup and silently
1195                                        // drop a pending bucket across a
1196                                        // stop→restart gate re-arm.
1197                                        continue;
1198                                    }
1199                                }
1200                                let ExchangeEnvelope {
1201                                    exchange,
1202                                    reply_tx,
1203                                    in_flight_claim,
1204                                } = envelope;
1205                                let _drain_guard = super::route_helpers::DrainGuard::new(Arc::clone(&drain_in_flight));
1206                                // drainclaim claim PROPAGATION: the
1207                                // envelope's claim travels WITH the
1208                                // exchange into the aggregator — a
1209                                // pending stash parks it in the bucket
1210                                // (counted until the bucket completes),
1211                                // and a sync completion returns it (with
1212                                // the rest of the bucket's claims) to be
1213                                // held across the post-pipeline below.
1214                                // Rejection paths drop it inside the
1215                                // service — rejected = released.
1216                                let pre_pipe = pre_pipeline.load();
1217                                // claimfamily (rc-e1a4f): split a sibling claim onto the exchange so
1218                                // residency inside stash sites embedded in the pre-pipeline stays
1219                                // counted after this oneshot resolves. Taken back from an in-band Ok
1220                                // result below — the envelope's own claim still travels with the
1221                                // exchange into the aggregator; a stash emission escapes with its
1222                                // sibling; an Err drops it with the exchange.
1223                                let mut exchange = exchange;
1224                                exchange.in_flight_claim = in_flight_claim.as_ref().map(InFlightClaim::split);
1225                                let ex = match pre_pipe.processor.clone_inner().oneshot(exchange).await {
1226                                    // claimfamily: in-band completion — reclaim the sibling so release
1227                                    // stays at this loop iteration.
1228                                    Ok(mut ex) => {
1229                                        ex.in_flight_claim = None;
1230                                        ex
1231                                    }
1232                                    Err(e) => {
1233                                        // rc-e2r9: the real error rides with the
1234                                        // result so a dropped receiver still gets
1235                                        // it (ConsumerStopping suppressed).
1236                                        send_reply_or_b_prime(
1237                                            reply_tx,
1238                                            Err(e),
1239                                            &metrics_for_reply_drop,
1240                                            &route_id_for_metrics,
1241                                            "aggregate:pre-pipeline",
1242                                        );
1243                                        continue;
1244                                    }
1245                                };
1246
1247                                let AggregationReceipt { reply, claims } =
1248                                    agg.submit_with_claim(ex, in_flight_claim).await;
1249                                // drainclaim: hold the completed bucket's
1250                                // claims across the post-pipeline
1251                                // continuation; dropped at iteration end or
1252                                // any `continue`/`return` (scope exit).
1253                                let _completed_bucket_claims = claims;
1254
1255                                match reply {
1256                                    Ok(ex) => {
1257                                        if !is_pending(&ex) {
1258                                            let post_pipe = post_pipeline.load();
1259                                            // claimfamily (rc-e1a4f): split a sibling from the completed
1260                                            // bucket's claims onto the aggregated output so residency inside
1261                                            // stash sites embedded in the post-pipeline stays counted after
1262                                            // this iteration. Taken back from an in-band result before the
1263                                            // reply — release stays at iteration end; a stash emission
1264                                            // escapes with its sibling; an Err drops it with the exchange.
1265                                            let mut ex = ex;
1266                                            ex.in_flight_claim = _completed_bucket_claims
1267                                                .iter()
1268                                                .flatten()
1269                                                .next()
1270                                                .map(InFlightClaim::split);
1271                                            let mut out = post_pipe.processor.clone_inner().oneshot(ex).await;
1272                                            if let Ok(ref mut out_ex) = out {
1273                                                out_ex.in_flight_claim = None;
1274                                            }
1275                                            // rc-e2r9 review, Important 1: this site
1276                                            // previously `let _ = send`ed the
1277                                            // post-pipeline result — an Err with an
1278                                            // abandoned receiver vanished silently.
1279                                            // The helper closes that gap.
1280                                            send_reply_or_b_prime(
1281                                                reply_tx,
1282                                                out,
1283                                                &metrics_for_reply_drop,
1284                                                &route_id_for_metrics,
1285                                                "aggregate:post-pipeline",
1286                                            );
1287                                        } else {
1288                                            // Pending Ok: a dropped reply is silent —
1289                                            // b′ is an ERROR signal.
1290                                            send_reply_or_b_prime(
1291                                                reply_tx,
1292                                                Ok(ex),
1293                                                &metrics_for_reply_drop,
1294                                                &route_id_for_metrics,
1295                                                "aggregate:pending",
1296                                            );
1297                                        }
1298                                    }
1299                                    Err(e) => {
1300                                        send_reply_or_b_prime(
1301                                            reply_tx,
1302                                            Err(e),
1303                                            &metrics_for_reply_drop,
1304                                            &route_id_for_metrics,
1305                                            "aggregate:pipeline",
1306                                        );
1307                                    }
1308                                }
1309                            }
1310                            None => return,
1311                        }
1312                    }
1313
1314                    _ = pipeline_cancel.cancelled() => {
1315                        agg.force_complete_all();
1316                        let mut rx_guard = late_rx.lock().await;
1317                        while let Ok(late_ex) = rx_guard.try_recv() {
1318                            // drainclaim: each forced emission's claims are
1319                            // held across its post-pipeline continuation —
1320                            // dropped after the oneshot completes (or
1321                            // immediately, which is also a release).
1322                            let AggregateEmission {
1323                                mut exchange,
1324                                claims: _in_flight_claims,
1325                            } = late_ex;
1326                            // claimfamily (rc-e1a4f): split a sibling onto the forced emission so
1327                            // residency inside stash sites embedded in the post-pipeline stays
1328                            // counted. The discarded result releases an in-band sibling with its
1329                            // drop; a stash emission escapes with its sibling.
1330                            exchange.in_flight_claim = _in_flight_claims
1331                                .iter()
1332                                .flatten()
1333                                .next()
1334                                .map(InFlightClaim::split);
1335                            let pipe = post_pipeline.load();
1336                            let _ = pipe.processor.clone_inner().oneshot(exchange).await;
1337                        }
1338                        break;
1339                    }
1340                }
1341            }
1342        });
1343        #[cfg(test)]
1344        emit_start_route_event("pipeline_spawned", route_id);
1345
1346        // Start consumer after pipeline loop is spawned to avoid startup races
1347        // where consumers emit exchanges before the route pipeline begins polling.
1348        // rc-kh7c cleanup parity: capture the consumer's cancel token before
1349        // consumer_ctx moves into the task so the failure arm can stop child
1350        // tasks spawned by consumer.start().
1351        let consumer_cancel_for_cleanup = consumer_ctx.cancel_token();
1352        let (consumer_handle, startup_rx, watcher_inputs, outer_inputs) =
1353            super::consumer_management::spawn_consumer_task(
1354                route_id.to_string(),
1355                consumer,
1356                consumer_ctx,
1357                crash_notifier,
1358                runtime_for_consumer,
1359                false,
1360            );
1361
1362        // rc-w1u9: await consumer startup handshake for aggregate routes too
1363        // so bind failures surface as route-start errors.
1364        // For Immediate consumers the receiver is pre-resolved (rc-slvd).
1365        let startup_result =
1366            super::consumer_management::await_consumer_startup(startup_rx, "startup").await;
1367        // rc-kh7c: on failure, abort the orphaned consumer task and cancel the
1368        // pipeline so neither runs detached. The aggregate pipeline loop would
1369        // eventually self-clean via rx-drop + late_tx-drop, but cancelling
1370        // pipeline_cancel also triggers force_complete_all (aggregate cleanup).
1371        if let Err(e) = startup_result {
1372            consumer_handle.abort();
1373            pipeline_cancel_for_monitor.cancel();
1374            // Deliberate Explicit-failure cleanup parity with the trait start arm (rc-kh7c).
1375            consumer_cancel_for_cleanup.cancel();
1376            return Err(e);
1377        }
1378
1379        // Detached failure watcher for Immediate consumers (rc-slvd).
1380        if let Some(inputs) = watcher_inputs {
1381            super::consumer_management::spawn_failure_watcher(inputs);
1382        }
1383
1384        // Detached outer-task watcher for Explicit consumers (rc-a7rh):
1385        // spawned only after the handshake resolved Ok — rollback
1386        // terminations (abort-then-cancel above) happen before this point
1387        // and are never watched. Explicit consumers on aggregate routes
1388        // get identical coverage — no aggregate carve-out.
1389        if let Some(outer) = outer_inputs {
1390            super::consumer_management::spawn_outer_task_watcher(outer);
1391        }
1392
1393        // Extend the stored consumer handle through aggregate force-completion.
1394        // While this monitor drains pending buckets, handle_is_running still reports
1395        // the Route as running because forced exchanges may still be in post-pipeline.
1396        //
1397        // bd rc-iioeq: a natural consumer exit (e.g. timer repeatCount
1398        // exhausted) must NOT destroy buckets whose inactivity timeout is
1399        // armed. With `force_completion_on_stop=false` (the default),
1400        // `force_complete_all` CANCELS the armed timeout task and silently
1401        // discards the bucket, so the inactivity emission never happens.
1402        // The forward loop stays up (the stored channel sender keeps the
1403        // input channel open), so an armed bucket still emits downstream
1404        // when its timeout fires. Buckets with no armed timeout task —
1405        // size/predicate-only, or timeout-configured but over the
1406        // `max_timeout_tasks` cap — can never complete after the consumer
1407        // exits; `release_unarmed_buckets` discards them eagerly so they
1408        // are not orphaned (the bucket_ttl sweep only runs inside the
1409        // pipeline's next exchange, which never arrives).
1410        let force_on_stop = agg_for_monitor.config().force_completion_on_stop;
1411        let consumer_handle = tokio::spawn(async move {
1412            let _ = consumer_handle.await;
1413            if !pipeline_cancel_for_monitor.is_cancelled() {
1414                if force_on_stop {
1415                    agg_for_monitor.force_complete_all();
1416                    pipeline_cancel_for_monitor.cancel();
1417                } else {
1418                    agg_for_monitor.release_unarmed_buckets();
1419                }
1420            }
1421        });
1422        #[cfg(test)]
1423        emit_start_route_event("consumer_spawned", route_id);
1424
1425        {
1426            let managed = self
1427                .routes
1428                .get_mut(route_id)
1429                .expect("invariant: route must exist"); // allow-unwrap
1430            managed.consumer_handle = Some(consumer_handle);
1431            managed.pipeline_handle = Some(pipeline_handle);
1432            managed.channel_sender = Some(tx_for_storage);
1433        }
1434
1435        info!(route_id = %route_id, "Route started (aggregate with timeout)");
1436        Ok(())
1437    }
1438
1439    /// Test-only: inject lifecycle handles into an existing route's pipeline
1440    /// assembly.  This makes the route lifecycle-bearing so that swap_pipeline
1441    /// rejects it, forcing callers (like reload_actions::apply_swap) to take
1442    /// the Restart path instead.
1443    #[cfg(test)]
1444    pub(crate) fn set_route_lifecycle_for_test(
1445        &mut self,
1446        route_id: &str,
1447        lifecycle: Vec<Arc<dyn StepLifecycle>>,
1448    ) -> Result<(), CamelError> {
1449        use super::pipeline_runtime::PipelineAssembly;
1450        use camel_api::SyncBoxProcessor;
1451        use std::sync::Arc;
1452
1453        let managed = self
1454            .routes
1455            .get_mut(route_id)
1456            .ok_or_else(|| CamelError::RouteError(format!("Route '{}' not found", route_id)))?;
1457        let old_processor = managed.pipeline.load().processor.clone_inner();
1458        managed.pipeline.store(Arc::new(PipelineAssembly::new(
1459            SyncBoxProcessor::new(old_processor),
1460            lifecycle,
1461        )));
1462        Ok(())
1463    }
1464}
1465
1466// ── rc-e2r9: b′ signal emission on reply-drop ──
1467// The shared `send_reply_or_b_prime` helper lives in
1468// `super::route_helpers` — it is used by both this module (aggregate drain
1469// loop) and `route_controller_trait` (Concurrent/Sequential pipelines).
1470
1471#[cfg(test)]
1472impl crate::hot_reload::ports::ReloadIntrospectionPort for DefaultRouteController {
1473    fn route_ids(&self) -> Vec<String> {
1474        self.route_ids() // inherent pub fn
1475    }
1476    fn route_from_uri(&self, route_id: &str) -> Option<String> {
1477        self.route_from_uri(route_id) // inherent pub fn
1478    }
1479    fn route_source_hash(&self, route_id: &str) -> Option<u64> {
1480        self.route_source_hash(route_id) // inherent pub fn
1481    }
1482}
1483
1484#[cfg(test)]
1485#[path = "route_controller_tests.rs"]
1486mod tests;
1487
1488#[cfg(test)]
1489#[path = "cohort_activation_regression.rs"]
1490mod cohort_activation_regression;