Skip to main content

camel_core/lifecycle/application/
runtime_bus.rs

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    /// Thread a metrics handle so infrastructure commands (ReloadTlsCerts,
88    /// ReloadTemplates) can record counters (rc-d3pj). When None (default),
89    /// no counters are recorded.
90    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    /// H8 boot reconciliation: fail any route still in a transient state
100    /// (`Starting` / `Stopping`) from a previous run. Called from
101    /// `CamelContext::start()` before `auto_startup_route_ids()`.
102    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    /// Journal-derived boot nonce from the recovered durable dedup store.
178    /// Yields `0` before journal recovery ran or when no unit-of-work is
179    /// installed (no durable journal support).
180    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        // ── TLS cert reload intercept ──────────────────────────────────────
189        // Infrastructure command — bypasses journal recovery + dedup.
190        // Reloads are idempotent and NOT journaled.
191        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                    // rc-d3pj: record reload counter once per successful reload.
200                    if let Some(metrics) = &self.metrics {
201                        // allow-open-label rc-ycts (scheme/host: TLS listener config values)
202                        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        // ── End TLS reload intercept ────────────────────────────────────────
222
223        // ── Template reload intercept ──────────────────────────────────────
224        // Infrastructure command — bypasses journal recovery + dedup (mirrors
225        // ReloadTlsCerts above). Reloads are idempotent and NOT journaled;
226        // RouteStatus is not mutated and ADR-0018 is not invoked. Dispatches to
227        // the registry in camel-component-api (NOT camel-template — that would
228        // invert the dependency).
229        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            // rc-d3pj: record reload counter once per successful reload.
235            if let Some(metrics) = &self.metrics {
236                // allow-open-label rc-xl5k (route label: user-defined route id)
237                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        // ── End template reload intercept ──────────────────────────────────
246
247        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        // Invalid canonical contract: empty route_id -> execute_command should fail.
730        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        // rc-slvd: bus-level dedup pin — the failure watcher retries with
795        // the SAME command_id; the second execute must return the dedup
796        // outcome and cause exactly ONE lifecycle transition to Failed.
797        // REAL bus (a fake would make dedup vacuously green).
798        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        // Exactly ONE transition: aggregate Failed at version 1 (0 was the
832        // register), and the status query agrees.
833        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}