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;