Skip to main content

prns_runtime_embassy/runtime/
embedded_persistence.rs

1use embedded_storage_async::nor_flash::NorFlash;
2use heapless::Vec as HeaplessVec;
3
4use crate::crypto::ratchets::SeedSelfRatchetsOutcome;
5use crate::engine::{EngineState, InstantMillis, Journaled, RouteSeedOutcome};
6use crate::identity::Zeroizing;
7use crate::interfaces::AttachedInterfaces;
8use crate::persistence::{
9    maximum_route_upsert_payload_len, read_routing_table_snapshot, read_self_ratchets_snapshot,
10    routing_table_snapshot_len, self_ratchets_snapshot_len, write_routing_table_snapshot,
11    write_self_ratchets_snapshot, FlashJournal, FlashJournalError, FlashJournalLayout,
12    FlashJournalRecord, FlashJournalRecordKind, FlashJournalWarning,
13    TIMEBASE_RECORD_INTERVAL_MILLIS,
14};
15use crate::routing::announce::emit::MAX_ANNOUNCE_APP_DATA_LEN;
16use crate::routing::AnnounceIdRing;
17use crate::storage::StorageLayout;
18use crate::wire::{DestinationHash, TRUNCATED_HASH_BYTE_LEN};
19
20const RECORD_SCRATCH_LEN: usize =
21    (maximum_route_upsert_payload_len(MAX_ANNOUNCE_APP_DATA_LEN, 0) + 3) & !3;
22const HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS: u64 = 24 * 60 * 60 * 1_000;
23
24#[derive(Debug, Clone, Copy, PartialEq, Eq)]
25pub struct EmbeddedCompactionPolicy {
26    minimum_interval_millis: u64,
27    critical_reserve_bytes: usize,
28}
29
30impl EmbeddedCompactionPolicy {
31    #[must_use]
32    pub const fn new(minimum_interval_millis: u64, critical_reserve_bytes: usize) -> Self {
33        Self {
34            minimum_interval_millis,
35            critical_reserve_bytes,
36        }
37    }
38
39    #[must_use]
40    pub const fn hopspot(critical_reserve_bytes: usize) -> Self {
41        Self::new(
42            HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS,
43            critical_reserve_bytes,
44        )
45    }
46}
47
48#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub struct EmbeddedPersistencePolicy {
50    first_route_commit_delay_millis: u64,
51    minimum_route_commit_interval_millis: u64,
52    ratchet_batch_delay_millis: u64,
53    retry_interval_millis: u64,
54    timebase_record_interval_millis: u64,
55    compaction: EmbeddedCompactionPolicy,
56}
57
58impl EmbeddedPersistencePolicy {
59    #[must_use]
60    pub const fn new(
61        first_route_commit_delay_millis: u64,
62        minimum_route_commit_interval_millis: u64,
63        ratchet_batch_delay_millis: u64,
64        retry_interval_millis: u64,
65        timebase_record_interval_millis: u64,
66        compaction: EmbeddedCompactionPolicy,
67    ) -> Self {
68        Self {
69            first_route_commit_delay_millis,
70            minimum_route_commit_interval_millis,
71            ratchet_batch_delay_millis,
72            retry_interval_millis,
73            timebase_record_interval_millis,
74            compaction,
75        }
76    }
77
78    #[must_use]
79    pub const fn hopspot_default(compaction: EmbeddedCompactionPolicy) -> Self {
80        Self::new(
81            2_000,
82            5 * 60 * 1_000,
83            2_000,
84            5 * 60 * 1_000,
85            TIMEBASE_RECORD_INTERVAL_MILLIS,
86            compaction,
87        )
88    }
89}
90
91#[derive(Debug, Clone, Copy, PartialEq, Eq)]
92pub enum RouteSnapshotKeyError {
93    Capacity,
94}
95
96pub trait RouteSnapshotKeys {
97    fn clear(&mut self);
98    fn push(&mut self, destination: DestinationHash) -> Result<(), RouteSnapshotKeyError>;
99    fn get(&self, index: usize) -> Option<DestinationHash>;
100}
101
102pub struct FixedRouteSnapshotKeys<const N: usize> {
103    keys: HeaplessVec<DestinationHash, N>,
104}
105
106impl<const N: usize> FixedRouteSnapshotKeys<N> {
107    #[must_use]
108    pub const fn new() -> Self {
109        Self {
110            keys: HeaplessVec::new(),
111        }
112    }
113}
114
115impl<const N: usize> Default for FixedRouteSnapshotKeys<N> {
116    fn default() -> Self {
117        Self::new()
118    }
119}
120
121impl<const N: usize> RouteSnapshotKeys for FixedRouteSnapshotKeys<N> {
122    fn clear(&mut self) {
123        self.keys.clear();
124    }
125
126    fn push(&mut self, destination: DestinationHash) -> Result<(), RouteSnapshotKeyError> {
127        self.keys
128            .push(destination)
129            .map_err(|_| RouteSnapshotKeyError::Capacity)
130    }
131
132    fn get(&self, index: usize) -> Option<DestinationHash> {
133        self.keys.get(index).copied()
134    }
135}
136
137#[derive(Debug, Clone, Copy, PartialEq, Eq)]
138pub enum EmbeddedPersistenceFailure {
139    Flash,
140    Codec,
141    Capacity,
142}
143
144#[derive(Debug, Clone, Copy, PartialEq, Eq)]
145pub enum EmbeddedPersistenceTarget {
146    Routes,
147    CriticalState,
148}
149
150#[derive(Debug, Clone, Copy, PartialEq, Eq)]
151pub struct EmbeddedPersistenceRestoreReport {
152    pub logical_start: InstantMillis,
153    pub route_seeded_count: u32,
154    pub route_refused_count: u32,
155    pub route_dropped_count: u32,
156    pub ratchet_seeded_count: u32,
157    pub ratchet_refused_count: u32,
158    pub warning: Option<FlashJournalWarning>,
159}
160
161#[derive(Debug, Clone, Copy, PartialEq, Eq)]
162pub enum EmbeddedPersistenceDiagnostic {
163    Restored(EmbeddedPersistenceRestoreReport),
164    BatchPersisted {
165        records: u32,
166        at: InstantMillis,
167        state_not_saved: bool,
168    },
169    CompactionStarted {
170        at: InstantMillis,
171        next_allowed_at: InstantMillis,
172    },
173    CompactionCompleted {
174        records: u32,
175        at: InstantMillis,
176        state_not_saved: bool,
177    },
178    DurabilityDeferred {
179        target: EmbeddedPersistenceTarget,
180        until: InstantMillis,
181    },
182    WriteFailed {
183        failure: EmbeddedPersistenceFailure,
184        retry_at: InstantMillis,
185    },
186}
187
188#[derive(Debug, Clone, Copy, PartialEq, Eq)]
189enum PendingRouteDelta {
190    RouteUpsert(DestinationHash),
191    RouteRemoval(DestinationHash),
192}
193
194impl PendingRouteDelta {
195    fn destination(self) -> DestinationHash {
196        match self {
197            Self::RouteUpsert(destination) | Self::RouteRemoval(destination) => destination,
198        }
199    }
200}
201
202#[derive(Debug, Clone, Copy, PartialEq, Eq)]
203enum BatchKind {
204    Routes,
205    Ratchets,
206    Compaction,
207}
208
209#[derive(Debug, Clone, Copy, PartialEq, Eq)]
210enum CompactionPhase {
211    RecordBudget { at: InstantMillis },
212    Erase { sector: usize },
213    Routes { index: usize },
214    Ratchets { index: usize },
215    Commit,
216}
217
218struct EncodedDelta {
219    kind: FlashJournalRecordKind,
220    payload: Zeroizing<[u8; RECORD_SCRATCH_LEN]>,
221    len: usize,
222}
223
224pub struct EmbeddedFlashPersistence<F, Keys, Observe, const PENDING: usize>
225where
226    F: NorFlash,
227    Keys: RouteSnapshotKeys,
228    Observe: FnMut(EmbeddedPersistenceDiagnostic),
229{
230    flash: Option<F>,
231    journal: Option<FlashJournal<F>>,
232    layout: FlashJournalLayout,
233    policy: EmbeddedPersistencePolicy,
234    observe_diagnostic: Observe,
235    pending_routes: HeaplessVec<PendingRouteDelta, PENDING>,
236    pending_ratchets: HeaplessVec<DestinationHash, PENDING>,
237    compaction_route_keys: Keys,
238    compaction_ratchet_keys: HeaplessVec<DestinationHash, PENDING>,
239    route_dirty_since: Option<InstantMillis>,
240    ratchet_dirty_since: Option<InstantMillis>,
241    last_route_success: Option<InstantMillis>,
242    last_timebase_success: Option<InstantMillis>,
243    retry_not_before: Option<InstantMillis>,
244    landing_batch: Option<BatchKind>,
245    landing_records: u32,
246    compaction: Option<CompactionPhase>,
247    compaction_target: Option<EmbeddedPersistenceTarget>,
248    snapshot_required: bool,
249    snapshot_target: EmbeddedPersistenceTarget,
250    next_compaction_not_before: Option<InstantMillis>,
251    deferred_target: Option<EmbeddedPersistenceTarget>,
252    deferred_until: Option<InstantMillis>,
253    write_failed: bool,
254}
255
256impl<F, Keys, Observe, const PENDING: usize> EmbeddedFlashPersistence<F, Keys, Observe, PENDING>
257where
258    F: NorFlash,
259    Keys: RouteSnapshotKeys,
260    Observe: FnMut(EmbeddedPersistenceDiagnostic),
261{
262    #[must_use]
263    pub fn new(
264        flash: F,
265        layout: FlashJournalLayout,
266        policy: EmbeddedPersistencePolicy,
267        compaction_route_keys: Keys,
268        observe_diagnostic: Observe,
269    ) -> Self {
270        Self {
271            flash: Some(flash),
272            journal: None,
273            layout,
274            policy,
275            observe_diagnostic,
276            pending_routes: HeaplessVec::new(),
277            pending_ratchets: HeaplessVec::new(),
278            compaction_route_keys,
279            compaction_ratchet_keys: HeaplessVec::new(),
280            route_dirty_since: None,
281            ratchet_dirty_since: None,
282            last_route_success: None,
283            last_timebase_success: None,
284            retry_not_before: None,
285            landing_batch: None,
286            landing_records: 0,
287            compaction: None,
288            compaction_target: None,
289            snapshot_required: false,
290            snapshot_target: EmbeddedPersistenceTarget::Routes,
291            next_compaction_not_before: None,
292            deferred_target: None,
293            deferred_until: None,
294            write_failed: false,
295        }
296    }
297
298    #[must_use]
299    pub fn state_not_saved(&self) -> bool {
300        self.write_failed || self.deferred_target.is_some()
301    }
302
303    pub async fn restore<S: StorageLayout>(
304        &mut self,
305        engine: &mut EngineState<S>,
306        raw_now: InstantMillis,
307    ) -> EmbeddedPersistenceRestoreReport {
308        let Some(mut flash) = self.flash.take() else {
309            return self.empty_restore_report(raw_now, Some(FlashJournalWarning::Corrupt));
310        };
311        let timebase_state = FlashJournal::inspect_timebase_state(&mut flash, self.layout)
312            .await
313            .ok();
314        let logical_start = timebase_state
315            .and_then(|state| state.high_water)
316            .map_or(raw_now, |high_water| high_water.max(raw_now));
317        let mut scratch = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
318        let mut report = EmbeddedPersistenceRestoreReport {
319            logical_start,
320            route_seeded_count: 0,
321            route_refused_count: 0,
322            route_dropped_count: 0,
323            ratchet_seeded_count: 0,
324            ratchet_refused_count: 0,
325            warning: None,
326        };
327        let opened = FlashJournal::open(flash, self.layout, &mut scratch[..], |record| {
328            apply_record(engine, logical_start, record, &mut report)
329        })
330        .await;
331        let Ok((mut journal, restored)) = opened else {
332            report.warning = Some(FlashJournalWarning::Corrupt);
333            (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::Restored(report));
334            return report;
335        };
336        report.warning = restored.warning;
337        let initialization_failed =
338            restored.active_epoch.is_none() && journal.initialize_empty().await.is_err();
339        self.next_compaction_not_before = timebase_state
340            .and_then(|state| state.last_compaction_attempt)
341            .map(|attempt| {
342                InstantMillis(
343                    attempt
344                        .0
345                        .saturating_add(self.policy.compaction.minimum_interval_millis),
346                )
347            });
348        self.journal = Some(journal);
349        (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::Restored(report));
350        if initialization_failed {
351            self.note_write_failure(raw_now, EmbeddedPersistenceFailure::Flash);
352        }
353        report
354    }
355
356    fn empty_restore_report(
357        &mut self,
358        logical_start: InstantMillis,
359        warning: Option<FlashJournalWarning>,
360    ) -> EmbeddedPersistenceRestoreReport {
361        let report = EmbeddedPersistenceRestoreReport {
362            logical_start,
363            route_seeded_count: 0,
364            route_refused_count: 0,
365            route_dropped_count: 0,
366            ratchet_seeded_count: 0,
367            ratchet_refused_count: 0,
368            warning,
369        };
370        (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::Restored(report));
371        report
372    }
373
374    fn observe_journaled(&mut self, journaled: &Journaled<'_>, now: InstantMillis) {
375        match journaled {
376            Journaled::AnnounceHeard { observation, .. } => {
377                self.queue_route(PendingRouteDelta::RouteUpsert(observation.destination), now);
378                if self.route_dirty_since.is_none() {
379                    self.route_dirty_since = Some(now);
380                }
381            }
382            Journaled::RouteRemoved { destination, .. } => {
383                self.queue_route(PendingRouteDelta::RouteRemoval(*destination), now);
384                if self.route_dirty_since.is_none() {
385                    self.route_dirty_since = Some(now);
386                }
387            }
388            Journaled::SelfRatchetRotated { destination } => {
389                self.queue_ratchet(*destination, now);
390                if self.ratchet_dirty_since.is_none() {
391                    self.ratchet_dirty_since = Some(now);
392                }
393            }
394            Journaled::AnnounceHeldDropped { .. }
395            | Journaled::Delivered(_)
396            | Journaled::CommandSettled { .. }
397            | Journaled::PersistenceFlushed { .. }
398            | Journaled::PersistenceFlushFailed { .. }
399            | Journaled::LinkEstablished(_)
400            | Journaled::PeerIdentified { .. }
401            | Journaled::RequestReceived { .. }
402            | Journaled::ResponseReceived { .. }
403            | Journaled::ResponseSegmentReceived { .. }
404            | Journaled::ChannelMessageReceived { .. }
405            | Journaled::LinkClosed { .. }
406            | Journaled::LinkInterfaceMismatch { .. }
407            | Journaled::ResourceReceived { .. }
408            | Journaled::ResourceFailed { .. }
409            | Journaled::ResourceNeedsDecompression { .. }
410            | Journaled::ResourceSegmentReceived { .. }
411            | Journaled::ResourceAssembled { .. } => {}
412        }
413    }
414
415    fn queue_route(&mut self, delta: PendingRouteDelta, now: InstantMillis) {
416        if self.snapshot_required {
417            return;
418        }
419        if let Some(existing) = self
420            .pending_routes
421            .iter_mut()
422            .find(|pending| pending.destination() == delta.destination())
423        {
424            *existing = delta;
425            return;
426        }
427        if self.pending_routes.push(delta).is_err() {
428            self.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
429        }
430    }
431
432    fn queue_ratchet(&mut self, destination: DestinationHash, now: InstantMillis) {
433        if self.pending_ratchets.contains(&destination) {
434            return;
435        }
436        if self.pending_ratchets.push(destination).is_err() {
437            self.require_snapshot(EmbeddedPersistenceTarget::CriticalState, now);
438        }
439    }
440
441    fn require_snapshot(&mut self, target: EmbeddedPersistenceTarget, now: InstantMillis) {
442        self.snapshot_required = true;
443        self.snapshot_target = match (self.snapshot_target, target) {
444            (EmbeddedPersistenceTarget::CriticalState, _)
445            | (_, EmbeddedPersistenceTarget::CriticalState) => {
446                EmbeddedPersistenceTarget::CriticalState
447            }
448            (EmbeddedPersistenceTarget::Routes, EmbeddedPersistenceTarget::Routes) => {
449                EmbeddedPersistenceTarget::Routes
450            }
451        };
452        self.pending_routes.clear();
453        if target == EmbeddedPersistenceTarget::CriticalState {
454            self.pending_ratchets.clear();
455        }
456        let Some(until) = self.next_compaction_not_before else {
457            return;
458        };
459        if now.0 >= until.0 {
460            return;
461        }
462        if self.deferred_target == Some(self.snapshot_target) && self.deferred_until == Some(until)
463        {
464            return;
465        }
466        self.deferred_target = Some(self.snapshot_target);
467        self.deferred_until = Some(until);
468        (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::DurabilityDeferred {
469            target: self.snapshot_target,
470            until,
471        });
472    }
473
474    fn next_deadline(&self, now: InstantMillis) -> Option<InstantMillis> {
475        self.journal.as_ref()?;
476        let mut deadline = if self.compaction.is_some() || self.landing_batch.is_some() {
477            Some(now)
478        } else {
479            None
480        };
481        let ratchet_ready = self.ratchet_ready_at();
482        let route_ready = self.route_ready_at();
483        deadline = earlier(deadline, ratchet_ready);
484        if self.snapshot_required {
485            let requested = match self.snapshot_target {
486                EmbeddedPersistenceTarget::Routes => route_ready,
487                EmbeddedPersistenceTarget::CriticalState => {
488                    earlier(route_ready, ratchet_ready).or(Some(now))
489                }
490            };
491            if let Some(requested) = requested {
492                let allowed = self.next_compaction_not_before.unwrap_or(requested);
493                deadline = earlier(deadline, Some(InstantMillis(requested.0.max(allowed.0))));
494            }
495        } else {
496            deadline = earlier(deadline, route_ready);
497        }
498        let timebase_ready = self.last_timebase_success.map_or(now, |last| {
499            InstantMillis(
500                last.0
501                    .saturating_add(self.policy.timebase_record_interval_millis),
502            )
503        });
504        deadline = earlier(deadline, Some(timebase_ready));
505        match (deadline, self.retry_not_before) {
506            (Some(deadline), Some(retry)) => Some(InstantMillis(deadline.0.max(retry.0))),
507            (deadline, None) => deadline,
508            (None, Some(_)) => None,
509        }
510    }
511
512    fn ratchet_ready_at(&self) -> Option<InstantMillis> {
513        self.ratchet_dirty_since.map(|dirty| {
514            InstantMillis(
515                dirty
516                    .0
517                    .saturating_add(self.policy.ratchet_batch_delay_millis),
518            )
519        })
520    }
521
522    fn route_ready_at(&self) -> Option<InstantMillis> {
523        self.route_dirty_since.map(|dirty| {
524            let first_ready = dirty
525                .0
526                .saturating_add(self.policy.first_route_commit_delay_millis);
527            let interval_ready = self.last_route_success.map_or(0, |last| {
528                last.0
529                    .saturating_add(self.policy.minimum_route_commit_interval_millis)
530            });
531            InstantMillis(first_ready.max(interval_ready))
532        })
533    }
534
535    async fn progress<S: StorageLayout>(
536        &mut self,
537        engine: &mut EngineState<S>,
538        now: InstantMillis,
539    ) {
540        if self.retry_not_before.is_some_and(|retry| now.0 < retry.0) {
541            return;
542        }
543        if self.compaction.is_some() {
544            self.progress_compaction(engine, now).await;
545            return;
546        }
547        if let Some(batch) = self.landing_batch {
548            let new_work = match batch {
549                BatchKind::Routes => !self.pending_routes.is_empty() || self.snapshot_required,
550                BatchKind::Ratchets => !self.pending_ratchets.is_empty(),
551                BatchKind::Compaction => {
552                    !self.pending_routes.is_empty()
553                        || !self.pending_ratchets.is_empty()
554                        || self.snapshot_required
555                }
556            };
557            if new_work {
558                self.landing_batch = None;
559            } else {
560                self.land_timebase(batch, now).await;
561                return;
562            }
563        }
564        let ratchet_due = self
565            .ratchet_ready_at()
566            .is_some_and(|ready| now.0 >= ready.0);
567        let route_due = self.route_ready_at().is_some_and(|ready| now.0 >= ready.0);
568        if ratchet_due {
569            if let Some(index) = (!self.pending_ratchets.is_empty()).then_some(0) {
570                self.append_ratchet(engine, index, now).await;
571                return;
572            }
573            if self.snapshot_required
574                && self.snapshot_target == EmbeddedPersistenceTarget::CriticalState
575            {
576                self.try_start_compaction(engine, now);
577                return;
578            }
579        }
580        if !route_due {
581            let timebase_due = self.last_timebase_success.is_none_or(|last| {
582                now.0.saturating_sub(last.0) >= self.policy.timebase_record_interval_millis
583            });
584            if timebase_due {
585                self.record_timebase(now).await;
586            }
587            return;
588        }
589        if self.snapshot_required {
590            self.try_start_compaction(engine, now);
591            return;
592        }
593        if !self.pending_routes.is_empty() {
594            self.append_route(engine, 0, now).await;
595        }
596    }
597
598    async fn append_route<S: StorageLayout>(
599        &mut self,
600        engine: &EngineState<S>,
601        index: usize,
602        now: InstantMillis,
603    ) {
604        let delta = self.pending_routes[index];
605        let encoded = encode_route_delta(engine, delta);
606        let Ok(payload) = encoded else {
607            self.note_codec_failure(now);
608            return;
609        };
610        let can_fit = self.journal.as_ref().is_some_and(|journal| {
611            journal.active_can_fit(payload.len, self.policy.compaction.critical_reserve_bytes)
612        });
613        if !can_fit {
614            self.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
615            self.try_start_compaction(engine, now);
616            return;
617        }
618        let Some(journal) = self.journal.as_mut() else {
619            return;
620        };
621        let result = journal
622            .append(payload.kind, &payload.payload[..payload.len])
623            .await;
624        match result {
625            Ok(()) => {
626                self.pending_routes.swap_remove(index);
627                self.landing_records = self.landing_records.saturating_add(1);
628                if self.pending_routes.is_empty() {
629                    self.landing_batch = Some(BatchKind::Routes);
630                }
631            }
632            Err(FlashJournalError::ArenaFull) => {
633                self.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
634                self.try_start_compaction(engine, now);
635            }
636            Err(error) => {
637                let failure = failure_from_journal(error);
638                self.note_write_failure(now, failure);
639            }
640        }
641    }
642
643    async fn append_ratchet<S: StorageLayout>(
644        &mut self,
645        engine: &EngineState<S>,
646        index: usize,
647        now: InstantMillis,
648    ) {
649        let destination = self.pending_ratchets[index];
650        let Ok(payload) = encode_ratchet(engine, destination) else {
651            self.note_codec_failure(now);
652            return;
653        };
654        let Some(journal) = self.journal.as_mut() else {
655            return;
656        };
657        let result = journal
658            .append(payload.kind, &payload.payload[..payload.len])
659            .await;
660        match result {
661            Ok(()) => {
662                self.pending_ratchets.swap_remove(index);
663                self.landing_records = self.landing_records.saturating_add(1);
664                if self.pending_ratchets.is_empty() {
665                    self.landing_batch = Some(BatchKind::Ratchets);
666                }
667            }
668            Err(FlashJournalError::ArenaFull) => {
669                self.require_snapshot(EmbeddedPersistenceTarget::CriticalState, now);
670                self.try_start_compaction(engine, now);
671            }
672            Err(error) => self.note_write_failure(now, failure_from_journal(error)),
673        }
674    }
675
676    fn try_start_compaction<S: StorageLayout>(
677        &mut self,
678        engine: &EngineState<S>,
679        now: InstantMillis,
680    ) {
681        let allowed = self.next_compaction_not_before.unwrap_or(now);
682        if now.0 < allowed.0 {
683            self.require_snapshot(self.snapshot_target, now);
684            return;
685        }
686        self.compaction_route_keys.clear();
687        let mut route_capacity_failed = false;
688        for destination in engine.persisted_route_destinations() {
689            if self.compaction_route_keys.push(destination).is_err() {
690                route_capacity_failed = true;
691                break;
692            }
693        }
694        if route_capacity_failed {
695            self.note_write_failure(now, EmbeddedPersistenceFailure::Capacity);
696            return;
697        }
698        self.compaction_ratchet_keys.clear();
699        let mut ratchet_capacity_failed = false;
700        for (destination, _, _) in engine.persisted_self_ratchet_rows() {
701            if self.compaction_ratchet_keys.push(destination).is_err() {
702                ratchet_capacity_failed = true;
703                break;
704            }
705        }
706        if ratchet_capacity_failed {
707            self.note_write_failure(now, EmbeddedPersistenceFailure::Capacity);
708            return;
709        }
710        let target = self.snapshot_target;
711        self.pending_routes.clear();
712        self.pending_ratchets.clear();
713        self.route_dirty_since = None;
714        self.ratchet_dirty_since = None;
715        self.snapshot_required = false;
716        self.snapshot_target = EmbeddedPersistenceTarget::Routes;
717        self.compaction_target = Some(target);
718        self.compaction = Some(CompactionPhase::RecordBudget { at: now });
719        self.landing_batch = None;
720        self.landing_records = 0;
721    }
722
723    async fn progress_compaction<S: StorageLayout>(
724        &mut self,
725        engine: &mut EngineState<S>,
726        now: InstantMillis,
727    ) {
728        let Some(phase) = self.compaction else {
729            return;
730        };
731        let Some(journal) = self.journal.as_mut() else {
732            return;
733        };
734        match phase {
735            CompactionPhase::RecordBudget { at } => {
736                match journal.record_compaction_budget(at).await {
737                    Ok(recorded_at) => {
738                        self.last_timebase_success = Some(at);
739                        let next_allowed_at = InstantMillis(
740                            recorded_at
741                                .0
742                                .saturating_add(self.policy.compaction.minimum_interval_millis),
743                        );
744                        self.next_compaction_not_before = Some(next_allowed_at);
745                        (self.observe_diagnostic)(
746                            EmbeddedPersistenceDiagnostic::CompactionStarted {
747                                at,
748                                next_allowed_at,
749                            },
750                        );
751                        self.compaction = Some(CompactionPhase::Erase { sector: 0 });
752                    }
753                    Err(error) => {
754                        self.note_write_failure(now, failure_from_journal(error));
755                    }
756                }
757            }
758            CompactionPhase::Erase { sector } => {
759                if sector < journal.inactive_sector_count() {
760                    match journal.erase_inactive_sector(sector).await {
761                        Ok(()) => {
762                            let next = sector + 1;
763                            if next == journal.inactive_sector_count() {
764                                if journal.begin_compaction().is_err() {
765                                    self.note_write_failure(
766                                        now,
767                                        EmbeddedPersistenceFailure::Capacity,
768                                    );
769                                    return;
770                                }
771                                self.compaction = Some(CompactionPhase::Routes { index: 0 });
772                            } else {
773                                self.compaction = Some(CompactionPhase::Erase { sector: next });
774                            }
775                        }
776                        Err(error) => {
777                            self.note_write_failure(now, failure_from_journal(error));
778                        }
779                    }
780                }
781            }
782            CompactionPhase::Routes { index } => {
783                let mut scratch = [0u8; RECORD_SCRATCH_LEN];
784                let Some(destination) = self.compaction_route_keys.get(index) else {
785                    self.compaction = Some(CompactionPhase::Ratchets { index: 0 });
786                    return;
787                };
788                let Some(row) = engine.persisted_route_row(&destination) else {
789                    self.compaction = Some(CompactionPhase::Routes { index: index + 1 });
790                    return;
791                };
792                let mut durable = row.clone();
793                durable.announce_id_ring = AnnounceIdRing::Table(&[]);
794                let required = routing_table_snapshot_len(core::iter::once(durable.clone()));
795                if required > scratch.len() {
796                    self.note_codec_failure(now);
797                    return;
798                }
799                let Ok(written) = write_routing_table_snapshot(
800                    core::iter::once(durable),
801                    &mut scratch[..required],
802                ) else {
803                    self.note_codec_failure(now);
804                    return;
805                };
806                match journal
807                    .append_compacted(FlashJournalRecordKind::RouteUpsert, &scratch[..written])
808                    .await
809                {
810                    Ok(()) => {
811                        self.landing_records = self.landing_records.saturating_add(1);
812                        self.compaction = Some(CompactionPhase::Routes { index: index + 1 });
813                    }
814                    Err(error) => {
815                        self.note_write_failure(now, failure_from_journal(error));
816                    }
817                }
818            }
819            CompactionPhase::Ratchets { index } => {
820                let Some(destination) = self.compaction_ratchet_keys.get(index).copied() else {
821                    self.compaction = Some(CompactionPhase::Commit);
822                    return;
823                };
824                let Some((last_rotated, secrets)) = engine.persisted_self_ratchet_row(&destination)
825                else {
826                    self.compaction = Some(CompactionPhase::Ratchets { index: index + 1 });
827                    return;
828                };
829                let mut scratch = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
830                scratch[..TRUNCATED_HASH_BYTE_LEN].copy_from_slice(destination.as_bytes());
831                let required = self_ratchets_snapshot_len(secrets.len());
832                let end = TRUNCATED_HASH_BYTE_LEN.saturating_add(required);
833                if end > scratch.len() {
834                    self.note_codec_failure(now);
835                    return;
836                }
837                let Ok(written) = write_self_ratchets_snapshot(
838                    last_rotated,
839                    secrets,
840                    &mut scratch[TRUNCATED_HASH_BYTE_LEN..end],
841                ) else {
842                    self.note_codec_failure(now);
843                    return;
844                };
845                match journal
846                    .append_compacted(
847                        FlashJournalRecordKind::SelfRatchet,
848                        &scratch[..TRUNCATED_HASH_BYTE_LEN + written],
849                    )
850                    .await
851                {
852                    Ok(()) => {
853                        self.landing_records = self.landing_records.saturating_add(1);
854                        self.compaction = Some(CompactionPhase::Ratchets { index: index + 1 });
855                    }
856                    Err(error) => {
857                        self.note_write_failure(now, failure_from_journal(error));
858                    }
859                }
860            }
861            CompactionPhase::Commit => match journal.commit_compaction().await {
862                Ok(()) => {
863                    self.compaction = None;
864                    self.compaction_target = None;
865                    let records = core::mem::take(&mut self.landing_records);
866                    self.retry_not_before = None;
867                    self.write_failed = false;
868                    if self.snapshot_required {
869                        self.require_snapshot(self.snapshot_target, now);
870                    } else {
871                        self.deferred_target = None;
872                        self.deferred_until = None;
873                    }
874                    let state_not_saved = self.state_not_saved();
875                    (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::CompactionCompleted {
876                        records,
877                        at: now,
878                        state_not_saved,
879                    });
880                    if !self.snapshot_required
881                        && self.pending_routes.is_empty()
882                        && self.pending_ratchets.is_empty()
883                    {
884                        self.landing_batch = Some(BatchKind::Compaction);
885                    }
886                }
887                Err(error) => {
888                    self.note_write_failure(now, failure_from_journal(error));
889                }
890            },
891        }
892    }
893
894    async fn land_timebase(&mut self, batch: BatchKind, now: InstantMillis) {
895        if !self.record_timebase(now).await {
896            return;
897        }
898        match batch {
899            BatchKind::Routes => {
900                self.route_dirty_since = None;
901                self.last_route_success = Some(now);
902            }
903            BatchKind::Ratchets => {
904                self.ratchet_dirty_since = None;
905            }
906            BatchKind::Compaction => {
907                self.route_dirty_since = None;
908                self.ratchet_dirty_since = None;
909                self.last_route_success = Some(now);
910            }
911        }
912        self.retry_not_before = None;
913        self.write_failed = false;
914        let records = core::mem::take(&mut self.landing_records);
915        self.landing_batch = None;
916        if batch != BatchKind::Compaction {
917            let state_not_saved = self.state_not_saved();
918            (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::BatchPersisted {
919                records,
920                at: now,
921                state_not_saved,
922            });
923        }
924    }
925
926    async fn record_timebase(&mut self, now: InstantMillis) -> bool {
927        let should_record = self.last_timebase_success.is_none_or(|last| {
928            now.0.saturating_sub(last.0) >= self.policy.timebase_record_interval_millis
929        });
930        if !should_record {
931            return true;
932        }
933        let Some(journal) = self.journal.as_mut() else {
934            return false;
935        };
936        if let Err(error) = journal.record_timebase(now).await {
937            self.note_write_failure(now, failure_from_journal(error));
938            return false;
939        }
940        self.last_timebase_success = Some(now);
941        self.retry_not_before = None;
942        self.write_failed = false;
943        true
944    }
945
946    fn note_codec_failure(&mut self, now: InstantMillis) {
947        self.note_write_failure(now, EmbeddedPersistenceFailure::Codec);
948    }
949
950    fn note_write_failure(&mut self, now: InstantMillis, failure: EmbeddedPersistenceFailure) {
951        let retry_at = InstantMillis(now.0.saturating_add(self.policy.retry_interval_millis));
952        self.retry_not_before = Some(retry_at);
953        self.write_failed = true;
954        if self.compaction.is_some() {
955            let target = self
956                .compaction_target
957                .unwrap_or(EmbeddedPersistenceTarget::Routes);
958            if let Some(journal) = self.journal.as_mut() {
959                journal.abort_compaction();
960            }
961            self.compaction = None;
962            self.compaction_target = None;
963            self.require_snapshot(target, now);
964            match target {
965                EmbeddedPersistenceTarget::Routes => {
966                    self.route_dirty_since.get_or_insert(now);
967                }
968                EmbeddedPersistenceTarget::CriticalState => {
969                    self.ratchet_dirty_since.get_or_insert(now);
970                }
971            }
972        }
973        (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::WriteFailed { failure, retry_at });
974    }
975}
976
977pub(crate) trait ManifoldPersistence<S: StorageLayout> {
978    fn observe(&mut self, journaled: &Journaled<'_>, now: InstantMillis);
979    fn deadline(&self, now: InstantMillis) -> Option<InstantMillis>;
980    async fn progress(&mut self, engine: &mut EngineState<S>, now: InstantMillis);
981}
982
983impl<S, F, Keys, Observe, const PENDING: usize> ManifoldPersistence<S>
984    for EmbeddedFlashPersistence<F, Keys, Observe, PENDING>
985where
986    S: StorageLayout,
987    F: NorFlash,
988    Keys: RouteSnapshotKeys,
989    Observe: FnMut(EmbeddedPersistenceDiagnostic),
990{
991    fn observe(&mut self, journaled: &Journaled<'_>, now: InstantMillis) {
992        self.observe_journaled(journaled, now);
993    }
994
995    fn deadline(&self, now: InstantMillis) -> Option<InstantMillis> {
996        self.next_deadline(now)
997    }
998
999    async fn progress(&mut self, engine: &mut EngineState<S>, now: InstantMillis) {
1000        self.progress(engine, now).await;
1001    }
1002}
1003
1004pub(crate) struct NoManifoldPersistence;
1005
1006impl<S: StorageLayout> ManifoldPersistence<S> for NoManifoldPersistence {
1007    fn observe(&mut self, _journaled: &Journaled<'_>, _now: InstantMillis) {}
1008
1009    fn deadline(&self, _now: InstantMillis) -> Option<InstantMillis> {
1010        None
1011    }
1012
1013    async fn progress(&mut self, _engine: &mut EngineState<S>, _now: InstantMillis) {}
1014}
1015
1016fn encode_route_delta<S: StorageLayout>(
1017    engine: &EngineState<S>,
1018    delta: PendingRouteDelta,
1019) -> Result<EncodedDelta, ()> {
1020    match delta {
1021        PendingRouteDelta::RouteUpsert(destination) => {
1022            let Some(row) = engine.persisted_route_row(&destination) else {
1023                return encode_tombstone(destination);
1024            };
1025            let mut durable = row.clone();
1026            durable.announce_id_ring = AnnounceIdRing::Table(&[]);
1027            let required = routing_table_snapshot_len(core::iter::once(durable.clone()));
1028            if required > RECORD_SCRATCH_LEN {
1029                return Err(());
1030            }
1031            let mut scratch = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
1032            let written =
1033                write_routing_table_snapshot(core::iter::once(durable), &mut scratch[..required])
1034                    .map_err(|_| ())?;
1035            Ok(EncodedDelta {
1036                kind: FlashJournalRecordKind::RouteUpsert,
1037                payload: scratch,
1038                len: written,
1039            })
1040        }
1041        PendingRouteDelta::RouteRemoval(destination) => encode_tombstone(destination),
1042    }
1043}
1044
1045fn encode_ratchet<S: StorageLayout>(
1046    engine: &EngineState<S>,
1047    destination: DestinationHash,
1048) -> Result<EncodedDelta, ()> {
1049    let Some((last_rotated, secrets)) = engine.persisted_self_ratchet_row(&destination) else {
1050        return Err(());
1051    };
1052    let required = self_ratchets_snapshot_len(secrets.len());
1053    if TRUNCATED_HASH_BYTE_LEN + required > RECORD_SCRATCH_LEN {
1054        return Err(());
1055    }
1056    let mut scratch = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
1057    scratch[..TRUNCATED_HASH_BYTE_LEN].copy_from_slice(destination.as_bytes());
1058    let written = write_self_ratchets_snapshot(
1059        last_rotated,
1060        secrets,
1061        &mut scratch[TRUNCATED_HASH_BYTE_LEN..TRUNCATED_HASH_BYTE_LEN + required],
1062    )
1063    .map_err(|_| ())?;
1064    Ok(EncodedDelta {
1065        kind: FlashJournalRecordKind::SelfRatchet,
1066        payload: scratch,
1067        len: TRUNCATED_HASH_BYTE_LEN + written,
1068    })
1069}
1070
1071fn encode_tombstone(destination: DestinationHash) -> Result<EncodedDelta, ()> {
1072    let mut payload = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
1073    payload[..TRUNCATED_HASH_BYTE_LEN].copy_from_slice(destination.as_bytes());
1074    Ok(EncodedDelta {
1075        kind: FlashJournalRecordKind::RouteRemoval,
1076        payload,
1077        len: TRUNCATED_HASH_BYTE_LEN,
1078    })
1079}
1080
1081fn apply_record<S: StorageLayout>(
1082    engine: &mut EngineState<S>,
1083    now: InstantMillis,
1084    record: FlashJournalRecord<'_>,
1085    report: &mut EmbeddedPersistenceRestoreReport,
1086) {
1087    match record.kind {
1088        FlashJournalRecordKind::ArenaCommit => {}
1089        FlashJournalRecordKind::RouteUpsert => {
1090            let Ok(mut rows) = read_routing_table_snapshot(record.payload) else {
1091                report.route_refused_count = report.route_refused_count.saturating_add(1);
1092                return;
1093            };
1094            let Some(Ok(row)) = rows.next() else {
1095                report.route_refused_count = report.route_refused_count.saturating_add(1);
1096                return;
1097            };
1098            if rows.next().is_some() {
1099                report.route_refused_count = report.route_refused_count.saturating_add(1);
1100                return;
1101            }
1102            let destination = row.destination;
1103            let Ok(pending) = engine.prepare_persisted_route(row) else {
1104                report.route_refused_count = report.route_refused_count.saturating_add(1);
1105                return;
1106            };
1107            let Ok(verified) = pending.verify() else {
1108                report.route_refused_count = report.route_refused_count.saturating_add(1);
1109                return;
1110            };
1111            let _ = engine.drop_route(&destination, AttachedInterfaces::new(&[]));
1112            match engine.seed_verified_route(verified, now) {
1113                RouteSeedOutcome::Seeded => {
1114                    report.route_seeded_count = report.route_seeded_count.saturating_add(1);
1115                }
1116                RouteSeedOutcome::RefusedDestinationMismatch
1117                | RouteSeedOutcome::RefusedBlackholedIdentity
1118                | RouteSeedOutcome::RefusedInvalidSignature => {
1119                    report.route_refused_count = report.route_refused_count.saturating_add(1);
1120                }
1121                RouteSeedOutcome::AlreadyPresent
1122                | RouteSeedOutcome::TableFull
1123                | RouteSeedOutcome::AppDataArenaFull => {
1124                    report.route_dropped_count = report.route_dropped_count.saturating_add(1);
1125                }
1126            }
1127        }
1128        FlashJournalRecordKind::RouteRemoval => {
1129            let Ok(bytes) = <[u8; TRUNCATED_HASH_BYTE_LEN]>::try_from(record.payload) else {
1130                report.route_refused_count = report.route_refused_count.saturating_add(1);
1131                return;
1132            };
1133            let destination = DestinationHash::new(bytes);
1134            let _ = engine.drop_route(&destination, AttachedInterfaces::new(&[]));
1135        }
1136        FlashJournalRecordKind::SelfRatchet => {
1137            let Some((destination, sealed)) = record
1138                .payload
1139                .split_first_chunk::<TRUNCATED_HASH_BYTE_LEN>()
1140            else {
1141                report.ratchet_refused_count = report.ratchet_refused_count.saturating_add(1);
1142                return;
1143            };
1144            let Ok(restored) = read_self_ratchets_snapshot(sealed) else {
1145                report.ratchet_refused_count = report.ratchet_refused_count.saturating_add(1);
1146                return;
1147            };
1148            match engine.replace_persisted_self_ratchets(
1149                &DestinationHash::new(*destination),
1150                restored.last_rotated,
1151                restored.secrets_newest_first(),
1152            ) {
1153                SeedSelfRatchetsOutcome::Seeded => {
1154                    report.ratchet_seeded_count = report.ratchet_seeded_count.saturating_add(1);
1155                }
1156                SeedSelfRatchetsOutcome::AlreadyMinted | SeedSelfRatchetsOutcome::Untracked => {
1157                    report.ratchet_refused_count = report.ratchet_refused_count.saturating_add(1);
1158                }
1159            }
1160        }
1161    }
1162}
1163
1164fn failure_from_journal<E>(error: FlashJournalError<E>) -> EmbeddedPersistenceFailure {
1165    match error {
1166        FlashJournalError::Flash(_) => EmbeddedPersistenceFailure::Flash,
1167        FlashJournalError::ArenaFull
1168        | FlashJournalError::OutOfBounds
1169        | FlashJournalError::Misaligned
1170        | FlashJournalError::Uninitialized
1171        | FlashJournalError::CompactionInProgress
1172        | FlashJournalError::NoCompaction
1173        | FlashJournalError::PayloadTooLarge
1174        | FlashJournalError::ScratchTooShort => EmbeddedPersistenceFailure::Capacity,
1175    }
1176}
1177
1178fn earlier(first: Option<InstantMillis>, second: Option<InstantMillis>) -> Option<InstantMillis> {
1179    match (first, second) {
1180        (Some(first), Some(second)) => Some(InstantMillis(first.0.min(second.0))),
1181        (Some(value), None) | (None, Some(value)) => Some(value),
1182        (None, None) => None,
1183    }
1184}
1185
1186#[cfg(test)]
1187mod tests {
1188    use super::*;
1189    use crate::persistence::TIMEBASE_HEADROOM_MILLIS;
1190    use embedded_storage::nor_flash::{ErrorType, NorFlashError, NorFlashErrorKind};
1191    use embedded_storage_async::nor_flash::ReadNorFlash;
1192    use std::cell::{Cell, RefCell};
1193    use std::rc::Rc;
1194    use std::vec::Vec;
1195
1196    const ERASE: usize = 512;
1197    const CAPACITY: usize = ERASE * 6;
1198    const LAYOUT: FlashJournalLayout = FlashJournalLayout::new(
1199        [0, ERASE as u32],
1200        [
1201            crate::persistence::FlashArenaRange::new((ERASE * 2) as u32, (ERASE * 4) as u32),
1202            crate::persistence::FlashArenaRange::new((ERASE * 4) as u32, (ERASE * 6) as u32),
1203        ],
1204    );
1205
1206    #[derive(Debug)]
1207    struct TestFlash {
1208        bytes: [u8; CAPACITY],
1209        sector_erases: [u32; CAPACITY / ERASE],
1210        fail_next_write: Rc<Cell<bool>>,
1211    }
1212
1213    impl TestFlash {
1214        fn new() -> Self {
1215            Self {
1216                bytes: [0xFF; CAPACITY],
1217                sector_erases: [0; CAPACITY / ERASE],
1218                fail_next_write: Rc::new(Cell::new(false)),
1219            }
1220        }
1221
1222        fn controlled() -> (Self, Rc<Cell<bool>>) {
1223            let flash = Self::new();
1224            let control = Rc::clone(&flash.fail_next_write);
1225            (flash, control)
1226        }
1227    }
1228
1229    #[derive(Debug)]
1230    struct TestFlashError;
1231
1232    impl NorFlashError for TestFlashError {
1233        fn kind(&self) -> NorFlashErrorKind {
1234            NorFlashErrorKind::Other
1235        }
1236    }
1237
1238    impl ErrorType for TestFlash {
1239        type Error = TestFlashError;
1240    }
1241
1242    impl ReadNorFlash for TestFlash {
1243        const READ_SIZE: usize = 4;
1244
1245        async fn read(&mut self, offset: u32, bytes: &mut [u8]) -> Result<(), Self::Error> {
1246            let start = offset as usize;
1247            let end = start + bytes.len();
1248            bytes.copy_from_slice(&self.bytes[start..end]);
1249            Ok(())
1250        }
1251
1252        fn capacity(&self) -> usize {
1253            CAPACITY
1254        }
1255    }
1256
1257    impl NorFlash for TestFlash {
1258        const WRITE_SIZE: usize = 4;
1259        const ERASE_SIZE: usize = ERASE;
1260
1261        async fn write(&mut self, offset: u32, bytes: &[u8]) -> Result<(), Self::Error> {
1262            if self.fail_next_write.replace(false) {
1263                return Err(TestFlashError);
1264            }
1265            let start = offset as usize;
1266            for (stored, written) in self.bytes[start..start + bytes.len()].iter_mut().zip(bytes) {
1267                *stored &= *written;
1268            }
1269            Ok(())
1270        }
1271
1272        async fn erase(&mut self, from: u32, to: u32) -> Result<(), Self::Error> {
1273            self.bytes[from as usize..to as usize].fill(0xFF);
1274            for sector in from as usize / ERASE..to as usize / ERASE {
1275                self.sector_erases[sector] = self.sector_erases[sector].saturating_add(1);
1276            }
1277            Ok(())
1278        }
1279    }
1280
1281    fn ready_with_observer<Observe>(
1282        observe: Observe,
1283    ) -> EmbeddedFlashPersistence<TestFlash, FixedRouteSnapshotKeys<8>, Observe, 4>
1284    where
1285        Observe: FnMut(EmbeddedPersistenceDiagnostic),
1286    {
1287        embassy_futures::block_on(async {
1288            let flash = TestFlash::new();
1289            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1290            let (mut journal, _) = FlashJournal::open(flash, LAYOUT, &mut scratch, |_| {})
1291                .await
1292                .unwrap();
1293            journal.initialize_empty().await.unwrap();
1294            journal
1295                .record_compaction_budget(InstantMillis(0))
1296                .await
1297                .unwrap();
1298            let mut persistence = EmbeddedFlashPersistence::new(
1299                TestFlash::new(),
1300                LAYOUT,
1301                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1302                FixedRouteSnapshotKeys::new(),
1303                observe,
1304            );
1305            persistence.flash = None;
1306            persistence.journal = Some(journal);
1307            persistence.next_compaction_not_before =
1308                Some(InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS));
1309            persistence.last_timebase_success = Some(InstantMillis(0));
1310            persistence
1311        })
1312    }
1313
1314    fn ready() -> EmbeddedFlashPersistence<
1315        TestFlash,
1316        FixedRouteSnapshotKeys<8>,
1317        fn(EmbeddedPersistenceDiagnostic),
1318        4,
1319    > {
1320        ready_with_observer((|_| {}) as fn(EmbeddedPersistenceDiagnostic))
1321    }
1322
1323    fn signed_route(secret: u8, app_data: &[u8]) -> crate::routing::PersistedRouteRow<'_> {
1324        use crate::identity::in_memory::InMemoryNodeIdentity;
1325        use crate::interfaces::InterfaceId;
1326        use crate::routing::announce::{Announce, AnnounceId, DottedNameHash};
1327        use crate::routing::routes::RouteEntry;
1328        use crate::routing::{AnnounceIdRing, NextHop, RouteResponsiveness};
1329
1330        let signer = InMemoryNodeIdentity::from_secret_key_bytes(&[secret; 64]);
1331        let announce = Announce::build_signed(
1332            &signer,
1333            DottedNameHash::new([secret; 10]),
1334            AnnounceId::from_wire([secret.wrapping_add(1); 10]),
1335            None,
1336            app_data,
1337        )
1338        .unwrap();
1339        crate::routing::PersistedRouteRow {
1340            destination: announce.destination,
1341            entry: RouteEntry {
1342                hops: secret,
1343                learned_at: InstantMillis(500),
1344                last_route_activity_at: InstantMillis(700),
1345                responsiveness: RouteResponsiveness::Responsive,
1346                receiving_interface: InterfaceId::new([secret; 8]),
1347                next_hop: NextHop::Direct,
1348            },
1349            public_keys: announce.public_keys,
1350            dotted_name_hash: announce.dotted_name_hash,
1351            announce_id: announce.announce_id,
1352            ratchet: announce.ratchet,
1353            signature: announce.signature,
1354            app_data,
1355            announce_id_ring: AnnounceIdRing::Wire(
1356                &[0; crate::routing::announce::ANNOUNCE_ID_WIRE_LEN],
1357            ),
1358        }
1359    }
1360
1361    #[test]
1362    fn exact_route_and_ratchet_deadlines_are_distinct() {
1363        let policy =
1364            EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(64));
1365        assert_eq!(policy.first_route_commit_delay_millis, 2_000);
1366        assert_eq!(policy.minimum_route_commit_interval_millis, 300_000);
1367        assert_eq!(policy.ratchet_batch_delay_millis, 2_000);
1368        assert_eq!(policy.retry_interval_millis, 300_000);
1369        assert_eq!(
1370            policy.timebase_record_interval_millis,
1371            TIMEBASE_RECORD_INTERVAL_MILLIS
1372        );
1373        assert_eq!(
1374            policy.compaction.minimum_interval_millis,
1375            HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS
1376        );
1377        assert_eq!(policy.compaction.critical_reserve_bytes, 64);
1378    }
1379
1380    #[test]
1381    fn deadline_formula_batches_first_write_then_honors_five_minutes() {
1382        let mut persistence = ready();
1383        persistence.route_dirty_since = Some(InstantMillis(1_000));
1384        assert_eq!(
1385            persistence.next_deadline(InstantMillis(1_500)),
1386            Some(InstantMillis(3_000))
1387        );
1388        persistence.last_route_success = Some(InstantMillis(10_000));
1389        persistence.route_dirty_since = Some(InstantMillis(11_000));
1390        assert_eq!(
1391            persistence.next_deadline(InstantMillis(12_000)),
1392            Some(InstantMillis(310_000))
1393        );
1394        persistence.ratchet_dirty_since = Some(InstantMillis(12_000));
1395        assert_eq!(
1396            persistence.next_deadline(InstantMillis(12_000)),
1397            Some(InstantMillis(14_000))
1398        );
1399        persistence.retry_not_before = Some(InstantMillis(400_000));
1400        assert_eq!(
1401            persistence.next_deadline(InstantMillis(12_000)),
1402            Some(InstantMillis(400_000))
1403        );
1404    }
1405
1406    #[test]
1407    fn repeated_route_and_ratchet_changes_coalesce_by_destination() {
1408        let mut persistence = ready();
1409        let destination = DestinationHash::new([0x11; TRUNCATED_HASH_BYTE_LEN]);
1410        persistence.queue_route(
1411            PendingRouteDelta::RouteUpsert(destination),
1412            InstantMillis(0),
1413        );
1414        persistence.queue_route(
1415            PendingRouteDelta::RouteUpsert(destination),
1416            InstantMillis(0),
1417        );
1418        persistence.queue_route(
1419            PendingRouteDelta::RouteRemoval(destination),
1420            InstantMillis(0),
1421        );
1422        persistence.queue_ratchet(destination, InstantMillis(0));
1423        persistence.queue_ratchet(destination, InstantMillis(0));
1424        assert_eq!(
1425            persistence.pending_routes.as_slice(),
1426            &[PendingRouteDelta::RouteRemoval(destination)]
1427        );
1428        assert_eq!(persistence.pending_ratchets.as_slice(), &[destination]);
1429    }
1430
1431    #[test]
1432    fn pending_overflow_waits_for_the_batch_deadline_before_compacting() {
1433        let mut persistence = ready();
1434        for byte in 0..5 {
1435            persistence.queue_route(
1436                PendingRouteDelta::RouteUpsert(DestinationHash::new(
1437                    [byte; TRUNCATED_HASH_BYTE_LEN],
1438                )),
1439                InstantMillis(1_000),
1440            );
1441        }
1442        persistence.route_dirty_since = Some(InstantMillis(1_000));
1443        assert!(persistence.snapshot_required);
1444        assert_eq!(
1445            persistence.deferred_target,
1446            Some(EmbeddedPersistenceTarget::Routes)
1447        );
1448        assert_eq!(
1449            persistence.next_deadline(InstantMillis(1_500)),
1450            Some(InstantMillis(TIMEBASE_RECORD_INTERVAL_MILLIS))
1451        );
1452
1453        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1454        embassy_futures::block_on(persistence.progress(&mut engine, InstantMillis(2_999)));
1455        assert_eq!(persistence.compaction, None);
1456        assert!(persistence.snapshot_required);
1457
1458        embassy_futures::block_on(persistence.progress(
1459            &mut engine,
1460            InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS),
1461        ));
1462        assert_eq!(
1463            persistence.compaction,
1464            Some(CompactionPhase::RecordBudget {
1465                at: InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS)
1466            })
1467        );
1468    }
1469
1470    #[test]
1471    fn failures_keep_dirty_state_and_raise_the_notice() {
1472        let mut persistence = ready();
1473        let destination = DestinationHash::new([0x22; TRUNCATED_HASH_BYTE_LEN]);
1474        persistence.queue_route(
1475            PendingRouteDelta::RouteUpsert(destination),
1476            InstantMillis(0),
1477        );
1478        persistence.note_codec_failure(InstantMillis(1_000));
1479        assert_eq!(persistence.pending_routes.len(), 1);
1480        assert_eq!(persistence.retry_not_before, Some(InstantMillis(301_000)));
1481        assert!(persistence.state_not_saved());
1482        assert!(!persistence.snapshot_required);
1483
1484        persistence.retry_not_before = None;
1485        persistence.compaction = Some(CompactionPhase::Commit);
1486        persistence.compaction_target = Some(EmbeddedPersistenceTarget::Routes);
1487        persistence.note_write_failure(InstantMillis(2_000), EmbeddedPersistenceFailure::Flash);
1488        assert_eq!(persistence.pending_routes.len(), 0);
1489        assert_eq!(persistence.retry_not_before, Some(InstantMillis(302_000)));
1490        assert!(persistence.state_not_saved());
1491        assert!(persistence.snapshot_required);
1492        assert_eq!(persistence.compaction, None);
1493    }
1494
1495    #[test]
1496    fn legacy_timebase_allows_one_needed_compaction_then_adopts_the_budget_marker() {
1497        embassy_futures::block_on(async {
1498            let (mut journal, _) = {
1499                let flash = TestFlash::new();
1500                let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1501                FlashJournal::open(flash, LAYOUT, &mut scratch, |_| {})
1502                    .await
1503                    .unwrap()
1504            };
1505            journal.initialize_empty().await.unwrap();
1506            journal
1507                .record_timebase(InstantMillis(10_000))
1508                .await
1509                .unwrap();
1510            let mut persistence =
1511                EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1512                    journal.release(),
1513                    LAYOUT,
1514                    EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(
1515                        0,
1516                    )),
1517                    FixedRouteSnapshotKeys::new(),
1518                    (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1519                );
1520            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1521            let report = persistence.restore(&mut engine, InstantMillis(0)).await;
1522            assert_eq!(persistence.next_compaction_not_before, None);
1523            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, report.logical_start);
1524            persistence.route_dirty_since = Some(InstantMillis(report.logical_start.0 - 2_000));
1525            persistence
1526                .progress(&mut engine, report.logical_start)
1527                .await;
1528            persistence
1529                .progress(&mut engine, report.logical_start)
1530                .await;
1531            assert!(matches!(
1532                persistence.compaction,
1533                Some(CompactionPhase::Erase { sector: 0 })
1534            ));
1535
1536            let flash = persistence.journal.take().unwrap().release();
1537            let mut restored = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1538                flash,
1539                LAYOUT,
1540                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1541                FixedRouteSnapshotKeys::new(),
1542                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1543            );
1544            restored.restore(&mut engine, InstantMillis(0)).await;
1545            assert!(restored.next_compaction_not_before.is_some());
1546        });
1547    }
1548
1549    #[test]
1550    fn restore_uses_the_later_of_flash_high_water_and_the_raw_clock() {
1551        embassy_futures::block_on(async {
1552            let recorded_at = InstantMillis(10_000);
1553            let flash_high_water = InstantMillis(recorded_at.0 + TIMEBASE_HEADROOM_MILLIS);
1554            let rtc_after_downtime = InstantMillis(flash_high_water.0 + 86_400_000);
1555
1556            for raw_now in [InstantMillis(0), rtc_after_downtime] {
1557                let flash = TestFlash::new();
1558                let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1559                let (mut journal, _) = FlashJournal::open(flash, LAYOUT, &mut scratch, |_| {})
1560                    .await
1561                    .unwrap();
1562                journal.initialize_empty().await.unwrap();
1563                journal.record_timebase(recorded_at).await.unwrap();
1564
1565                let mut persistence =
1566                    EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1567                        journal.release(),
1568                        LAYOUT,
1569                        EmbeddedPersistencePolicy::hopspot_default(
1570                            EmbeddedCompactionPolicy::hopspot(0),
1571                        ),
1572                        FixedRouteSnapshotKeys::new(),
1573                        (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1574                    );
1575                let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1576                let report = persistence.restore(&mut engine, raw_now).await;
1577
1578                assert_eq!(report.logical_start, raw_now.max(flash_high_water));
1579            }
1580        });
1581    }
1582
1583    #[test]
1584    fn idle_persistence_advances_the_flash_timebase_on_schedule() {
1585        embassy_futures::block_on(async {
1586            let mut persistence =
1587                EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1588                    TestFlash::new(),
1589                    LAYOUT,
1590                    EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(
1591                        0,
1592                    )),
1593                    FixedRouteSnapshotKeys::new(),
1594                    (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1595                );
1596            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1597            let start = persistence.restore(&mut engine, InstantMillis(1_000)).await;
1598
1599            assert_eq!(
1600                persistence.next_deadline(start.logical_start),
1601                Some(start.logical_start)
1602            );
1603            persistence.progress(&mut engine, start.logical_start).await;
1604            assert_eq!(persistence.last_timebase_success, Some(start.logical_start));
1605
1606            let next = InstantMillis(
1607                start
1608                    .logical_start
1609                    .0
1610                    .saturating_add(TIMEBASE_RECORD_INTERVAL_MILLIS),
1611            );
1612            assert_eq!(persistence.next_deadline(start.logical_start), Some(next));
1613            persistence
1614                .progress(&mut engine, InstantMillis(next.0 - 1))
1615                .await;
1616            assert_eq!(persistence.last_timebase_success, Some(start.logical_start));
1617            persistence.progress(&mut engine, next).await;
1618            assert_eq!(persistence.last_timebase_success, Some(next));
1619
1620            let mut flash = persistence.journal.take().unwrap().release();
1621            assert_eq!(
1622                FlashJournal::inspect_timebase(&mut flash, LAYOUT)
1623                    .await
1624                    .unwrap(),
1625                Some(InstantMillis(next.0 + TIMEBASE_HEADROOM_MILLIS))
1626            );
1627        });
1628    }
1629
1630    #[test]
1631    fn failed_compaction_attempt_consumes_the_daily_budget() {
1632        let diagnostics = Rc::new(RefCell::new(Vec::new()));
1633        let observed = Rc::clone(&diagnostics);
1634        let mut persistence = ready_with_observer(move |diagnostic| {
1635            observed.borrow_mut().push(diagnostic);
1636        });
1637        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1638        let first = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1639        persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, first);
1640        persistence.route_dirty_since = Some(InstantMillis(first.0 - 2_000));
1641        embassy_futures::block_on(persistence.progress(&mut engine, first));
1642        let second = InstantMillis(first.0 + HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1643        assert_eq!(persistence.next_compaction_not_before, Some(first));
1644        embassy_futures::block_on(persistence.progress(&mut engine, first));
1645        assert_eq!(persistence.next_compaction_not_before, Some(second));
1646        assert_eq!(
1647            persistence.compaction,
1648            Some(CompactionPhase::Erase { sector: 0 })
1649        );
1650
1651        persistence.note_write_failure(first, EmbeddedPersistenceFailure::Flash);
1652        assert_eq!(persistence.compaction, None);
1653        assert!(persistence.snapshot_required);
1654        assert_eq!(
1655            persistence.next_deadline(first),
1656            Some(InstantMillis(first.0 + TIMEBASE_RECORD_INTERVAL_MILLIS))
1657        );
1658        embassy_futures::block_on(persistence.progress(&mut engine, InstantMillis(second.0 - 1)));
1659        assert_eq!(persistence.compaction, None);
1660        embassy_futures::block_on(persistence.progress(&mut engine, second));
1661        assert_eq!(
1662            persistence.compaction,
1663            Some(CompactionPhase::RecordBudget { at: second })
1664        );
1665        embassy_futures::block_on(persistence.progress(&mut engine, second));
1666        assert_eq!(
1667            persistence.compaction,
1668            Some(CompactionPhase::Erase { sector: 0 })
1669        );
1670
1671        let starts = diagnostics
1672            .borrow()
1673            .iter()
1674            .filter(|diagnostic| {
1675                matches!(
1676                    diagnostic,
1677                    EmbeddedPersistenceDiagnostic::CompactionStarted { .. }
1678                )
1679            })
1680            .count();
1681        assert_eq!(starts, 2);
1682    }
1683
1684    #[test]
1685    fn recorded_compaction_budget_survives_reboot() {
1686        embassy_futures::block_on(async {
1687            let mut persistence = ready();
1688            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1689            let attempt = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1690            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, attempt);
1691            persistence.route_dirty_since = Some(InstantMillis(attempt.0 - 2_000));
1692            persistence.progress(&mut engine, attempt).await;
1693            persistence.progress(&mut engine, attempt).await;
1694            assert_eq!(
1695                persistence.compaction,
1696                Some(CompactionPhase::Erase { sector: 0 })
1697            );
1698
1699            let flash = persistence.journal.take().unwrap().release();
1700            let mut restored = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1701                flash,
1702                LAYOUT,
1703                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1704                FixedRouteSnapshotKeys::new(),
1705                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1706            );
1707            let report = restored.restore(&mut engine, InstantMillis(0)).await;
1708            assert!(report.logical_start.0 >= attempt.0);
1709            assert_eq!(
1710                restored.next_compaction_not_before,
1711                Some(InstantMillis(
1712                    attempt
1713                        .0
1714                        .saturating_add(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS)
1715                ))
1716            );
1717        });
1718    }
1719
1720    #[test]
1721    fn marker_write_failure_does_not_consume_the_compaction_budget() {
1722        embassy_futures::block_on(async {
1723            let (flash, fail_next_write) = TestFlash::controlled();
1724            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1725            let (mut journal, _) = FlashJournal::open(flash, LAYOUT, &mut scratch, |_| {})
1726                .await
1727                .unwrap();
1728            journal.initialize_empty().await.unwrap();
1729            let mut persistence =
1730                EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1731                    TestFlash::new(),
1732                    LAYOUT,
1733                    EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(
1734                        0,
1735                    )),
1736                    FixedRouteSnapshotKeys::new(),
1737                    (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1738                );
1739            persistence.flash = None;
1740            persistence.journal = Some(journal);
1741            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1742            let attempt = InstantMillis(2_000);
1743            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, attempt);
1744            persistence.route_dirty_since = Some(InstantMillis(0));
1745            persistence.progress(&mut engine, attempt).await;
1746            fail_next_write.set(true);
1747            persistence.progress(&mut engine, attempt).await;
1748            assert_eq!(persistence.next_compaction_not_before, None);
1749            assert_eq!(persistence.compaction, None);
1750
1751            let flash = persistence.journal.take().unwrap().release();
1752            let mut restored = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1753                flash,
1754                LAYOUT,
1755                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1756                FixedRouteSnapshotKeys::new(),
1757                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1758            );
1759            restored.restore(&mut engine, InstantMillis(0)).await;
1760            assert_eq!(restored.next_compaction_not_before, None);
1761        });
1762    }
1763
1764    #[test]
1765    fn timebase_writes_and_repeated_reboots_do_not_move_the_compaction_deadline() {
1766        embassy_futures::block_on(async {
1767            let mut persistence = ready();
1768            let deadline = persistence.next_compaction_not_before;
1769            persistence
1770                .journal
1771                .as_mut()
1772                .unwrap()
1773                .record_timebase(InstantMillis(3 * 60 * 60 * 1_000))
1774                .await
1775                .unwrap();
1776            let flash = persistence.journal.take().unwrap().release();
1777            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1778
1779            let mut first = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1780                flash,
1781                LAYOUT,
1782                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1783                FixedRouteSnapshotKeys::new(),
1784                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1785            );
1786            first.restore(&mut engine, InstantMillis(0)).await;
1787            assert_eq!(first.next_compaction_not_before, deadline);
1788
1789            let flash = first.journal.take().unwrap().release();
1790            let mut second = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1791                flash,
1792                LAYOUT,
1793                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1794                FixedRouteSnapshotKeys::new(),
1795                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1796            );
1797            second.restore(&mut engine, InstantMillis(0)).await;
1798            assert_eq!(second.next_compaction_not_before, deadline);
1799        });
1800    }
1801
1802    #[test]
1803    fn overflow_during_compaction_commits_once_and_defers_the_next_snapshot() {
1804        let diagnostics = Rc::new(RefCell::new(Vec::new()));
1805        let observed = Rc::clone(&diagnostics);
1806        let mut persistence = ready_with_observer(move |diagnostic| {
1807            observed.borrow_mut().push(diagnostic);
1808        });
1809        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1810        let now = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1811        persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
1812        persistence.route_dirty_since = Some(InstantMillis(now.0 - 2_000));
1813        embassy_futures::block_on(persistence.progress(&mut engine, now));
1814
1815        for byte in 0..6 {
1816            persistence.queue_route(
1817                PendingRouteDelta::RouteUpsert(DestinationHash::new(
1818                    [byte; TRUNCATED_HASH_BYTE_LEN],
1819                )),
1820                now,
1821            );
1822        }
1823        persistence.route_dirty_since = Some(now);
1824        for _ in 0..8 {
1825            embassy_futures::block_on(persistence.progress(&mut engine, now));
1826        }
1827
1828        assert_eq!(persistence.compaction, None);
1829        assert!(persistence.snapshot_required);
1830        assert!(persistence.state_not_saved());
1831        assert_eq!(
1832            persistence.next_deadline(now),
1833            Some(InstantMillis(now.0 + TIMEBASE_RECORD_INTERVAL_MILLIS))
1834        );
1835        let diagnostics = diagnostics.borrow();
1836        assert_eq!(
1837            diagnostics
1838                .iter()
1839                .filter(|diagnostic| matches!(
1840                    diagnostic,
1841                    EmbeddedPersistenceDiagnostic::CompactionStarted { .. }
1842                ))
1843                .count(),
1844            1
1845        );
1846        assert_eq!(
1847            diagnostics
1848                .iter()
1849                .filter(|diagnostic| matches!(
1850                    diagnostic,
1851                    EmbeddedPersistenceDiagnostic::CompactionCompleted { .. }
1852                ))
1853                .count(),
1854            1
1855        );
1856        assert_eq!(
1857            diagnostics
1858                .iter()
1859                .filter(|diagnostic| matches!(
1860                    diagnostic,
1861                    EmbeddedPersistenceDiagnostic::DurabilityDeferred { .. }
1862                ))
1863                .count(),
1864            1
1865        );
1866    }
1867
1868    #[test]
1869    fn captured_route_keys_survive_slot_shifts_and_new_routes_land_after_compaction() {
1870        embassy_futures::block_on(async {
1871            let mut persistence = ready();
1872            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1873            let rows = [signed_route(0x31, &[0xA1]), signed_route(0x32, &[0xA2])];
1874            for row in &rows {
1875                assert_eq!(
1876                    engine.seed_route(row, InstantMillis(1_000)),
1877                    RouteSeedOutcome::Seeded
1878                );
1879            }
1880            let now = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1881            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
1882            persistence.route_dirty_since = Some(InstantMillis(now.0 - 2_000));
1883            persistence.progress(&mut engine, now).await;
1884            persistence.progress(&mut engine, now).await;
1885            persistence.progress(&mut engine, now).await;
1886
1887            let removed = rows[0].destination;
1888            let retained = rows[1].destination;
1889            let _ = engine.drop_route(&removed, AttachedInterfaces::new(&[]));
1890            let added = signed_route(0x34, &[0xA4]);
1891            assert_eq!(engine.seed_route(&added, now), RouteSeedOutcome::Seeded);
1892            persistence.queue_route(PendingRouteDelta::RouteUpsert(added.destination), now);
1893            persistence.route_dirty_since = Some(now);
1894
1895            for _ in 0..8 {
1896                persistence.progress(&mut engine, now).await;
1897            }
1898            assert_eq!(persistence.compaction, None);
1899            assert_eq!(
1900                (
1901                    persistence.pending_routes.len(),
1902                    persistence.snapshot_required,
1903                    persistence.write_failed,
1904                    persistence.route_dirty_since,
1905                    persistence.landing_batch,
1906                ),
1907                (1, false, false, Some(now), None)
1908            );
1909            let correction_at = InstantMillis(now.0 + 2_000);
1910            persistence.progress(&mut engine, correction_at).await;
1911            persistence.progress(&mut engine, correction_at).await;
1912
1913            let flash = persistence.journal.take().unwrap().release();
1914            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1915            let mut restored = Vec::new();
1916            let _ = FlashJournal::open(flash, LAYOUT, &mut scratch, |record| {
1917                if record.kind != FlashJournalRecordKind::RouteUpsert {
1918                    return;
1919                }
1920                let mut rows = read_routing_table_snapshot(record.payload).unwrap();
1921                restored.push(rows.next().unwrap().unwrap().destination);
1922            })
1923            .await
1924            .unwrap();
1925            assert_eq!(restored, vec![retained, added.destination]);
1926        });
1927    }
1928
1929    #[test]
1930    fn sixteen_route_records_restore_eight_and_report_capacity_drops() {
1931        type EightRouteStorage =
1932            crate::storage::TestFixedStorage<8, 8, 256, 2, 2, 16, 4, 4, 4, 4, 4, 16>;
1933        let now = InstantMillis(1_000);
1934        let mut engine = EngineState::<EightRouteStorage>::default();
1935        let mut report = EmbeddedPersistenceRestoreReport {
1936            logical_start: now,
1937            route_seeded_count: 0,
1938            route_refused_count: 0,
1939            route_dropped_count: 0,
1940            ratchet_seeded_count: 0,
1941            ratchet_refused_count: 0,
1942            warning: None,
1943        };
1944
1945        for secret in 1..=16 {
1946            let row = signed_route(secret, &[]);
1947            let required = routing_table_snapshot_len(core::iter::once(row.clone()));
1948            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1949            let written =
1950                write_routing_table_snapshot(core::iter::once(row), &mut scratch[..required])
1951                    .unwrap();
1952            apply_record(
1953                &mut engine,
1954                now,
1955                FlashJournalRecord {
1956                    epoch: 1,
1957                    kind: FlashJournalRecordKind::RouteUpsert,
1958                    payload: &scratch[..written],
1959                },
1960                &mut report,
1961            );
1962        }
1963
1964        assert_eq!(
1965            report,
1966            EmbeddedPersistenceRestoreReport {
1967                logical_start: now,
1968                route_seeded_count: 8,
1969                route_refused_count: 0,
1970                route_dropped_count: 8,
1971                ratchet_seeded_count: 0,
1972                ratchet_refused_count: 0,
1973                warning: None,
1974            }
1975        );
1976    }
1977
1978    #[test]
1979    fn thirty_days_of_pressure_erase_each_arena_sector_at_most_fifteen_times() {
1980        let diagnostics = Rc::new(RefCell::new(Vec::new()));
1981        let observed = Rc::clone(&diagnostics);
1982        let mut persistence = ready_with_observer(move |diagnostic| {
1983            observed.borrow_mut().push(diagnostic);
1984        });
1985        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1986        embassy_futures::block_on(async {
1987            for day in 1..=30 {
1988                let now = InstantMillis(day * HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1989                persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
1990                persistence.route_dirty_since = Some(InstantMillis(now.0 - 2_000));
1991                for _ in 0..8 {
1992                    persistence.progress(&mut engine, now).await;
1993                }
1994                assert_eq!(persistence.compaction, None);
1995            }
1996        });
1997        assert_eq!(
1998            diagnostics
1999                .borrow()
2000                .iter()
2001                .filter(|diagnostic| matches!(
2002                    diagnostic,
2003                    EmbeddedPersistenceDiagnostic::CompactionStarted { .. }
2004                ))
2005                .count(),
2006            30
2007        );
2008        let flash = persistence.journal.take().unwrap().release();
2009        assert_eq!(
2010            [
2011                flash.sector_erases[2] - 1,
2012                flash.sector_erases[3] - 1,
2013                flash.sector_erases[4],
2014                flash.sector_erases[5],
2015            ],
2016            [15, 15, 15, 15]
2017        );
2018    }
2019}