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/// Journal-derived boot nonce from the recorded command IDs of the durable
416/// dedup store: `1 + max(P)` where `P` collects, for every recorded command
417/// ID whose final two `':'`-separated segments both parse as `u64`, its
418/// penultimate segment. Tail scanning (right-to-left via `rsplit`) needs no
419/// prefix or segment-count parsing, so legacy IDs written for routes whose
420/// route IDs contain colons are covered by construction.
421///
422/// No qualifying recorded ID yields `Ok(0)`. A penultimate value of
423/// `u64::MAX` (only reachable via an adversarial recorded command ID)
424/// exhausts the deterministic nonce space and fails closed: the error names
425/// the offending recorded ID and instructs the operator to clean or rotate
426/// the journal.
427fn derive_boot_nonce<'a>(ids: impl IntoIterator<Item = &'a str>) -> Result<u64, DomainError> {
428    let mut max_penultimate: Option<(u64, &str)> = None;
429    for id in ids {
430        let mut segments = id.rsplit(':');
431        let (Some(last), Some(penultimate)) = (segments.next(), segments.next()) else {
432            continue;
433        };
434        let (Ok(_), Ok(nonce_seg)) = (last.parse::<u64>(), penultimate.parse::<u64>()) else {
435            continue;
436        };
437        if max_penultimate.as_ref().is_none_or(|&(m, _)| nonce_seg > m) {
438            max_penultimate = Some((nonce_seg, id));
439        }
440    }
441
442    match max_penultimate {
443        None => Ok(0),
444        Some((u64::MAX, offending)) => Err(DomainError::InvalidState(format!(
445            "deterministic boot nonce space exhausted by recorded command ID '{offending}'; clean or rotate the journal"
446        ))),
447        Some((m, _)) => Ok(m + 1),
448    }
449}
450
451#[async_trait]
452impl RuntimeUnitOfWorkPort for InMemoryRuntimeStore {
453    async fn persist_upsert(
454        &self,
455        aggregate: RouteRuntimeAggregate,
456        expected_version: Option<u64>,
457        projection: RouteStatusProjection,
458        events: &[RuntimeEvent],
459    ) -> Result<(), DomainError> {
460        let mut guard = self.inner.lock().await;
461        if let Some(expected) = expected_version {
462            let route_id = aggregate.route_id().to_string();
463            let current = guard.routes.get(&route_id).ok_or_else(|| {
464                DomainError::InvalidState(format!(
465                    "optimistic lock conflict for route '{route_id}': route not found"
466                ))
467            })?;
468            if current.version() != expected {
469                return Err(DomainError::InvalidState(format!(
470                    "optimistic lock conflict for route '{route_id}': expected version {expected}, actual {}",
471                    current.version()
472                )));
473            }
474        }
475
476        if let Some(journal) = &self.journal {
477            journal.append_batch(events).await?;
478        }
479
480        guard
481            .routes
482            .insert(aggregate.route_id().to_string(), aggregate);
483        let route_id = projection.route_id.clone();
484        let state = projection.status.clone();
485        guard
486            .statuses
487            .insert(projection.route_id.clone(), projection);
488        guard.events.extend(events.iter().cloned());
489        self.emit_route_state(&route_id, &state);
490        Ok(())
491    }
492
493    async fn persist_delete(
494        &self,
495        route_id: &str,
496        events: &[RuntimeEvent],
497    ) -> Result<(), DomainError> {
498        let mut guard = self.inner.lock().await;
499        if let Some(journal) = &self.journal {
500            journal.append_batch(events).await?;
501        }
502        guard.routes.remove(route_id);
503        guard.statuses.remove(route_id);
504        guard.events.extend(events.iter().cloned());
505        if let Some(metrics) = &self.metrics {
506            metrics.clear_route_state(route_id);
507        }
508        Ok(())
509    }
510
511    async fn recover_from_journal(&self) -> Result<(), DomainError> {
512        let Some(journal) = &self.journal else {
513            return Ok(());
514        };
515
516        let replayed_events = journal.load_all().await?;
517        let replayed_command_ids = journal.load_command_ids().await?;
518
519        let mut guard = self.inner.lock().await;
520        guard.routes.clear();
521        guard.statuses.clear();
522        guard.events.clear();
523        guard.seen.clear();
524        guard.recovered.clear();
525
526        for event in &replayed_events {
527            apply_replayed_event(&mut guard, event);
528        }
529        // Every aggregate that survived replay is adoptable exactly once by the
530        // first re-registration of its id (declarative boot), so a durable
531        // journal no longer turns a clean restart into an AlreadyExists storm.
532        guard.recovered = guard.routes.keys().cloned().collect();
533        // Seed the route-state gauges from the recovered projections:
534        // replay writes statuses directly, bypassing the emit sites.
535        if let Some(metrics) = &self.metrics {
536            for status in guard.statuses.values() {
537                metrics.set_route_state(&status.route_id, &status.status);
538            }
539        }
540        guard.events = replayed_events;
541        for command_id in replayed_command_ids {
542            guard.seen.insert(command_id);
543        }
544        Ok(())
545    }
546
547    async fn recovered_boot_nonce(&self) -> Result<u64, DomainError> {
548        let guard = self.inner.lock().await;
549        // After replay, `seen` holds exactly the replayed durable command IDs.
550        derive_boot_nonce(guard.seen.iter().map(String::as_str))
551    }
552}
553
554#[cfg(test)]
555mod tests {
556    use super::*;
557    use std::sync::Arc;
558
559    #[derive(Clone)]
560    struct ReplayJournal {
561        events: Vec<RuntimeEvent>,
562    }
563
564    #[async_trait]
565    impl RuntimeEventJournalPort for ReplayJournal {
566        async fn append_batch(&self, _events: &[RuntimeEvent]) -> Result<(), DomainError> {
567            Ok(())
568        }
569
570        async fn load_all(&self) -> Result<Vec<RuntimeEvent>, DomainError> {
571            Ok(self.events.clone())
572        }
573    }
574
575    #[tokio::test]
576    async fn repo_roundtrip_works() {
577        let repo = InMemoryRouteRepository::default();
578        repo.save(RouteRuntimeAggregate::new("r1")).await.unwrap();
579        assert!(repo.load("r1").await.unwrap().is_some());
580
581        let updated = RouteRuntimeAggregate::from_snapshot(
582            "r1",
583            crate::lifecycle::domain::RouteRuntimeState::Started,
584            1,
585        );
586        repo.save_if_version(updated.clone(), 0).await.unwrap();
587        let loaded = repo.load("r1").await.unwrap().unwrap();
588        assert_eq!(loaded.version(), 1);
589
590        let conflict = repo.save_if_version(updated, 0).await.unwrap_err();
591        assert!(
592            conflict.to_string().contains("optimistic lock conflict"),
593            "unexpected conflict error: {conflict}"
594        );
595
596        repo.delete("r1").await.unwrap();
597        assert!(repo.load("r1").await.unwrap().is_none());
598    }
599
600    #[tokio::test]
601    async fn projection_roundtrip_works() {
602        let store = InMemoryProjectionStore::default();
603        store
604            .upsert_status(RouteStatusProjection {
605                route_id: "r1".into(),
606                status: "Started".into(),
607            })
608            .await
609            .unwrap();
610
611        let status = store.get_status("r1").await.unwrap();
612        assert!(status.is_some());
613        assert_eq!(status.unwrap().status, "Started");
614        store.remove_status("r1").await.unwrap();
615        assert!(store.get_status("r1").await.unwrap().is_none());
616    }
617
618    #[tokio::test]
619    async fn event_publisher_stores_events() {
620        let publisher = InMemoryEventPublisher::default();
621        publisher
622            .publish(&[RuntimeEvent::RouteStarted {
623                route_id: "r1".into(),
624            }])
625            .await
626            .unwrap();
627
628        let events = publisher.snapshot().await;
629        assert_eq!(events.len(), 1);
630    }
631
632    #[tokio::test]
633    async fn command_dedup_detects_duplicates() {
634        let dedup = InMemoryCommandDedup::default();
635        assert!(dedup.first_seen("c1").await.unwrap());
636        assert!(!dedup.first_seen("c1").await.unwrap());
637        dedup.forget_seen("c1").await.unwrap();
638        assert!(dedup.first_seen("c1").await.unwrap());
639        assert!(dedup.first_seen("c2").await.unwrap());
640    }
641
642    #[tokio::test]
643    async fn runtime_store_uow_persists_all_three_writes() {
644        let store = InMemoryRuntimeStore::default();
645        let aggregate = RouteRuntimeAggregate::new("uow-r1");
646        let projection = RouteStatusProjection {
647            route_id: "uow-r1".to_string(),
648            status: "Registered".to_string(),
649        };
650        let events = vec![RuntimeEvent::RouteRegistered {
651            route_id: "uow-r1".to_string(),
652        }];
653
654        store
655            .persist_upsert(aggregate, None, projection.clone(), &events)
656            .await
657            .unwrap();
658
659        assert!(store.load("uow-r1").await.unwrap().is_some());
660        assert_eq!(
661            store.get_status("uow-r1").await.unwrap().unwrap(),
662            projection
663        );
664        assert_eq!(store.snapshot_events().await, events);
665    }
666
667    #[tokio::test]
668    async fn runtime_store_uow_enforces_expected_version() {
669        let store = InMemoryRuntimeStore::default();
670        let initial = RouteRuntimeAggregate::new("uow-r2");
671        let initial_projection = RouteStatusProjection {
672            route_id: "uow-r2".to_string(),
673            status: "Registered".to_string(),
674        };
675        store
676            .persist_upsert(
677                initial,
678                None,
679                initial_projection,
680                &[RuntimeEvent::RouteRegistered {
681                    route_id: "uow-r2".to_string(),
682                }],
683            )
684            .await
685            .unwrap();
686
687        let started = RouteRuntimeAggregate::from_snapshot(
688            "uow-r2",
689            crate::lifecycle::domain::RouteRuntimeState::Started,
690            1,
691        );
692        let err = store
693            .persist_upsert(
694                started,
695                Some(99),
696                RouteStatusProjection {
697                    route_id: "uow-r2".to_string(),
698                    status: "Started".to_string(),
699                },
700                &[RuntimeEvent::RouteStarted {
701                    route_id: "uow-r2".to_string(),
702                }],
703            )
704            .await
705            .unwrap_err()
706            .to_string();
707        assert!(
708            err.contains("optimistic lock conflict"),
709            "unexpected error: {err}"
710        );
711    }
712
713    #[tokio::test]
714    async fn replay_start_requested_only_advances_version_once() {
715        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
716            events: vec![
717                RuntimeEvent::RouteRegistered {
718                    route_id: "replay-r1".to_string(),
719                },
720                RuntimeEvent::RouteStartRequested {
721                    route_id: "replay-r1".to_string(),
722                },
723            ],
724        }));
725
726        store.recover_from_journal().await.unwrap();
727        let aggregate = store.load("replay-r1").await.unwrap().unwrap();
728
729        assert_eq!(aggregate.state(), &RouteRuntimeState::Starting);
730        assert_eq!(aggregate.version(), 1);
731    }
732
733    #[tokio::test]
734    async fn replay_start_requested_then_started_keeps_single_command_version() {
735        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
736            events: vec![
737                RuntimeEvent::RouteRegistered {
738                    route_id: "replay-r2".to_string(),
739                },
740                RuntimeEvent::RouteStartRequested {
741                    route_id: "replay-r2".to_string(),
742                },
743                RuntimeEvent::RouteStarted {
744                    route_id: "replay-r2".to_string(),
745                },
746            ],
747        }));
748
749        store.recover_from_journal().await.unwrap();
750        let aggregate = store.load("replay-r2").await.unwrap().unwrap();
751
752        assert_eq!(aggregate.state(), &RouteRuntimeState::Started);
753        assert_eq!(aggregate.version(), 1);
754    }
755
756    #[tokio::test]
757    async fn recover_from_journal_marks_replayed_routes_adoptable_once() {
758        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
759            events: vec![
760                RuntimeEvent::RouteRegistered {
761                    route_id: "replay-r3".to_string(),
762                },
763                RuntimeEvent::RouteStarted {
764                    route_id: "replay-r3".to_string(),
765                },
766            ],
767        }));
768
769        store.recover_from_journal().await.unwrap();
770
771        // First re-registration adopts the recovered aggregate.
772        assert!(store.take_recovered("replay-r3").await.unwrap());
773        // Marker is consumed: a same-session duplicate is no longer adoptable.
774        assert!(!store.take_recovered("replay-r3").await.unwrap());
775        // A route that never existed in the journal is never adoptable.
776        assert!(!store.take_recovered("never-seen").await.unwrap());
777    }
778
779    #[tokio::test]
780    async fn recover_from_journal_leaves_no_marker_for_removed_routes() {
781        let store = InMemoryRuntimeStore::default().with_journal(Arc::new(ReplayJournal {
782            events: vec![
783                RuntimeEvent::RouteRegistered {
784                    route_id: "replay-r4".to_string(),
785                },
786                RuntimeEvent::RouteRemoved {
787                    route_id: "replay-r4".to_string(),
788                },
789            ],
790        }));
791
792        store.recover_from_journal().await.unwrap();
793
794        // A route removed before the restart is gone from the store and thus
795        // not adoptable — a fresh registration proceeds as first-time.
796        assert!(store.load("replay-r4").await.unwrap().is_none());
797        assert!(!store.take_recovered("replay-r4").await.unwrap());
798    }
799
800    #[test]
801    fn derive_boot_nonce_empty_is_zero() {
802        assert_eq!(derive_boot_nonce([]), Ok(0));
803    }
804
805    #[test]
806    fn derive_boot_nonce_strictly_above_recorded_penultimates() {
807        let ids = [
808            "context:start:r:7:0",
809            "context:start:r:3:5",
810            "context:stop:r:7:1",
811        ];
812        // 1 + max penultimate 7
813        assert_eq!(derive_boot_nonce(ids), Ok(8));
814    }
815
816    #[test]
817    fn derive_boot_nonce_ignores_non_numeric_tails() {
818        // Legacy four-segment IDs whose final-two rule fails: `hello`/`0`
819        // (`hello` does not parse) and a lone trailing segment.
820        let ids = ["context:start:hello:0", "context:start:foo"];
821        assert_eq!(derive_boot_nonce(ids), Ok(0));
822    }
823
824    #[test]
825    fn derive_boot_nonce_legacy_colon_route_is_forbidden() {
826        // Legacy write for route `foo:0`: the tail scan collects penultimate
827        // `0`, so the derived nonce must be strictly above it — never `Ok(0)`.
828        assert_eq!(derive_boot_nonce(["context:start:foo:0:0"]), Ok(1));
829    }
830
831    #[test]
832    fn derive_boot_nonce_max_penultimate_fails_closed() {
833        let offending = "context:start:r:18446744073709551615:0";
834        let err = derive_boot_nonce([offending]).unwrap_err();
835        let msg = err.to_string();
836        assert!(
837            msg.contains(offending),
838            "error must name the offending recorded ID verbatim: {msg}"
839        );
840        assert!(
841            msg.contains("clean or rotate the journal"),
842            "error must instruct the operator to clean or rotate the journal: {msg}"
843        );
844    }
845
846    /// Task 3.1 follow-up (holistic review): journal recovery seeds the
847    /// route-state gauges from recovered projections — the replay loop
848    /// writes statuses directly, bypassing the emit sites, so the seed
849    /// after the loop is the only path a recovered route's series has.
850    #[tokio::test]
851    async fn recovery_seeds_route_state_gauges() {
852        struct StateRecorder(std::sync::Mutex<Vec<(String, String)>>);
853        impl MetricsCollector for StateRecorder {
854            fn increment_errors(&self, _r: &str, _e: &str) {}
855            fn record_exchange_duration(&self, _r: &str, _s: std::time::Duration) {}
856            fn set_route_state(&self, route: &str, state: &str) {
857                self.0
858                    .lock()
859                    .expect("states lock") // allow-unwrap
860                    .push((route.to_string(), state.to_string()));
861            }
862            fn increment_exchanges(&self, _r: &str) {}
863            fn set_queue_depth(&self, _q: &str, _d: usize) {}
864            fn record_circuit_breaker_change(&self, _r: &str, _f: &str, _t: &str) {}
865        }
866        let recorder = Arc::new(StateRecorder(std::sync::Mutex::new(Vec::new())));
867        let store = InMemoryRuntimeStore::default()
868            .with_metrics(recorder.clone() as Arc<dyn MetricsCollector>)
869            .with_journal(Arc::new(ReplayJournal {
870                events: vec![
871                    RuntimeEvent::RouteRegistered {
872                        route_id: "seed-r1".to_string(),
873                    },
874                    RuntimeEvent::RouteStarted {
875                        route_id: "seed-r1".to_string(),
876                    },
877                ],
878            }));
879
880        store.recover_from_journal().await.unwrap();
881
882        let states = recorder.0.lock().expect("states lock").clone();
883        assert_eq!(
884            states,
885            vec![("seed-r1".to_string(), "Started".to_string())],
886            "recovery must seed the route-state gauge for each recovered projection"
887        );
888    }
889}