Skip to main content

camel_core/lifecycle/adapters/
in_memory.rs

1use std::collections::{HashMap, HashSet};
2use std::sync::Arc;
3
4use async_trait::async_trait;
5use camel_api::MetricsCollector;
6use tokio::sync::{Mutex, RwLock};
7
8use crate::lifecycle::domain::DomainError;
9
10use crate::lifecycle::application::ports::{
11    CommandDedupPort, EventPublisherPort, ProjectionStorePort, RouteRepositoryPort,
12    RouteStatusProjection, RuntimeEventJournalPort, RuntimeUnitOfWorkPort,
13};
14use crate::lifecycle::domain::{RouteRuntimeAggregate, RouteRuntimeState, RuntimeEvent};
15
16#[derive(Default, Clone)]
17pub struct InMemoryRouteRepository {
18    routes: Arc<RwLock<HashMap<String, RouteRuntimeAggregate>>>,
19}
20
21#[async_trait]
22impl RouteRepositoryPort for InMemoryRouteRepository {
23    async fn load(&self, route_id: &str) -> Result<Option<RouteRuntimeAggregate>, DomainError> {
24        let routes = self.routes.read().await;
25        Ok(routes.get(route_id).cloned())
26    }
27
28    async fn save(&self, aggregate: RouteRuntimeAggregate) -> Result<(), DomainError> {
29        let mut routes = self.routes.write().await;
30        routes.insert(aggregate.route_id().to_string(), aggregate);
31        Ok(())
32    }
33
34    async fn save_if_version(
35        &self,
36        aggregate: RouteRuntimeAggregate,
37        expected_version: u64,
38    ) -> Result<(), DomainError> {
39        let mut routes = self.routes.write().await;
40        let route_id = aggregate.route_id().to_string();
41        let current = routes.get(&route_id).ok_or_else(|| {
42            DomainError::InvalidState(format!(
43                "optimistic lock conflict for route '{route_id}': route not found"
44            ))
45        })?;
46
47        if current.version() != expected_version {
48            return Err(DomainError::InvalidState(format!(
49                "optimistic lock conflict for route '{route_id}': expected version {expected_version}, actual {}",
50                current.version()
51            )));
52        }
53
54        routes.insert(route_id, aggregate);
55        Ok(())
56    }
57
58    async fn delete(&self, route_id: &str) -> Result<(), DomainError> {
59        let mut routes = self.routes.write().await;
60        routes.remove(route_id);
61        Ok(())
62    }
63}
64
65#[derive(Default, Clone)]
66pub struct InMemoryProjectionStore {
67    statuses: Arc<RwLock<HashMap<String, RouteStatusProjection>>>,
68}
69
70#[async_trait]
71impl ProjectionStorePort for InMemoryProjectionStore {
72    async fn upsert_status(&self, status: RouteStatusProjection) -> Result<(), DomainError> {
73        let mut statuses = self.statuses.write().await;
74        statuses.insert(status.route_id.clone(), status);
75        Ok(())
76    }
77
78    async fn get_status(
79        &self,
80        route_id: &str,
81    ) -> Result<Option<RouteStatusProjection>, DomainError> {
82        let statuses = self.statuses.read().await;
83        Ok(statuses.get(route_id).cloned())
84    }
85
86    async fn list_statuses(&self) -> Result<Vec<RouteStatusProjection>, DomainError> {
87        let statuses = self.statuses.read().await;
88        Ok(statuses.values().cloned().collect())
89    }
90
91    async fn remove_status(&self, route_id: &str) -> Result<(), DomainError> {
92        let mut statuses = self.statuses.write().await;
93        statuses.remove(route_id);
94        Ok(())
95    }
96}
97
98#[derive(Default, Clone)]
99pub struct InMemoryEventPublisher {
100    events: Arc<RwLock<Vec<RuntimeEvent>>>,
101}
102
103impl InMemoryEventPublisher {
104    pub async fn snapshot(&self) -> Vec<RuntimeEvent> {
105        self.events.read().await.clone()
106    }
107}
108
109#[async_trait]
110impl EventPublisherPort for InMemoryEventPublisher {
111    async fn publish(&self, events: &[RuntimeEvent]) -> Result<(), DomainError> {
112        let mut stored = self.events.write().await;
113        stored.extend(events.iter().cloned());
114        Ok(())
115    }
116}
117
118#[derive(Default, Clone)]
119pub struct InMemoryCommandDedup {
120    seen: Arc<RwLock<HashSet<String>>>,
121}
122
123#[async_trait]
124impl CommandDedupPort for InMemoryCommandDedup {
125    async fn first_seen(&self, command_id: &str) -> Result<bool, DomainError> {
126        let mut seen = self.seen.write().await;
127        Ok(seen.insert(command_id.to_string()))
128    }
129
130    async fn forget_seen(&self, command_id: &str) -> Result<(), DomainError> {
131        let mut seen = self.seen.write().await;
132        seen.remove(command_id);
133        Ok(())
134    }
135}
136
137#[derive(Clone)]
138pub struct InMemoryRuntimeStore {
139    inner: Arc<Mutex<RuntimeStoreState>>,
140    journal: Option<Arc<dyn RuntimeEventJournalPort>>,
141    /// Shared late-bound metrics handle (dashboard-observability T3.1).
142    /// Seeded once by `CamelContextBuilder::build()` — never a fresh
143    /// collector — so every projection transition published through this
144    /// store emits `MetricsCollector::set_route_state`.
145    metrics: Option<Arc<dyn MetricsCollector>>,
146}
147
148#[derive(Default)]
149struct RuntimeStoreState {
150    routes: HashMap<String, RouteRuntimeAggregate>,
151    statuses: HashMap<String, RouteStatusProjection>,
152    events: Vec<RuntimeEvent>,
153    seen: HashSet<String>,
154    /// Route ids reconstructed by the most recent journal replay that have not
155    /// yet been re-registered in this process. Drained by `take_recovered` so
156    /// the first declarative-boot registration adopts the persisted aggregate
157    /// while a later same-session duplicate is still rejected.
158    recovered: HashSet<String>,
159}
160
161impl InMemoryRuntimeStore {
162    pub fn with_journal(mut self, journal: Arc<dyn RuntimeEventJournalPort>) -> Self {
163        self.journal = Some(journal);
164        self
165    }
166
167    /// Seed the shared late-bound metrics handle. Called from
168    /// `CamelContextBuilder::build()` with the same handle threaded to the
169    /// route controller and the runtime bus (no new collector instances).
170    pub fn with_metrics(mut self, metrics: Arc<dyn MetricsCollector>) -> Self {
171        self.metrics = Some(metrics);
172        self
173    }
174
175    /// Emit a route-state gauge transition for a projection write. `None`
176    /// handle (tests, standalone stores) = no-op.
177    fn emit_route_state(&self, route_id: &str, state: &str) {
178        if let Some(metrics) = &self.metrics {
179            metrics.set_route_state(route_id, state);
180        }
181    }
182
183    pub async fn snapshot_events(&self) -> Vec<RuntimeEvent> {
184        self.inner.lock().await.events.clone()
185    }
186}
187
188fn upsert_replayed_route(
189    state: &mut RuntimeStoreState,
190    route_id: &str,
191    next_state: RouteRuntimeState,
192    status: &str,
193    increment_version: bool,
194) {
195    let current_version = state
196        .routes
197        .get(route_id)
198        .map(|agg| agg.version())
199        .unwrap_or(0);
200    let next_version = if increment_version {
201        current_version.saturating_add(1)
202    } else {
203        current_version
204    };
205    state.routes.insert(
206        route_id.to_string(),
207        RouteRuntimeAggregate::from_snapshot(route_id, next_state, next_version),
208    );
209    state.statuses.insert(
210        route_id.to_string(),
211        RouteStatusProjection {
212            route_id: route_id.to_string(),
213            status: status.to_string(),
214        },
215    );
216}
217
218fn state_label(state: &RouteRuntimeState) -> &'static str {
219    match state {
220        RouteRuntimeState::Registered => "Registered",
221        RouteRuntimeState::Starting => "Starting",
222        RouteRuntimeState::Started => "Started",
223        RouteRuntimeState::Suspended => "Suspended",
224        RouteRuntimeState::Stopping => "Stopping",
225        RouteRuntimeState::Stopped => "Stopped",
226        RouteRuntimeState::Failed(_) => "Failed",
227    }
228}
229
230fn apply_replayed_event(state: &mut RuntimeStoreState, event: &RuntimeEvent) {
231    match event {
232        RuntimeEvent::RouteRegistered { route_id } => {
233            state.routes.insert(
234                route_id.clone(),
235                RouteRuntimeAggregate::new(route_id.clone()),
236            );
237            state.statuses.insert(
238                route_id.clone(),
239                RouteStatusProjection {
240                    route_id: route_id.clone(),
241                    status: "Registered".to_string(),
242                },
243            );
244        }
245        RuntimeEvent::RouteRemoved { route_id } => {
246            state.routes.remove(route_id);
247            state.statuses.remove(route_id);
248        }
249        _ => {
250            // Use the domain's state machine to derive the next state from the event.
251            let Some(next_state) = RouteRuntimeAggregate::state_from_event(event) else {
252                return;
253            };
254            let route_id = match event {
255                RuntimeEvent::RouteStartRequested { route_id }
256                | RuntimeEvent::RouteStarted { route_id }
257                | RuntimeEvent::RouteFailed { route_id, .. }
258                | RuntimeEvent::RouteStopped { route_id }
259                | RuntimeEvent::RouteSuspended { route_id }
260                | RuntimeEvent::RouteResumed { route_id }
261                | RuntimeEvent::RouteReloaded { route_id } => route_id,
262                _ => return,
263            };
264            let status = state_label(&next_state);
265            let increment_version = !matches!(
266                (event, state.routes.get(route_id).map(|agg| agg.state())),
267                (
268                    RuntimeEvent::RouteStarted { .. },
269                    Some(RouteRuntimeState::Starting)
270                )
271            );
272            upsert_replayed_route(state, route_id, next_state, status, increment_version);
273        }
274    }
275}
276
277impl Default for InMemoryRuntimeStore {
278    fn default() -> Self {
279        Self {
280            inner: Arc::new(Mutex::new(RuntimeStoreState::default())),
281            journal: None,
282            metrics: None,
283        }
284    }
285}
286
287#[async_trait]
288impl RouteRepositoryPort for InMemoryRuntimeStore {
289    async fn load(&self, route_id: &str) -> Result<Option<RouteRuntimeAggregate>, DomainError> {
290        let guard = self.inner.lock().await;
291        Ok(guard.routes.get(route_id).cloned())
292    }
293
294    async fn save(&self, aggregate: RouteRuntimeAggregate) -> Result<(), DomainError> {
295        let mut guard = self.inner.lock().await;
296        guard
297            .routes
298            .insert(aggregate.route_id().to_string(), aggregate);
299        Ok(())
300    }
301
302    async fn save_if_version(
303        &self,
304        aggregate: RouteRuntimeAggregate,
305        expected_version: u64,
306    ) -> Result<(), DomainError> {
307        let mut guard = self.inner.lock().await;
308        let route_id = aggregate.route_id().to_string();
309        let current = guard.routes.get(&route_id).ok_or_else(|| {
310            DomainError::InvalidState(format!(
311                "optimistic lock conflict for route '{route_id}': route not found"
312            ))
313        })?;
314
315        if current.version() != expected_version {
316            return Err(DomainError::InvalidState(format!(
317                "optimistic lock conflict for route '{route_id}': expected version {expected_version}, actual {}",
318                current.version()
319            )));
320        }
321
322        guard.routes.insert(route_id, aggregate);
323        Ok(())
324    }
325
326    async fn delete(&self, route_id: &str) -> Result<(), DomainError> {
327        let mut guard = self.inner.lock().await;
328        guard.routes.remove(route_id);
329        guard.recovered.remove(route_id);
330        Ok(())
331    }
332
333    async fn take_recovered(&self, route_id: &str) -> Result<bool, DomainError> {
334        let mut guard = self.inner.lock().await;
335        Ok(guard.recovered.remove(route_id))
336    }
337}
338
339#[async_trait]
340impl ProjectionStorePort for InMemoryRuntimeStore {
341    async fn upsert_status(&self, status: RouteStatusProjection) -> Result<(), DomainError> {
342        let route_id = status.route_id.clone();
343        let state = status.status.clone();
344        let mut guard = self.inner.lock().await;
345        guard.statuses.insert(status.route_id.clone(), status);
346        self.emit_route_state(&route_id, &state);
347        Ok(())
348    }
349
350    async fn get_status(
351        &self,
352        route_id: &str,
353    ) -> Result<Option<RouteStatusProjection>, DomainError> {
354        let guard = self.inner.lock().await;
355        Ok(guard.statuses.get(route_id).cloned())
356    }
357
358    async fn list_statuses(&self) -> Result<Vec<RouteStatusProjection>, DomainError> {
359        let guard = self.inner.lock().await;
360        Ok(guard.statuses.values().cloned().collect())
361    }
362
363    async fn remove_status(&self, route_id: &str) -> Result<(), DomainError> {
364        let mut guard = self.inner.lock().await;
365        guard.statuses.remove(route_id);
366        // Undeployed route: drop its gauge series so a scrape reflects
367        // only routes that exist.
368        if let Some(metrics) = &self.metrics {
369            metrics.clear_route_state(route_id);
370        }
371        Ok(())
372    }
373}
374
375#[async_trait]
376impl EventPublisherPort for InMemoryRuntimeStore {
377    async fn publish(&self, events: &[RuntimeEvent]) -> Result<(), DomainError> {
378        let mut guard = self.inner.lock().await;
379        if let Some(journal) = &self.journal {
380            journal.append_batch(events).await?;
381        }
382        guard.events.extend(events.iter().cloned());
383        Ok(())
384    }
385}
386
387#[async_trait]
388impl CommandDedupPort for InMemoryRuntimeStore {
389    async fn first_seen(&self, command_id: &str) -> Result<bool, DomainError> {
390        let mut guard = self.inner.lock().await;
391        if !guard.seen.insert(command_id.to_string()) {
392            return Ok(false);
393        }
394
395        if let Some(journal) = &self.journal
396            && let Err(err) = journal.append_command_id(command_id).await
397        {
398            guard.seen.remove(command_id);
399            return Err(err);
400        }
401
402        Ok(true)
403    }
404
405    async fn forget_seen(&self, command_id: &str) -> Result<(), DomainError> {
406        let mut guard = self.inner.lock().await;
407        let removed = guard.seen.remove(command_id);
408        if removed && let Some(journal) = &self.journal {
409            journal.remove_command_id(command_id).await?;
410        }
411        Ok(())
412    }
413}
414
415#[async_trait]
416impl RuntimeUnitOfWorkPort for InMemoryRuntimeStore {
417    async fn persist_upsert(
418        &self,
419        aggregate: RouteRuntimeAggregate,
420        expected_version: Option<u64>,
421        projection: RouteStatusProjection,
422        events: &[RuntimeEvent],
423    ) -> Result<(), DomainError> {
424        let mut guard = self.inner.lock().await;
425        if let Some(expected) = expected_version {
426            let route_id = aggregate.route_id().to_string();
427            let current = guard.routes.get(&route_id).ok_or_else(|| {
428                DomainError::InvalidState(format!(
429                    "optimistic lock conflict for route '{route_id}': route not found"
430                ))
431            })?;
432            if current.version() != expected {
433                return Err(DomainError::InvalidState(format!(
434                    "optimistic lock conflict for route '{route_id}': expected version {expected}, actual {}",
435                    current.version()
436                )));
437            }
438        }
439
440        if let Some(journal) = &self.journal {
441            journal.append_batch(events).await?;
442        }
443
444        guard
445            .routes
446            .insert(aggregate.route_id().to_string(), aggregate);
447        let route_id = projection.route_id.clone();
448        let state = projection.status.clone();
449        guard
450            .statuses
451            .insert(projection.route_id.clone(), projection);
452        guard.events.extend(events.iter().cloned());
453        self.emit_route_state(&route_id, &state);
454        Ok(())
455    }
456
457    async fn persist_delete(
458        &self,
459        route_id: &str,
460        events: &[RuntimeEvent],
461    ) -> Result<(), DomainError> {
462        let mut guard = self.inner.lock().await;
463        if let Some(journal) = &self.journal {
464            journal.append_batch(events).await?;
465        }
466        guard.routes.remove(route_id);
467        guard.statuses.remove(route_id);
468        guard.events.extend(events.iter().cloned());
469        if let Some(metrics) = &self.metrics {
470            metrics.clear_route_state(route_id);
471        }
472        Ok(())
473    }
474
475    async fn recover_from_journal(&self) -> Result<(), DomainError> {
476        let Some(journal) = &self.journal else {
477            return Ok(());
478        };
479
480        let replayed_events = journal.load_all().await?;
481        let replayed_command_ids = journal.load_command_ids().await?;
482
483        let mut guard = self.inner.lock().await;
484        guard.routes.clear();
485        guard.statuses.clear();
486        guard.events.clear();
487        guard.seen.clear();
488        guard.recovered.clear();
489
490        for event in &replayed_events {
491            apply_replayed_event(&mut guard, event);
492        }
493        // Every aggregate that survived replay is adoptable exactly once by the
494        // first re-registration of its id (declarative boot), so a durable
495        // journal no longer turns a clean restart into an AlreadyExists storm.
496        guard.recovered = guard.routes.keys().cloned().collect();
497        // Seed the route-state gauges from the recovered projections:
498        // replay writes statuses directly, bypassing the emit sites.
499        if let Some(metrics) = &self.metrics {
500            for status in guard.statuses.values() {
501                metrics.set_route_state(&status.route_id, &status.status);
502            }
503        }
504        guard.events = replayed_events;
505        for command_id in replayed_command_ids {
506            guard.seen.insert(command_id);
507        }
508        Ok(())
509    }
510}
511
512#[cfg(test)]
513mod tests {
514    use super::*;
515    use std::sync::Arc;
516
517    #[derive(Clone)]
518    struct ReplayJournal {
519        events: Vec<RuntimeEvent>,
520    }
521
522    #[async_trait]
523    impl RuntimeEventJournalPort for ReplayJournal {
524        async fn append_batch(&self, _events: &[RuntimeEvent]) -> Result<(), DomainError> {
525            Ok(())
526        }
527
528        async fn load_all(&self) -> Result<Vec<RuntimeEvent>, DomainError> {
529            Ok(self.events.clone())
530        }
531    }
532
533    #[tokio::test]
534    async fn repo_roundtrip_works() {
535        let repo = InMemoryRouteRepository::default();
536        repo.save(RouteRuntimeAggregate::new("r1")).await.unwrap();
537        assert!(repo.load("r1").await.unwrap().is_some());
538
539        let updated = RouteRuntimeAggregate::from_snapshot(
540            "r1",
541            crate::lifecycle::domain::RouteRuntimeState::Started,
542            1,
543        );
544        repo.save_if_version(updated.clone(), 0).await.unwrap();
545        let loaded = repo.load("r1").await.unwrap().unwrap();
546        assert_eq!(loaded.version(), 1);
547
548        let conflict = repo.save_if_version(updated, 0).await.unwrap_err();
549        assert!(
550            conflict.to_string().contains("optimistic lock conflict"),
551            "unexpected conflict error: {conflict}"
552        );
553
554        repo.delete("r1").await.unwrap();
555        assert!(repo.load("r1").await.unwrap().is_none());
556    }
557
558    #[tokio::test]
559    async fn projection_roundtrip_works() {
560        let store = InMemoryProjectionStore::default();
561        store
562            .upsert_status(RouteStatusProjection {
563                route_id: "r1".into(),
564                status: "Started".into(),
565            })
566            .await
567            .unwrap();
568
569        let status = store.get_status("r1").await.unwrap();
570        assert!(status.is_some());
571        assert_eq!(status.unwrap().status, "Started");
572        store.remove_status("r1").await.unwrap();
573        assert!(store.get_status("r1").await.unwrap().is_none());
574    }
575
576    #[tokio::test]
577    async fn event_publisher_stores_events() {
578        let publisher = InMemoryEventPublisher::default();
579        publisher
580            .publish(&[RuntimeEvent::RouteStarted {
581                route_id: "r1".into(),
582            }])
583            .await
584            .unwrap();
585
586        let events = publisher.snapshot().await;
587        assert_eq!(events.len(), 1);
588    }
589
590    #[tokio::test]
591    async fn command_dedup_detects_duplicates() {
592        let dedup = InMemoryCommandDedup::default();
593        assert!(dedup.first_seen("c1").await.unwrap());
594        assert!(!dedup.first_seen("c1").await.unwrap());
595        dedup.forget_seen("c1").await.unwrap();
596        assert!(dedup.first_seen("c1").await.unwrap());
597        assert!(dedup.first_seen("c2").await.unwrap());
598    }
599
600    #[tokio::test]
601    async fn runtime_store_uow_persists_all_three_writes() {
602        let store = InMemoryRuntimeStore::default();
603        let aggregate = RouteRuntimeAggregate::new("uow-r1");
604        let projection = RouteStatusProjection {
605            route_id: "uow-r1".to_string(),
606            status: "Registered".to_string(),
607        };
608        let events = vec![RuntimeEvent::RouteRegistered {
609            route_id: "uow-r1".to_string(),
610        }];
611
612        store
613            .persist_upsert(aggregate, None, projection.clone(), &events)
614            .await
615            .unwrap();
616
617        assert!(store.load("uow-r1").await.unwrap().is_some());
618        assert_eq!(
619            store.get_status("uow-r1").await.unwrap().unwrap(),
620            projection
621        );
622        assert_eq!(store.snapshot_events().await, events);
623    }
624
625    #[tokio::test]
626    async fn runtime_store_uow_enforces_expected_version() {
627        let store = InMemoryRuntimeStore::default();
628        let initial = RouteRuntimeAggregate::new("uow-r2");
629        let initial_projection = RouteStatusProjection {
630            route_id: "uow-r2".to_string(),
631            status: "Registered".to_string(),
632        };
633        store
634            .persist_upsert(
635                initial,
636                None,
637                initial_projection,
638                &[RuntimeEvent::RouteRegistered {
639                    route_id: "uow-r2".to_string(),
640                }],
641            )
642            .await
643            .unwrap();
644
645        let started = RouteRuntimeAggregate::from_snapshot(
646            "uow-r2",
647            crate::lifecycle::domain::RouteRuntimeState::Started,
648            1,
649        );
650        let err = store
651            .persist_upsert(
652                started,
653                Some(99),
654                RouteStatusProjection {
655                    route_id: "uow-r2".to_string(),
656                    status: "Started".to_string(),
657                },
658                &[RuntimeEvent::RouteStarted {
659                    route_id: "uow-r2".to_string(),
660                }],
661            )
662            .await
663            .unwrap_err()
664            .to_string();
665        assert!(
666            err.contains("optimistic lock conflict"),
667            "unexpected error: {err}"
668        );
669    }
670
671    #[tokio::test]
672    async fn replay_start_requested_only_advances_version_once() {
673        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
674            events: vec![
675                RuntimeEvent::RouteRegistered {
676                    route_id: "replay-r1".to_string(),
677                },
678                RuntimeEvent::RouteStartRequested {
679                    route_id: "replay-r1".to_string(),
680                },
681            ],
682        }));
683
684        store.recover_from_journal().await.unwrap();
685        let aggregate = store.load("replay-r1").await.unwrap().unwrap();
686
687        assert_eq!(aggregate.state(), &RouteRuntimeState::Starting);
688        assert_eq!(aggregate.version(), 1);
689    }
690
691    #[tokio::test]
692    async fn replay_start_requested_then_started_keeps_single_command_version() {
693        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
694            events: vec![
695                RuntimeEvent::RouteRegistered {
696                    route_id: "replay-r2".to_string(),
697                },
698                RuntimeEvent::RouteStartRequested {
699                    route_id: "replay-r2".to_string(),
700                },
701                RuntimeEvent::RouteStarted {
702                    route_id: "replay-r2".to_string(),
703                },
704            ],
705        }));
706
707        store.recover_from_journal().await.unwrap();
708        let aggregate = store.load("replay-r2").await.unwrap().unwrap();
709
710        assert_eq!(aggregate.state(), &RouteRuntimeState::Started);
711        assert_eq!(aggregate.version(), 1);
712    }
713
714    #[tokio::test]
715    async fn recover_from_journal_marks_replayed_routes_adoptable_once() {
716        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
717            events: vec![
718                RuntimeEvent::RouteRegistered {
719                    route_id: "replay-r3".to_string(),
720                },
721                RuntimeEvent::RouteStarted {
722                    route_id: "replay-r3".to_string(),
723                },
724            ],
725        }));
726
727        store.recover_from_journal().await.unwrap();
728
729        // First re-registration adopts the recovered aggregate.
730        assert!(store.take_recovered("replay-r3").await.unwrap());
731        // Marker is consumed: a same-session duplicate is no longer adoptable.
732        assert!(!store.take_recovered("replay-r3").await.unwrap());
733        // A route that never existed in the journal is never adoptable.
734        assert!(!store.take_recovered("never-seen").await.unwrap());
735    }
736
737    #[tokio::test]
738    async fn recover_from_journal_leaves_no_marker_for_removed_routes() {
739        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
740            events: vec![
741                RuntimeEvent::RouteRegistered {
742                    route_id: "replay-r4".to_string(),
743                },
744                RuntimeEvent::RouteRemoved {
745                    route_id: "replay-r4".to_string(),
746                },
747            ],
748        }));
749
750        store.recover_from_journal().await.unwrap();
751
752        // A route removed before the restart is gone from the store and thus
753        // not adoptable — a fresh registration proceeds as first-time.
754        assert!(store.load("replay-r4").await.unwrap().is_none());
755        assert!(!store.take_recovered("replay-r4").await.unwrap());
756    }
757
758    /// Task 3.1 follow-up (holistic review): journal recovery seeds the
759    /// route-state gauges from recovered projections — the replay loop
760    /// writes statuses directly, bypassing the emit sites, so the seed
761    /// after the loop is the only path a recovered route's series has.
762    #[tokio::test]
763    async fn recovery_seeds_route_state_gauges() {
764        struct StateRecorder(std::sync::Mutex<Vec<(String, String)>>);
765        impl MetricsCollector for StateRecorder {
766            fn increment_errors(&self, _r: &str, _e: &str) {}
767            fn record_exchange_duration(&self, _r: &str, _s: std::time::Duration) {}
768            fn set_route_state(&self, route: &str, state: &str) {
769                self.0
770                    .lock()
771                    .expect("states lock") // allow-unwrap
772                    .push((route.to_string(), state.to_string()));
773            }
774            fn increment_exchanges(&self, _r: &str) {}
775            fn set_queue_depth(&self, _q: &str, _d: usize) {}
776            fn record_circuit_breaker_change(&self, _r: &str, _f: &str, _t: &str) {}
777        }
778        let recorder = Arc::new(StateRecorder(std::sync::Mutex::new(Vec::new())));
779        let store = InMemoryRuntimeStore::default()
780            .with_metrics(recorder.clone() as Arc<dyn MetricsCollector>)
781            .with_journal(Arc::new(ReplayJournal {
782                events: vec![
783                    RuntimeEvent::RouteRegistered {
784                        route_id: "seed-r1".to_string(),
785                    },
786                    RuntimeEvent::RouteStarted {
787                        route_id: "seed-r1".to_string(),
788                    },
789                ],
790            }));
791
792        store.recover_from_journal().await.unwrap();
793
794        let states = recorder.0.lock().expect("states lock").clone();
795        assert_eq!(
796            states,
797            vec![("seed-r1".to_string(), "Started".to_string())],
798            "recovery must seed the route-state gauge for each recovered projection"
799        );
800    }
801}