1use std::sync::Arc;
2
3use async_trait::async_trait;
4use tokio::sync::OnceCell;
5
6use camel_api::{
7 CamelError, MetricsCollector, RuntimeCommand, RuntimeCommandBus, RuntimeCommandResult,
8 RuntimeQuery, RuntimeQueryBus, RuntimeQueryResult,
9};
10
11use crate::lifecycle::application::commands::{
12 CommandDeps, execute_command, handle_register_internal,
13};
14use crate::lifecycle::application::ports::RouteRegistrationPort;
15use crate::lifecycle::application::ports::{
16 CommandDedupPort, EventPublisherPort, InFlightCountResult, ProjectionStorePort,
17 RouteRepositoryPort, RuntimeExecutionPort, RuntimeUnitOfWorkPort,
18};
19use crate::lifecycle::application::queries::{QueryDeps, execute_query};
20use crate::lifecycle::application::route_definition::RouteDefinition;
21use crate::lifecycle::domain::DomainError;
22use camel_component_api::HealthCheckRegistry as HealthCheckRegistryTrait;
23
24impl From<InFlightCountResult> for RuntimeQueryResult {
25 fn from(r: InFlightCountResult) -> Self {
26 match r {
27 InFlightCountResult::InFlightCount { route_id, count } => {
28 RuntimeQueryResult::InFlightCount { route_id, count }
29 }
30 InFlightCountResult::RouteNotFound { route_id } => {
31 RuntimeQueryResult::RouteNotFound { route_id }
32 }
33 }
34 }
35}
36
37pub struct RuntimeBus {
38 repo: Arc<dyn RouteRepositoryPort>,
39 projections: Arc<dyn ProjectionStorePort>,
40 events: Arc<dyn EventPublisherPort>,
41 dedup: Arc<dyn CommandDedupPort>,
42 uow: Option<Arc<dyn RuntimeUnitOfWorkPort>>,
43 execution: Option<Arc<dyn RuntimeExecutionPort>>,
44 health_registry: Option<Arc<dyn HealthCheckRegistryTrait>>,
45 metrics: Option<Arc<dyn MetricsCollector>>,
46 journal_recovered_once: OnceCell<u64>,
47}
48
49impl RuntimeBus {
50 pub fn new(
51 repo: Arc<dyn RouteRepositoryPort>,
52 projections: Arc<dyn ProjectionStorePort>,
53 events: Arc<dyn EventPublisherPort>,
54 dedup: Arc<dyn CommandDedupPort>,
55 ) -> Self {
56 Self {
57 repo,
58 projections,
59 events,
60 dedup,
61 uow: None,
62 execution: None,
63 health_registry: None,
64 metrics: None,
65 journal_recovered_once: OnceCell::new(),
66 }
67 }
68
69 pub fn with_uow(mut self, uow: Arc<dyn RuntimeUnitOfWorkPort>) -> Self {
70 self.uow = Some(uow);
71 self
72 }
73
74 pub fn with_execution(mut self, execution: Arc<dyn RuntimeExecutionPort>) -> Self {
75 self.execution = Some(execution);
76 self
77 }
78
79 pub fn with_health_registry(
80 mut self,
81 health_registry: Arc<crate::health_registry::HealthCheckRegistry>,
82 ) -> Self {
83 self.health_registry = Some(health_registry);
84 self
85 }
86
87 pub fn with_metrics(mut self, metrics: Arc<dyn MetricsCollector>) -> Self {
91 self.metrics = Some(metrics);
92 self
93 }
94
95 pub fn repo(&self) -> &Arc<dyn RouteRepositoryPort> {
96 &self.repo
97 }
98
99 pub async fn reconcile_transient_states(&self) -> Result<(), CamelError> {
103 self.ensure_journal_recovered().await?;
104 let deps = self.deps();
105 crate::lifecycle::application::commands::reconcile_transient_states(&deps).await
106 }
107
108 pub(crate) async fn register_aggregate_only(&self, route_id: String) -> Result<(), CamelError> {
109 self.ensure_journal_recovered().await?;
110 let deps = self.deps();
111 if deps.repo.load(&route_id).await?.is_some() {
112 return Err(CamelError::RouteError(format!(
113 "route '{route_id}' already registered"
114 )));
115 }
116 let (aggregate, events) =
117 crate::lifecycle::domain::RouteRuntimeAggregate::register(route_id.clone());
118 if let Some(uow) = &deps.uow {
119 uow.persist_upsert(
120 aggregate.clone(),
121 None,
122 crate::lifecycle::application::commands::project_from_aggregate(&aggregate),
123 &events,
124 )
125 .await?;
126 } else {
127 deps.repo.save(aggregate.clone()).await?;
128 if let Some(primary_error) =
129 crate::lifecycle::application::commands::upsert_projection_with_reconciliation(
130 &*deps.projections,
131 crate::lifecycle::application::commands::project_from_aggregate(&aggregate),
132 )
133 .await?
134 {
135 deps.events.publish(&events).await?;
136 return Err(CamelError::RouteError(format!(
137 "post-effect reconciliation recovered after runtime persistence error: {primary_error}"
138 )));
139 }
140 deps.events.publish(&events).await?;
141 }
142 Ok(())
143 }
144
145 fn deps(&self) -> CommandDeps {
146 CommandDeps {
147 repo: Arc::clone(&self.repo),
148 projections: Arc::clone(&self.projections),
149 events: Arc::clone(&self.events),
150 uow: self.uow.clone(),
151 execution: self.execution.clone(),
152 health_registry: self.health_registry.clone(),
153 }
154 }
155
156 fn query_deps(&self) -> QueryDeps {
157 QueryDeps {
158 projections: Arc::clone(&self.projections),
159 }
160 }
161
162 async fn ensure_journal_recovered(&self) -> Result<(), CamelError> {
163 let Some(uow) = &self.uow else {
164 return Ok(());
165 };
166
167 self.journal_recovered_once
168 .get_or_try_init(|| async {
169 uow.recover_from_journal().await?;
170 let nonce = uow.recovered_boot_nonce().await?;
171 Ok::<u64, CamelError>(nonce)
172 })
173 .await?;
174 Ok(())
175 }
176
177 pub(crate) fn boot_nonce(&self) -> u64 {
181 self.journal_recovered_once.get().copied().unwrap_or(0)
182 }
183}
184
185#[async_trait]
186impl RuntimeCommandBus for RuntimeBus {
187 async fn execute(&self, cmd: RuntimeCommand) -> Result<RuntimeCommandResult, CamelError> {
188 if let RuntimeCommand::ReloadTlsCerts {
192 scheme, host, port, ..
193 } = &cmd
194 {
195 let registry = camel_component_api::tls_source::TlsReloadRegistry::global();
196 match registry.find(scheme, host, *port) {
197 Some(handler) => {
198 handler.reload().await?;
199 if let Some(metrics) = &self.metrics {
201 metrics.record_counter(
203 "tls_reloads_total",
204 1.0,
205 &[("scheme", scheme), ("host", host)],
206 );
207 }
208 return Ok(RuntimeCommandResult::TlsCertsReloaded {
209 scheme: scheme.clone(),
210 host: host.clone(),
211 port: *port,
212 });
213 }
214 None => {
215 return Err(CamelError::Config(format!(
216 "no TLS server found for {scheme}://{host}:{port}"
217 )));
218 }
219 }
220 }
221 if let RuntimeCommand::ReloadTemplates { route_id, .. } = &cmd {
230 let route_id = route_id.clone();
231 camel_component_api::template_reload::TemplateReloadRegistry::global()
232 .reload_route(&route_id)
233 .await?;
234 if let Some(metrics) = &self.metrics {
236 metrics.record_counter(
238 "template_reloads_total",
239 1.0,
240 &[("route_id", route_id.as_str())],
241 );
242 }
243 return Ok(RuntimeCommandResult::TemplatesReloaded { route_id });
244 }
245 self.ensure_journal_recovered().await?;
248 let command_id = cmd.command_id().to_string();
249 if !self.dedup.first_seen(&command_id).await? {
250 return Ok(RuntimeCommandResult::Duplicate { command_id });
251 }
252 let deps = self.deps();
253 match execute_command(&deps, cmd).await {
254 Ok(result) => Ok(result),
255 Err(err) => {
256 let _ = self.dedup.forget_seen(&command_id).await;
257 Err(err)
258 }
259 }
260 }
261}
262
263#[async_trait]
264impl RuntimeQueryBus for RuntimeBus {
265 async fn ask(&self, query: RuntimeQuery) -> Result<RuntimeQueryResult, CamelError> {
266 self.ensure_journal_recovered().await?;
267
268 match query {
269 RuntimeQuery::InFlightCount { route_id } => {
270 if let Some(execution) = &self.execution {
271 execution
272 .in_flight_count(&route_id)
273 .await
274 .map(|r| r.into())
275 .map_err(Into::into)
276 } else {
277 Ok(RuntimeQueryResult::RouteNotFound { route_id })
278 }
279 }
280 other => {
281 let deps = self.query_deps();
282 execute_query(&deps, other).await
283 }
284 }
285 }
286}
287
288#[async_trait]
289impl RouteRegistrationPort for RuntimeBus {
290 async fn register_route(&self, def: RouteDefinition) -> Result<(), DomainError> {
291 self.ensure_journal_recovered()
292 .await
293 .map_err(|e| DomainError::InvalidState(e.to_string()))?;
294 let deps = self.deps();
295 handle_register_internal(&deps, def)
296 .await
297 .map(|_| ())
298 .map_err(|e| match e {
299 CamelError::RouteError(msg) => DomainError::InvalidState(msg),
300 other => DomainError::InvalidState(other.to_string()),
301 })
302 }
303}
304
305#[cfg(test)]
306mod tests {
307 use crate::lifecycle::domain::DomainError;
308
309 use super::*;
310 use std::collections::{HashMap, HashSet};
311 use std::sync::Mutex;
312
313 use crate::lifecycle::application::ports::RouteRegistrationPort as InternalRuntimeCommandBus;
314 use crate::lifecycle::application::ports::RouteStatusProjection;
315 use crate::lifecycle::application::route_definition::RouteDefinition;
316 use crate::lifecycle::domain::{RouteRuntimeAggregate, RuntimeEvent};
317
318 #[derive(Clone, Default)]
319 struct InMemoryTestRepo {
320 routes: Arc<Mutex<HashMap<String, RouteRuntimeAggregate>>>,
321 }
322
323 #[async_trait]
324 impl RouteRepositoryPort for InMemoryTestRepo {
325 async fn load(&self, route_id: &str) -> Result<Option<RouteRuntimeAggregate>, DomainError> {
326 Ok(self
327 .routes
328 .lock()
329 .expect("lock test routes")
330 .get(route_id)
331 .cloned())
332 }
333
334 async fn save(&self, aggregate: RouteRuntimeAggregate) -> Result<(), DomainError> {
335 self.routes
336 .lock()
337 .expect("lock test routes")
338 .insert(aggregate.route_id().to_string(), aggregate);
339 Ok(())
340 }
341
342 async fn save_if_version(
343 &self,
344 aggregate: RouteRuntimeAggregate,
345 expected_version: u64,
346 ) -> Result<(), DomainError> {
347 let route_id = aggregate.route_id().to_string();
348 let mut routes = self.routes.lock().expect("lock test routes");
349 let current = routes.get(&route_id).ok_or_else(|| {
350 DomainError::InvalidState(format!(
351 "optimistic lock conflict for route '{route_id}': route not found"
352 ))
353 })?;
354
355 if current.version() != expected_version {
356 return Err(DomainError::InvalidState(format!(
357 "optimistic lock conflict for route '{route_id}': expected version {expected_version}, actual {}",
358 current.version()
359 )));
360 }
361
362 routes.insert(route_id, aggregate);
363 Ok(())
364 }
365
366 async fn delete(&self, route_id: &str) -> Result<(), DomainError> {
367 self.routes
368 .lock()
369 .expect("lock test routes")
370 .remove(route_id);
371 Ok(())
372 }
373 }
374
375 #[derive(Clone, Default)]
376 struct InMemoryTestProjectionStore {
377 statuses: Arc<Mutex<HashMap<String, RouteStatusProjection>>>,
378 }
379
380 #[async_trait]
381 impl ProjectionStorePort for InMemoryTestProjectionStore {
382 async fn upsert_status(&self, status: RouteStatusProjection) -> Result<(), DomainError> {
383 self.statuses
384 .lock()
385 .expect("lock test statuses")
386 .insert(status.route_id.clone(), status);
387 Ok(())
388 }
389
390 async fn get_status(
391 &self,
392 route_id: &str,
393 ) -> Result<Option<RouteStatusProjection>, DomainError> {
394 Ok(self
395 .statuses
396 .lock()
397 .expect("lock test statuses")
398 .get(route_id)
399 .cloned())
400 }
401
402 async fn list_statuses(&self) -> Result<Vec<RouteStatusProjection>, DomainError> {
403 Ok(self
404 .statuses
405 .lock()
406 .expect("lock test statuses")
407 .values()
408 .cloned()
409 .collect())
410 }
411
412 async fn remove_status(&self, route_id: &str) -> Result<(), DomainError> {
413 self.statuses
414 .lock()
415 .expect("lock test statuses")
416 .remove(route_id);
417 Ok(())
418 }
419 }
420
421 #[derive(Clone, Default)]
422 struct InMemoryTestEventPublisher;
423
424 #[async_trait]
425 impl EventPublisherPort for InMemoryTestEventPublisher {
426 async fn publish(&self, _events: &[RuntimeEvent]) -> Result<(), DomainError> {
427 Ok(())
428 }
429 }
430
431 #[derive(Clone, Default)]
432 struct InMemoryTestDedup {
433 seen: Arc<Mutex<HashSet<String>>>,
434 }
435
436 #[derive(Clone, Default)]
437 struct InspectableDedup {
438 seen: Arc<Mutex<HashSet<String>>>,
439 forget_calls: Arc<Mutex<u32>>,
440 }
441
442 #[async_trait]
443 impl CommandDedupPort for InMemoryTestDedup {
444 async fn first_seen(&self, command_id: &str) -> Result<bool, DomainError> {
445 let mut seen = self.seen.lock().expect("lock dedup set");
446 Ok(seen.insert(command_id.to_string()))
447 }
448
449 async fn forget_seen(&self, command_id: &str) -> Result<(), DomainError> {
450 self.seen.lock().expect("lock dedup set").remove(command_id);
451 Ok(())
452 }
453 }
454
455 #[async_trait]
456 impl CommandDedupPort for InspectableDedup {
457 async fn first_seen(&self, command_id: &str) -> Result<bool, DomainError> {
458 let mut seen = self.seen.lock().expect("lock dedup set");
459 Ok(seen.insert(command_id.to_string()))
460 }
461
462 async fn forget_seen(&self, command_id: &str) -> Result<(), DomainError> {
463 self.seen.lock().expect("lock dedup set").remove(command_id);
464 let mut calls = self.forget_calls.lock().expect("forget calls");
465 *calls += 1;
466 Ok(())
467 }
468 }
469
470 fn build_test_runtime_bus() -> RuntimeBus {
471 let repo: Arc<dyn RouteRepositoryPort> = Arc::new(InMemoryTestRepo::default());
472 let projections: Arc<dyn ProjectionStorePort> =
473 Arc::new(InMemoryTestProjectionStore::default());
474 let events: Arc<dyn EventPublisherPort> = Arc::new(InMemoryTestEventPublisher);
475 let dedup: Arc<dyn CommandDedupPort> = Arc::new(InMemoryTestDedup::default());
476 RuntimeBus::new(repo, projections, events, dedup)
477 }
478
479 #[derive(Default)]
480 struct CountingUow {
481 recover_calls: Arc<Mutex<u32>>,
482 }
483
484 #[derive(Default)]
485 struct FailingRecoverUow;
486
487 #[async_trait]
488 impl RuntimeUnitOfWorkPort for CountingUow {
489 async fn persist_upsert(
490 &self,
491 _aggregate: RouteRuntimeAggregate,
492 _expected_version: Option<u64>,
493 _projection: RouteStatusProjection,
494 _events: &[RuntimeEvent],
495 ) -> Result<(), DomainError> {
496 Ok(())
497 }
498
499 async fn persist_delete(
500 &self,
501 _route_id: &str,
502 _events: &[RuntimeEvent],
503 ) -> Result<(), DomainError> {
504 Ok(())
505 }
506
507 async fn recover_from_journal(&self) -> Result<(), DomainError> {
508 let mut calls = self.recover_calls.lock().expect("recover_calls");
509 *calls += 1;
510 Ok(())
511 }
512 }
513
514 #[async_trait]
515 impl RuntimeUnitOfWorkPort for FailingRecoverUow {
516 async fn persist_upsert(
517 &self,
518 _aggregate: RouteRuntimeAggregate,
519 _expected_version: Option<u64>,
520 _projection: RouteStatusProjection,
521 _events: &[RuntimeEvent],
522 ) -> Result<(), DomainError> {
523 Ok(())
524 }
525
526 async fn persist_delete(
527 &self,
528 _route_id: &str,
529 _events: &[RuntimeEvent],
530 ) -> Result<(), DomainError> {
531 Ok(())
532 }
533
534 async fn recover_from_journal(&self) -> Result<(), DomainError> {
535 Err(DomainError::InvalidState("recover failed".into()))
536 }
537 }
538
539 #[derive(Default)]
540 struct InFlightExecutionPort;
541
542 #[async_trait]
543 impl RuntimeExecutionPort for InFlightExecutionPort {
544 async fn register_route(&self, _definition: RouteDefinition) -> Result<(), DomainError> {
545 Ok(())
546 }
547 async fn start_route(&self, _route_id: &str) -> Result<(), DomainError> {
548 Ok(())
549 }
550 async fn stop_route(&self, _route_id: &str) -> Result<(), DomainError> {
551 Ok(())
552 }
553 async fn suspend_route(&self, _route_id: &str) -> Result<(), DomainError> {
554 Ok(())
555 }
556 async fn resume_route(&self, _route_id: &str) -> Result<(), DomainError> {
557 Ok(())
558 }
559 async fn reload_route(&self, _route_id: &str) -> Result<(), DomainError> {
560 Ok(())
561 }
562 async fn remove_route(&self, _route_id: &str) -> Result<(), DomainError> {
563 Ok(())
564 }
565 async fn in_flight_count(
566 &self,
567 route_id: &str,
568 ) -> Result<InFlightCountResult, DomainError> {
569 if route_id == "known" {
570 Ok(InFlightCountResult::InFlightCount {
571 route_id: route_id.to_string(),
572 count: 3,
573 })
574 } else {
575 Ok(InFlightCountResult::RouteNotFound {
576 route_id: route_id.to_string(),
577 })
578 }
579 }
580 }
581
582 #[tokio::test]
583 async fn runtime_bus_implements_internal_command_bus() {
584 let bus = build_test_runtime_bus();
585 let def = RouteDefinition::new("timer:test", vec![]).with_route_id("internal-route");
586 let result = InternalRuntimeCommandBus::register_route(&bus, def).await;
587 assert!(
588 result.is_ok(),
589 "internal bus registration failed: {:?}",
590 result
591 );
592
593 let status = bus
594 .ask(RuntimeQuery::GetRouteStatus {
595 route_id: "internal-route".to_string(),
596 })
597 .await
598 .unwrap();
599 match status {
600 RuntimeQueryResult::RouteStatus { status, .. } => {
601 assert_eq!(status, "Registered");
602 }
603 _ => panic!("unexpected query result"),
604 }
605 }
606
607 #[tokio::test]
608 async fn execute_returns_duplicate_for_replayed_command_id() {
609 use camel_api::runtime::{CanonicalRouteSpec, CanonicalStepSpec, RuntimeCommand};
610
611 let bus = build_test_runtime_bus();
612
613 let mut spec = CanonicalRouteSpec::new("dup-route", "timer:tick");
614 spec.steps = vec![CanonicalStepSpec::Stop];
615
616 let cmd = RuntimeCommand::RegisterRoute {
617 spec: spec.clone(),
618 command_id: "dup-cmd".into(),
619 causation_id: None,
620 };
621 let first = bus.execute(cmd).await.unwrap();
622 assert!(matches!(
623 first,
624 RuntimeCommandResult::RouteRegistered { route_id } if route_id == "dup-route"
625 ));
626
627 let second = bus
628 .execute(RuntimeCommand::RegisterRoute {
629 spec,
630 command_id: "dup-cmd".into(),
631 causation_id: None,
632 })
633 .await
634 .unwrap();
635 assert!(matches!(
636 second,
637 RuntimeCommandResult::Duplicate { command_id } if command_id == "dup-cmd"
638 ));
639 }
640
641 #[tokio::test]
642 async fn ask_in_flight_count_without_execution_returns_route_not_found() {
643 let bus = build_test_runtime_bus();
644 let res = bus
645 .ask(RuntimeQuery::InFlightCount {
646 route_id: "missing".into(),
647 })
648 .await
649 .unwrap();
650 assert!(matches!(
651 res,
652 RuntimeQueryResult::RouteNotFound { route_id } if route_id == "missing"
653 ));
654 }
655
656 #[tokio::test]
657 async fn ask_in_flight_count_with_execution_delegates_to_adapter() {
658 let repo: Arc<dyn RouteRepositoryPort> = Arc::new(InMemoryTestRepo::default());
659 let projections: Arc<dyn ProjectionStorePort> =
660 Arc::new(InMemoryTestProjectionStore::default());
661 let events: Arc<dyn EventPublisherPort> = Arc::new(InMemoryTestEventPublisher);
662 let dedup: Arc<dyn CommandDedupPort> = Arc::new(InMemoryTestDedup::default());
663 let execution: Arc<dyn RuntimeExecutionPort> = Arc::new(InFlightExecutionPort);
664 let bus = RuntimeBus::new(repo, projections, events, dedup).with_execution(execution);
665
666 let known = bus
667 .ask(RuntimeQuery::InFlightCount {
668 route_id: "known".into(),
669 })
670 .await
671 .unwrap();
672 assert!(matches!(
673 known,
674 RuntimeQueryResult::InFlightCount { route_id, count }
675 if route_id == "known" && count == 3
676 ));
677 }
678
679 #[tokio::test]
680 async fn journal_recovery_runs_once_even_with_multiple_commands() {
681 use camel_api::runtime::{CanonicalRouteSpec, CanonicalStepSpec, RuntimeCommand};
682
683 let repo: Arc<dyn RouteRepositoryPort> = Arc::new(InMemoryTestRepo::default());
684 let projections: Arc<dyn ProjectionStorePort> =
685 Arc::new(InMemoryTestProjectionStore::default());
686 let events: Arc<dyn EventPublisherPort> = Arc::new(InMemoryTestEventPublisher);
687 let dedup: Arc<dyn CommandDedupPort> = Arc::new(InMemoryTestDedup::default());
688 let uow = Arc::new(CountingUow::default());
689 let bus = RuntimeBus::new(repo, projections, events, dedup).with_uow(uow.clone());
690
691 let mut spec_a = CanonicalRouteSpec::new("a", "timer:a");
692 spec_a.steps = vec![CanonicalStepSpec::Stop];
693 let mut spec_b = CanonicalRouteSpec::new("b", "timer:b");
694 spec_b.steps = vec![CanonicalStepSpec::Stop];
695
696 bus.execute(RuntimeCommand::RegisterRoute {
697 spec: spec_a,
698 command_id: "c-a".into(),
699 causation_id: None,
700 })
701 .await
702 .unwrap();
703
704 bus.execute(RuntimeCommand::RegisterRoute {
705 spec: spec_b,
706 command_id: "c-b".into(),
707 causation_id: None,
708 })
709 .await
710 .unwrap();
711
712 let calls = *uow.recover_calls.lock().expect("recover calls");
713 assert_eq!(calls, 1, "journal recovery should run once");
714 }
715
716 #[tokio::test]
717 async fn execute_on_command_error_forgets_dedup_marker() {
718 use camel_api::runtime::{CanonicalRouteSpec, RuntimeCommand};
719
720 let repo: Arc<dyn RouteRepositoryPort> = Arc::new(InMemoryTestRepo::default());
721 let projections: Arc<dyn ProjectionStorePort> =
722 Arc::new(InMemoryTestProjectionStore::default());
723 let events: Arc<dyn EventPublisherPort> = Arc::new(InMemoryTestEventPublisher);
724 let dedup = Arc::new(InspectableDedup::default());
725 let dedup_port: Arc<dyn CommandDedupPort> = dedup.clone();
726
727 let bus = RuntimeBus::new(repo, projections, events, dedup_port);
728
729 let cmd = RuntimeCommand::RegisterRoute {
731 spec: CanonicalRouteSpec::new("", "timer:tick"),
732 command_id: "err-cmd".into(),
733 causation_id: None,
734 };
735
736 let err = bus.execute(cmd).await.expect_err("must fail");
737 assert!(err.to_string().contains("route_id cannot be empty"));
738
739 assert_eq!(*dedup.forget_calls.lock().expect("forget calls"), 1);
740 assert!(!dedup.seen.lock().expect("seen").contains("err-cmd"));
741 }
742
743 #[tokio::test]
744 async fn execute_propagates_recover_error_from_uow() {
745 use camel_api::runtime::{CanonicalRouteSpec, CanonicalStepSpec, RuntimeCommand};
746
747 let repo: Arc<dyn RouteRepositoryPort> = Arc::new(InMemoryTestRepo::default());
748 let projections: Arc<dyn ProjectionStorePort> =
749 Arc::new(InMemoryTestProjectionStore::default());
750 let events: Arc<dyn EventPublisherPort> = Arc::new(InMemoryTestEventPublisher);
751 let dedup: Arc<dyn CommandDedupPort> = Arc::new(InMemoryTestDedup::default());
752 let uow: Arc<dyn RuntimeUnitOfWorkPort> = Arc::new(FailingRecoverUow);
753
754 let bus = RuntimeBus::new(repo, projections, events, dedup).with_uow(uow);
755
756 let mut spec = CanonicalRouteSpec::new("x", "timer:x");
757 spec.steps = vec![CanonicalStepSpec::Stop];
758 let err = bus
759 .execute(RuntimeCommand::RegisterRoute {
760 spec,
761 command_id: "recover-err".into(),
762 causation_id: None,
763 })
764 .await
765 .expect_err("recover should fail");
766
767 assert!(err.to_string().contains("recover failed"));
768 }
769
770 #[tokio::test]
771 async fn ask_in_flight_count_with_execution_handles_unknown_route() {
772 let repo: Arc<dyn RouteRepositoryPort> = Arc::new(InMemoryTestRepo::default());
773 let projections: Arc<dyn ProjectionStorePort> =
774 Arc::new(InMemoryTestProjectionStore::default());
775 let events: Arc<dyn EventPublisherPort> = Arc::new(InMemoryTestEventPublisher);
776 let dedup: Arc<dyn CommandDedupPort> = Arc::new(InMemoryTestDedup::default());
777 let execution: Arc<dyn RuntimeExecutionPort> = Arc::new(InFlightExecutionPort);
778 let bus = RuntimeBus::new(repo, projections, events, dedup).with_execution(execution);
779
780 let unknown = bus
781 .ask(RuntimeQuery::InFlightCount {
782 route_id: "unknown".into(),
783 })
784 .await
785 .unwrap();
786 assert!(matches!(
787 unknown,
788 RuntimeQueryResult::RouteNotFound { route_id } if route_id == "unknown"
789 ));
790 }
791
792 #[tokio::test]
793 async fn watcher_duplicate_failroute_is_noop() {
794 use crate::lifecycle::domain::RouteRuntimeState;
799 use camel_api::RuntimeCommand;
800
801 let bus = build_test_runtime_bus();
802 let def = RouteDefinition::new("timer:test", vec![]).with_route_id("dup-fail-route");
803 InternalRuntimeCommandBus::register_route(&bus, def)
804 .await
805 .expect("register route");
806
807 let cmd = RuntimeCommand::FailRoute {
808 route_id: "dup-fail-route".into(),
809 error: "watcher".into(),
810 command_id: "fail-once".into(),
811 causation_id: None,
812 };
813 let first = bus.execute(cmd.clone()).await.unwrap();
814 assert!(
815 matches!(
816 &first,
817 RuntimeCommandResult::RouteStateChanged { status, .. } if status == "Failed"
818 ),
819 "first FailRoute must transition to Failed, got {first:?}"
820 );
821
822 let second = bus.execute(cmd).await.unwrap();
823 assert!(
824 matches!(
825 &second,
826 RuntimeCommandResult::Duplicate { command_id } if command_id == "fail-once"
827 ),
828 "duplicate command_id must return the dedup no-op, got {second:?}"
829 );
830
831 let aggregate = bus
834 .repo()
835 .load("dup-fail-route")
836 .await
837 .unwrap()
838 .expect("route exists");
839 assert!(matches!(aggregate.state(), RouteRuntimeState::Failed(_)));
840 assert_eq!(
841 aggregate.version(),
842 1,
843 "the duplicate must not apply a second transition"
844 );
845
846 let status = bus
847 .ask(RuntimeQuery::GetRouteStatus {
848 route_id: "dup-fail-route".into(),
849 })
850 .await
851 .unwrap();
852 match status {
853 RuntimeQueryResult::RouteStatus { status, .. } => assert_eq!(status, "Failed"),
854 other => panic!("unexpected query result: {other:?}"),
855 }
856 }
857}