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 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 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 pub fn with_metrics(mut self, metrics: Arc<dyn MetricsCollector>) -> Self {
171 self.metrics = Some(metrics);
172 self
173 }
174
175 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 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 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 guard.recovered = guard.routes.keys().cloned().collect();
497 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 assert!(store.take_recovered("replay-r3").await.unwrap());
731 assert!(!store.take_recovered("replay-r3").await.unwrap());
733 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 assert!(store.load("replay-r4").await.unwrap().is_none());
755 assert!(!store.take_recovered("replay-r4").await.unwrap());
756 }
757
758 #[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") .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}