camel_core/lifecycle/adapters/
route_controller_trait.rs1use std::sync::Arc;
7use std::time::Duration;
8
9use tokio::sync::mpsc;
10use tokio_util::sync::CancellationToken;
11use tower::Service;
12use tracing::{error, info, warn};
13
14use camel_api::{CamelError, NoOpMetrics, StepLifecycle, StepShutdownReason};
15use camel_component_api::{ConcurrencyModel, ConsumerContext, consumer::ExchangeEnvelope};
16
17use crate::lifecycle::adapters::consumer_management;
18use crate::lifecycle::adapters::controller_component_context::ControllerComponentContext;
19use crate::lifecycle::adapters::route_compiler::CANCEL_TOKEN;
20use crate::lifecycle::adapters::route_controller::DefaultRouteController;
21#[cfg(test)]
22use crate::lifecycle::adapters::route_helpers::emit_start_route_event;
23use crate::lifecycle::adapters::route_helpers::{
24 DrainGuard, handle_is_running, inferred_lifecycle_label, ready_with_backoff,
25};
26use crate::lifecycle::adapters::route_registry::DEFAULT_SHUTDOWN_TIMEOUT;
27
28async fn rollback_started(route_id: &str, handles: &[Arc<dyn StepLifecycle>]) {
39 for handle in handles.iter().rev() {
40 if let Err(e) = handle.shutdown(StepShutdownReason::RouteStop).await {
41 warn!(
42 route_id = %route_id,
43 step = handle.name(),
44 error = %e,
45 "best-effort step shutdown during start rollback failed"
46 );
47 }
48 }
49}
50
51#[async_trait::async_trait]
52impl camel_api::RouteController for DefaultRouteController {
53 async fn start_route(&mut self, route_id: &str) -> Result<(), CamelError> {
54 {
56 let managed = self
57 .routes
58 .get_mut(route_id)
59 .ok_or_else(|| CamelError::RouteError(format!("Route '{}' not found", route_id)))?;
60
61 let consumer_running = handle_is_running(&managed.consumer_handle);
62 let pipeline_running = handle_is_running(&managed.pipeline_handle);
63 if consumer_running && pipeline_running {
64 return Ok(());
65 }
66 if !consumer_running && pipeline_running {
67 return Err(CamelError::RouteError(format!(
68 "Route '{}' is suspended; use resume_route() to resume, or stop_route() then start_route() for full restart",
69 route_id
70 )));
71 }
72 if consumer_running && !pipeline_running {
73 return Err(CamelError::RouteError(format!(
74 "Route '{}' has inconsistent execution state; stop_route() then retry start_route()",
75 route_id
76 )));
77 }
78 }
79
80 info!(route_id = %route_id, "Starting route");
81
82 let (from_uri, pipeline, concurrency) = {
84 let managed = self
85 .routes
86 .get(route_id)
87 .expect("invariant: route must exist after prior existence check"); (
89 managed.from_uri.clone(),
90 Arc::clone(&managed.pipeline),
91 managed.concurrency.clone(),
92 )
93 };
94
95 let lifecycle_handles: Vec<Arc<dyn StepLifecycle>> = pipeline.load().lifecycle.clone();
101 for (idx, handle) in lifecycle_handles.iter().enumerate() {
102 if let Err(start_err) = handle.start().await {
103 warn!(
104 route_id = %route_id,
105 step = handle.name(),
106 "step start failed; rolling back already-started steps"
107 );
108 rollback_started(route_id, &lifecycle_handles[0..idx]).await;
110 return Err(start_err);
111 }
112 }
113
114 let crash_notifier = self.crash_notifier.clone();
116 let runtime_for_consumer = self.runtime.clone();
117
118 let consumer_component_ctx = Arc::new(ControllerComponentContext::new(
119 Arc::clone(&self.registry),
120 Arc::clone(&self.languages),
121 self.tracer_metrics
122 .clone()
123 .unwrap_or_else(|| Arc::new(NoOpMetrics)),
124 Arc::clone(&self.platform_service),
125 self.health_registry(),
126 Some(route_id.to_string()),
127 ));
128 let consumer_rt: Arc<dyn camel_component_api::RuntimeObservability> =
129 Arc::clone(&consumer_component_ctx) as Arc<_>;
130 let (mut consumer, consumer_concurrency) = match consumer_management::create_route_consumer(
131 consumer_rt,
132 &self.registry,
133 &from_uri,
134 consumer_component_ctx.as_ref(),
135 ) {
136 Ok(v) => v,
137 Err(e) => {
141 rollback_started(route_id, &lifecycle_handles).await;
142 return Err(e);
143 }
144 };
145
146 let effective_concurrency = concurrency.unwrap_or(consumer_concurrency);
148
149 let managed = self
151 .routes
152 .get_mut(route_id)
153 .expect("invariant: route must exist after prior existence check"); if let (Some(sp_config), Some(authenticator)) = (
157 managed.compiled.security_policy.as_ref(),
158 managed.compiled.security_authenticator.as_ref(),
159 ) {
160 use camel_component_api::SecurityContext;
161 let sec_ctx =
162 SecurityContext::from_arc(Arc::clone(&sp_config.policy), Arc::clone(authenticator))
163 .with_credential_sources(sp_config.credential_sources.clone());
164 consumer.set_security_context(sec_ctx);
165 }
166
167 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(256);
169 let consumer_cancel = managed.consumer_cancel_token.child_token();
171 let pipeline_cancel = managed.pipeline_cancel_token.child_token();
172 let drain_in_flight = Arc::clone(&managed.drain_in_flight);
173 let tx_for_storage = tx.clone();
175 let consumer_ctx = ConsumerContext::new(tx, consumer_cancel.clone(), route_id.to_string());
176
177 let split_clone = managed.aggregate_split.clone();
179 if let Some(split) = split_clone {
180 let result = self
181 .start_aggregate_route(
182 route_id,
183 split,
184 consumer,
185 consumer_ctx,
186 rx,
187 crash_notifier,
188 runtime_for_consumer,
189 tx_for_storage,
190 pipeline_cancel,
191 drain_in_flight,
192 )
193 .await;
194 if result.is_err() {
197 if let Some(managed) = self.routes.get_mut(route_id) {
200 managed.consumer_cancel_token.cancel();
201 }
202 rollback_started(route_id, &lifecycle_handles).await;
203 }
204 return result;
205 }
206 let pipeline_cancel_for_cleanup = pipeline_cancel.clone();
212
213 let pipeline_handle = match effective_concurrency {
215 ConcurrencyModel::Concurrent { max } => {
216 let sem = max.map(|n| Arc::new(tokio::sync::Semaphore::new(n)));
217 tokio::spawn(async move {
218 loop {
219 let permit = match &sem {
222 Some(s) => {
223 let acquired = tokio::select! {
224 p = Arc::clone(s).acquire_owned() => p.expect("semaphore closed"), _ = pipeline_cancel.cancelled() => return,
226 };
227 Some(acquired)
228 }
229 None => None,
230 };
231
232 let envelope = tokio::select! {
233 envelope = rx.recv() => match envelope {
234 Some(e) => e,
235 None => return,
236 },
237 _ = pipeline_cancel.cancelled() => return,
238 };
239 let ExchangeEnvelope { exchange, reply_tx } = envelope;
240 let pipe_ref = Arc::clone(&pipeline);
241 let cancel = pipeline_cancel.clone();
242 let drain_clone = Arc::clone(&drain_in_flight);
243 tokio::spawn(async move {
244 let _permit = permit;
246 let _drain_guard = DrainGuard::new(drain_clone);
247
248 let mut pipe = pipe_ref.load().processor.clone_inner();
250
251 if let Err(e) = ready_with_backoff(&mut pipe, &cancel).await {
253 if let Some(tx) = reply_tx {
254 let _ = tx.send(Err(e));
255 }
256 return;
257 }
258
259 let result = CANCEL_TOKEN
262 .scope(cancel, async move { pipe.call(exchange).await })
263 .await;
264 if let Some(tx) = reply_tx {
265 let _ = tx.send(result);
266 } else if let Err(ref e) = result {
267 error!("Pipeline error: {e}");
269 }
270 });
271 }
272 })
273 }
274 _ => {
280 tokio::spawn(async move {
281 loop {
282 let envelope = tokio::select! {
284 envelope = rx.recv() => match envelope {
285 Some(e) => e,
286 None => return, },
288 _ = pipeline_cancel.cancelled() => {
289 return;
291 }
292 };
293 let ExchangeEnvelope { exchange, reply_tx } = envelope;
294
295 let mut pipeline = pipeline.load().processor.clone_inner();
297
298 if let Err(e) = ready_with_backoff(&mut pipeline, &pipeline_cancel).await {
299 if let Some(tx) = reply_tx {
300 let _ = tx.send(Err(e));
301 }
302 return;
303 }
304
305 let cancel = pipeline_cancel.clone();
311 let _drain_guard = DrainGuard::new(Arc::clone(&drain_in_flight));
312 let result = CANCEL_TOKEN
313 .scope(cancel, async move { pipeline.call(exchange).await })
314 .await;
315 if let Some(tx) = reply_tx {
316 let _ = tx.send(result);
317 } else if let Err(ref e) = result {
318 error!("Pipeline error: {e}");
320 }
321 }
322 })
323 }
324 };
325 #[cfg(test)]
326 emit_start_route_event("pipeline_spawned");
327
328 let (consumer_handle, startup_rx) = consumer_management::spawn_consumer_task(
331 route_id.to_string(),
332 consumer,
333 consumer_ctx,
334 crash_notifier,
335 runtime_for_consumer,
336 false,
337 );
338 #[cfg(test)]
339 emit_start_route_event("consumer_spawned");
340
341 match consumer_management::await_consumer_startup(startup_rx, "startup").await {
346 Ok(()) => {}
347 Err(e) => {
348 consumer_handle.abort();
355 pipeline_cancel_for_cleanup.cancel();
356 consumer_cancel.cancel();
359 rollback_started(route_id, &lifecycle_handles).await;
360 return Err(e);
361 }
362 }
363
364 let managed = self
366 .routes
367 .get_mut(route_id)
368 .expect("invariant: route must exist after prior existence check"); managed.consumer_handle = Some(consumer_handle);
370 managed.pipeline_handle = Some(pipeline_handle);
371 managed.channel_sender = Some(tx_for_storage);
372
373 info!(route_id = %route_id, "Route started");
374 self.health_registry().mark_route_started(route_id);
375 Ok(())
376 }
377
378 async fn stop_route(&mut self, route_id: &str) -> Result<(), CamelError> {
379 self.stop_route_internal(route_id).await?;
380 self.health_registry().mark_route_stopped(route_id);
381 Ok(())
382 }
383
384 async fn restart_route(&mut self, route_id: &str) -> Result<(), CamelError> {
385 self.stop_route(route_id).await?;
386 tokio::time::sleep(Duration::from_millis(100)).await;
387 self.start_route(route_id).await
388 }
389
390 async fn suspend_route(&mut self, route_id: &str) -> Result<(), CamelError> {
391 let managed = self
393 .routes
394 .get_mut(route_id)
395 .ok_or_else(|| CamelError::RouteError(format!("Route '{}' not found", route_id)))?;
396
397 let consumer_running = handle_is_running(&managed.consumer_handle);
398 let pipeline_running = handle_is_running(&managed.pipeline_handle);
399
400 if !consumer_running || !pipeline_running {
402 return Err(CamelError::RouteError(format!(
403 "Cannot suspend route '{}' with execution lifecycle {}",
404 route_id,
405 inferred_lifecycle_label(managed)
406 )));
407 }
408
409 info!(route_id = %route_id, "Suspending route (consumer only, keeping pipeline)");
410
411 let managed = self
413 .routes
414 .get_mut(route_id)
415 .expect("invariant: route must exist after prior existence check"); managed.consumer_cancel_token.cancel();
417
418 let managed = self
420 .routes
421 .get_mut(route_id)
422 .expect("invariant: route must exist after prior existence check"); let consumer_handle = managed.consumer_handle.take();
424
425 let timeout_result = tokio::time::timeout(DEFAULT_SHUTDOWN_TIMEOUT, async {
427 if let Some(handle) = consumer_handle {
428 let _ = handle.await;
429 }
430 })
431 .await;
432
433 if timeout_result.is_err() {
434 warn!(route_id = %route_id, "Consumer shutdown timed out during suspend");
435 }
436
437 let managed = self
439 .routes
440 .get_mut(route_id)
441 .expect("invariant: route must exist after prior existence check"); managed.consumer_cancel_token = CancellationToken::new();
445
446 info!(route_id = %route_id, "Route suspended (pipeline still running)");
447 self.health_registry().mark_route_stopped(route_id);
448 Ok(())
449 }
450
451 async fn resume_route(&mut self, route_id: &str) -> Result<(), CamelError> {
452 let managed = self
454 .routes
455 .get(route_id)
456 .ok_or_else(|| CamelError::RouteError(format!("Route '{}' not found", route_id)))?;
457
458 let consumer_running = handle_is_running(&managed.consumer_handle);
459 let pipeline_running = handle_is_running(&managed.pipeline_handle);
460 if consumer_running || !pipeline_running {
461 return Err(CamelError::RouteError(format!(
462 "Cannot resume route '{}' with execution lifecycle {} (expected Suspended)",
463 route_id,
464 inferred_lifecycle_label(managed)
465 )));
466 }
467
468 let sender = managed.channel_sender.clone().ok_or_else(|| {
470 CamelError::RouteError("Suspended route has no channel sender".into())
471 })?;
472
473 let from_uri = managed.from_uri.clone();
475
476 info!(route_id = %route_id, "Resuming route (spawning consumer only)");
477
478 let consumer_component_ctx = Arc::new(ControllerComponentContext::new(
479 Arc::clone(&self.registry),
480 Arc::clone(&self.languages),
481 self.tracer_metrics
482 .clone()
483 .unwrap_or_else(|| Arc::new(NoOpMetrics)),
484 Arc::clone(&self.platform_service),
485 self.health_registry(),
486 Some(route_id.to_string()),
487 ));
488 let consumer_rt: Arc<dyn camel_component_api::RuntimeObservability> =
489 Arc::clone(&consumer_component_ctx) as Arc<_>;
490 let (mut consumer, _) = consumer_management::create_route_consumer(
491 consumer_rt,
492 &self.registry,
493 &from_uri,
494 consumer_component_ctx.as_ref(),
495 )?;
496
497 let managed = self
499 .routes
500 .get(route_id)
501 .expect("invariant: route must exist after prior existence check"); if let (Some(sp_config), Some(authenticator)) = (
503 managed.compiled.security_policy.as_ref(),
504 managed.compiled.security_authenticator.as_ref(),
505 ) {
506 use camel_component_api::SecurityContext;
507 let sec_ctx =
508 SecurityContext::from_arc(Arc::clone(&sp_config.policy), Arc::clone(authenticator))
509 .with_credential_sources(sp_config.credential_sources.clone());
510 consumer.set_security_context(sec_ctx);
511 }
512
513 let managed = self
515 .routes
516 .get_mut(route_id)
517 .expect("invariant: route must exist after prior existence check"); let consumer_cancel = managed.consumer_cancel_token.child_token();
521
522 let crash_notifier = self.crash_notifier.clone();
523 let runtime_for_consumer = self.runtime.clone();
524
525 let consumer_ctx =
527 ConsumerContext::new(sender, consumer_cancel.clone(), route_id.to_string());
528
529 let (consumer_handle, startup_rx) = consumer_management::spawn_consumer_task(
531 route_id.to_string(),
532 consumer,
533 consumer_ctx,
534 crash_notifier,
535 runtime_for_consumer,
536 true,
537 );
538
539 consumer_management::await_consumer_startup(startup_rx, "resume").await?;
542
543 let managed = self
545 .routes
546 .get_mut(route_id)
547 .expect("invariant: route must exist after prior existence check"); managed.consumer_handle = Some(consumer_handle);
549
550 info!(route_id = %route_id, "Route resumed");
551 self.health_registry().mark_route_started(route_id);
552 Ok(())
553 }
554
555 async fn start_all_routes(&mut self) -> Result<(), CamelError> {
556 let route_ids: Vec<String> = {
559 let pairs = self.routes.auto_startup_sorted();
560 pairs.into_iter().map(|(id, _)| id).collect()
561 };
562
563 info!("Starting {} auto-startup routes", route_ids.len());
564
565 let mut errors: Vec<String> = Vec::new();
567 for route_id in route_ids {
568 if let Err(e) = self.start_route(&route_id).await {
569 errors.push(format!("Route '{}': {}", route_id, e));
570 }
571 }
572
573 if !errors.is_empty() {
574 return Err(CamelError::RouteError(format!(
575 "Failed to start routes: {}",
576 errors.join(", ")
577 )));
578 }
579
580 info!("All auto-startup routes started");
581 Ok(())
582 }
583
584 async fn stop_all_routes(&mut self) -> Result<(), CamelError> {
585 let route_ids: Vec<String> = {
587 let pairs = self.routes.shutdown_sorted();
588 pairs.into_iter().map(|(id, _)| id).collect()
589 };
590
591 info!("Stopping {} routes", route_ids.len());
592
593 for route_id in route_ids {
594 let _ = self.stop_route(&route_id).await;
595 }
596
597 info!("All routes stopped");
598 Ok(())
599 }
600}