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
415fn 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 guard.recovered = guard.routes.keys().cloned().collect();
533 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 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 assert!(store.take_recovered("replay-r3").await.unwrap());
773 assert!(!store.take_recovered("replay-r3").await.unwrap());
775 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 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 assert_eq!(derive_boot_nonce(ids), Ok(8));
814 }
815
816 #[test]
817 fn derive_boot_nonce_ignores_non_numeric_tails() {
818 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 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 #[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") .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}