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            .unwrap_or(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        match (deadline, self.retry_not_before) {
499            (Some(deadline), Some(retry)) => Some(InstantMillis(deadline.0.max(retry.0))),
500            (deadline, None) => deadline,
501            (None, Some(_)) => None,
502        }
503    }
504
505    fn ratchet_ready_at(&self) -> Option<InstantMillis> {
506        self.ratchet_dirty_since.map(|dirty| {
507            InstantMillis(
508                dirty
509                    .0
510                    .saturating_add(self.policy.ratchet_batch_delay_millis),
511            )
512        })
513    }
514
515    fn route_ready_at(&self) -> Option<InstantMillis> {
516        self.route_dirty_since.map(|dirty| {
517            let first_ready = dirty
518                .0
519                .saturating_add(self.policy.first_route_commit_delay_millis);
520            let interval_ready = self.last_route_success.map_or(0, |last| {
521                last.0
522                    .saturating_add(self.policy.minimum_route_commit_interval_millis)
523            });
524            InstantMillis(first_ready.max(interval_ready))
525        })
526    }
527
528    async fn progress<S: StorageLayout>(
529        &mut self,
530        engine: &mut EngineState<S>,
531        now: InstantMillis,
532    ) {
533        if self.retry_not_before.is_some_and(|retry| now.0 < retry.0) {
534            return;
535        }
536        if self.compaction.is_some() {
537            self.progress_compaction(engine, now).await;
538            return;
539        }
540        if let Some(batch) = self.landing_batch {
541            let new_work = match batch {
542                BatchKind::Routes => !self.pending_routes.is_empty() || self.snapshot_required,
543                BatchKind::Ratchets => !self.pending_ratchets.is_empty(),
544                BatchKind::Compaction => {
545                    !self.pending_routes.is_empty()
546                        || !self.pending_ratchets.is_empty()
547                        || self.snapshot_required
548                }
549            };
550            if new_work {
551                self.landing_batch = None;
552            } else {
553                self.land_timebase(batch, now).await;
554                return;
555            }
556        }
557        let ratchet_due = self
558            .ratchet_ready_at()
559            .is_some_and(|ready| now.0 >= ready.0);
560        let route_due = self.route_ready_at().is_some_and(|ready| now.0 >= ready.0);
561        if ratchet_due {
562            if let Some(index) = (!self.pending_ratchets.is_empty()).then_some(0) {
563                self.append_ratchet(engine, index, now).await;
564                return;
565            }
566            if self.snapshot_required
567                && self.snapshot_target == EmbeddedPersistenceTarget::CriticalState
568            {
569                self.try_start_compaction(engine, now);
570                return;
571            }
572        }
573        if !route_due {
574            return;
575        }
576        if self.snapshot_required {
577            self.try_start_compaction(engine, now);
578            return;
579        }
580        if !self.pending_routes.is_empty() {
581            self.append_route(engine, 0, now).await;
582        }
583    }
584
585    async fn append_route<S: StorageLayout>(
586        &mut self,
587        engine: &EngineState<S>,
588        index: usize,
589        now: InstantMillis,
590    ) {
591        let delta = self.pending_routes[index];
592        let encoded = encode_route_delta(engine, delta);
593        let Ok(payload) = encoded else {
594            self.note_codec_failure(now);
595            return;
596        };
597        let can_fit = self.journal.as_ref().is_some_and(|journal| {
598            journal.active_can_fit(payload.len, self.policy.compaction.critical_reserve_bytes)
599        });
600        if !can_fit {
601            self.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
602            self.try_start_compaction(engine, now);
603            return;
604        }
605        let Some(journal) = self.journal.as_mut() else {
606            return;
607        };
608        let result = journal
609            .append(payload.kind, &payload.payload[..payload.len])
610            .await;
611        match result {
612            Ok(()) => {
613                self.pending_routes.swap_remove(index);
614                self.landing_records = self.landing_records.saturating_add(1);
615                if self.pending_routes.is_empty() {
616                    self.landing_batch = Some(BatchKind::Routes);
617                }
618            }
619            Err(FlashJournalError::ArenaFull) => {
620                self.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
621                self.try_start_compaction(engine, now);
622            }
623            Err(error) => {
624                let failure = failure_from_journal(error);
625                self.note_write_failure(now, failure);
626            }
627        }
628    }
629
630    async fn append_ratchet<S: StorageLayout>(
631        &mut self,
632        engine: &EngineState<S>,
633        index: usize,
634        now: InstantMillis,
635    ) {
636        let destination = self.pending_ratchets[index];
637        let Ok(payload) = encode_ratchet(engine, destination) else {
638            self.note_codec_failure(now);
639            return;
640        };
641        let Some(journal) = self.journal.as_mut() else {
642            return;
643        };
644        let result = journal
645            .append(payload.kind, &payload.payload[..payload.len])
646            .await;
647        match result {
648            Ok(()) => {
649                self.pending_ratchets.swap_remove(index);
650                self.landing_records = self.landing_records.saturating_add(1);
651                if self.pending_ratchets.is_empty() {
652                    self.landing_batch = Some(BatchKind::Ratchets);
653                }
654            }
655            Err(FlashJournalError::ArenaFull) => {
656                self.require_snapshot(EmbeddedPersistenceTarget::CriticalState, now);
657                self.try_start_compaction(engine, now);
658            }
659            Err(error) => self.note_write_failure(now, failure_from_journal(error)),
660        }
661    }
662
663    fn try_start_compaction<S: StorageLayout>(
664        &mut self,
665        engine: &EngineState<S>,
666        now: InstantMillis,
667    ) {
668        let allowed = self.next_compaction_not_before.unwrap_or(now);
669        if now.0 < allowed.0 {
670            self.require_snapshot(self.snapshot_target, now);
671            return;
672        }
673        self.compaction_route_keys.clear();
674        let mut route_capacity_failed = false;
675        for destination in engine.persisted_route_destinations() {
676            if self.compaction_route_keys.push(destination).is_err() {
677                route_capacity_failed = true;
678                break;
679            }
680        }
681        if route_capacity_failed {
682            self.note_write_failure(now, EmbeddedPersistenceFailure::Capacity);
683            return;
684        }
685        self.compaction_ratchet_keys.clear();
686        let mut ratchet_capacity_failed = false;
687        for (destination, _, _) in engine.persisted_self_ratchet_rows() {
688            if self.compaction_ratchet_keys.push(destination).is_err() {
689                ratchet_capacity_failed = true;
690                break;
691            }
692        }
693        if ratchet_capacity_failed {
694            self.note_write_failure(now, EmbeddedPersistenceFailure::Capacity);
695            return;
696        }
697        let target = self.snapshot_target;
698        self.pending_routes.clear();
699        self.pending_ratchets.clear();
700        self.route_dirty_since = None;
701        self.ratchet_dirty_since = None;
702        self.snapshot_required = false;
703        self.snapshot_target = EmbeddedPersistenceTarget::Routes;
704        self.compaction_target = Some(target);
705        self.compaction = Some(CompactionPhase::RecordBudget { at: now });
706        self.landing_batch = None;
707        self.landing_records = 0;
708    }
709
710    async fn progress_compaction<S: StorageLayout>(
711        &mut self,
712        engine: &mut EngineState<S>,
713        now: InstantMillis,
714    ) {
715        let Some(phase) = self.compaction else {
716            return;
717        };
718        let Some(journal) = self.journal.as_mut() else {
719            return;
720        };
721        match phase {
722            CompactionPhase::RecordBudget { at } => {
723                match journal.record_compaction_budget(at).await {
724                    Ok(recorded_at) => {
725                        self.last_timebase_success = Some(at);
726                        let next_allowed_at = InstantMillis(
727                            recorded_at
728                                .0
729                                .saturating_add(self.policy.compaction.minimum_interval_millis),
730                        );
731                        self.next_compaction_not_before = Some(next_allowed_at);
732                        (self.observe_diagnostic)(
733                            EmbeddedPersistenceDiagnostic::CompactionStarted {
734                                at,
735                                next_allowed_at,
736                            },
737                        );
738                        self.compaction = Some(CompactionPhase::Erase { sector: 0 });
739                    }
740                    Err(error) => {
741                        self.note_write_failure(now, failure_from_journal(error));
742                    }
743                }
744            }
745            CompactionPhase::Erase { sector } => {
746                if sector < journal.inactive_sector_count() {
747                    match journal.erase_inactive_sector(sector).await {
748                        Ok(()) => {
749                            let next = sector + 1;
750                            if next == journal.inactive_sector_count() {
751                                if journal.begin_compaction().is_err() {
752                                    self.note_write_failure(
753                                        now,
754                                        EmbeddedPersistenceFailure::Capacity,
755                                    );
756                                    return;
757                                }
758                                self.compaction = Some(CompactionPhase::Routes { index: 0 });
759                            } else {
760                                self.compaction = Some(CompactionPhase::Erase { sector: next });
761                            }
762                        }
763                        Err(error) => {
764                            self.note_write_failure(now, failure_from_journal(error));
765                        }
766                    }
767                }
768            }
769            CompactionPhase::Routes { index } => {
770                let mut scratch = [0u8; RECORD_SCRATCH_LEN];
771                let Some(destination) = self.compaction_route_keys.get(index) else {
772                    self.compaction = Some(CompactionPhase::Ratchets { index: 0 });
773                    return;
774                };
775                let Some(row) = engine.persisted_route_row(&destination) else {
776                    self.compaction = Some(CompactionPhase::Routes { index: index + 1 });
777                    return;
778                };
779                let mut durable = row.clone();
780                durable.announce_id_ring = AnnounceIdRing::Table(&[]);
781                let required = routing_table_snapshot_len(core::iter::once(durable.clone()));
782                if required > scratch.len() {
783                    self.note_codec_failure(now);
784                    return;
785                }
786                let Ok(written) = write_routing_table_snapshot(
787                    core::iter::once(durable),
788                    &mut scratch[..required],
789                ) else {
790                    self.note_codec_failure(now);
791                    return;
792                };
793                match journal
794                    .append_compacted(FlashJournalRecordKind::RouteUpsert, &scratch[..written])
795                    .await
796                {
797                    Ok(()) => {
798                        self.landing_records = self.landing_records.saturating_add(1);
799                        self.compaction = Some(CompactionPhase::Routes { index: index + 1 });
800                    }
801                    Err(error) => {
802                        self.note_write_failure(now, failure_from_journal(error));
803                    }
804                }
805            }
806            CompactionPhase::Ratchets { index } => {
807                let Some(destination) = self.compaction_ratchet_keys.get(index).copied() else {
808                    self.compaction = Some(CompactionPhase::Commit);
809                    return;
810                };
811                let Some((last_rotated, secrets)) = engine.persisted_self_ratchet_row(&destination)
812                else {
813                    self.compaction = Some(CompactionPhase::Ratchets { index: index + 1 });
814                    return;
815                };
816                let mut scratch = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
817                scratch[..TRUNCATED_HASH_BYTE_LEN].copy_from_slice(destination.as_bytes());
818                let required = self_ratchets_snapshot_len(secrets.len());
819                let end = TRUNCATED_HASH_BYTE_LEN.saturating_add(required);
820                if end > scratch.len() {
821                    self.note_codec_failure(now);
822                    return;
823                }
824                let Ok(written) = write_self_ratchets_snapshot(
825                    last_rotated,
826                    secrets,
827                    &mut scratch[TRUNCATED_HASH_BYTE_LEN..end],
828                ) else {
829                    self.note_codec_failure(now);
830                    return;
831                };
832                match journal
833                    .append_compacted(
834                        FlashJournalRecordKind::SelfRatchet,
835                        &scratch[..TRUNCATED_HASH_BYTE_LEN + written],
836                    )
837                    .await
838                {
839                    Ok(()) => {
840                        self.landing_records = self.landing_records.saturating_add(1);
841                        self.compaction = Some(CompactionPhase::Ratchets { index: index + 1 });
842                    }
843                    Err(error) => {
844                        self.note_write_failure(now, failure_from_journal(error));
845                    }
846                }
847            }
848            CompactionPhase::Commit => match journal.commit_compaction().await {
849                Ok(()) => {
850                    self.compaction = None;
851                    self.compaction_target = None;
852                    let records = core::mem::take(&mut self.landing_records);
853                    self.retry_not_before = None;
854                    self.write_failed = false;
855                    if self.snapshot_required {
856                        self.require_snapshot(self.snapshot_target, now);
857                    } else {
858                        self.deferred_target = None;
859                        self.deferred_until = None;
860                    }
861                    let state_not_saved = self.state_not_saved();
862                    (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::CompactionCompleted {
863                        records,
864                        at: now,
865                        state_not_saved,
866                    });
867                    if !self.snapshot_required
868                        && self.pending_routes.is_empty()
869                        && self.pending_ratchets.is_empty()
870                    {
871                        self.landing_batch = Some(BatchKind::Compaction);
872                    }
873                }
874                Err(error) => {
875                    self.note_write_failure(now, failure_from_journal(error));
876                }
877            },
878        }
879    }
880
881    async fn land_timebase(&mut self, batch: BatchKind, now: InstantMillis) {
882        let should_record = self.last_timebase_success.is_none_or(|last| {
883            now.0.saturating_sub(last.0) >= self.policy.timebase_record_interval_millis
884        });
885        if should_record {
886            let Some(journal) = self.journal.as_mut() else {
887                return;
888            };
889            if let Err(error) = journal.record_timebase(now).await {
890                self.note_write_failure(now, failure_from_journal(error));
891                return;
892            }
893            self.last_timebase_success = Some(now);
894        }
895        match batch {
896            BatchKind::Routes => {
897                self.route_dirty_since = None;
898                self.last_route_success = Some(now);
899            }
900            BatchKind::Ratchets => {
901                self.ratchet_dirty_since = None;
902            }
903            BatchKind::Compaction => {
904                self.route_dirty_since = None;
905                self.ratchet_dirty_since = None;
906                self.last_route_success = Some(now);
907            }
908        }
909        self.retry_not_before = None;
910        self.write_failed = false;
911        let records = core::mem::take(&mut self.landing_records);
912        self.landing_batch = None;
913        if batch != BatchKind::Compaction {
914            let state_not_saved = self.state_not_saved();
915            (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::BatchPersisted {
916                records,
917                at: now,
918                state_not_saved,
919            });
920        }
921    }
922
923    fn note_codec_failure(&mut self, now: InstantMillis) {
924        self.note_write_failure(now, EmbeddedPersistenceFailure::Codec);
925    }
926
927    fn note_write_failure(&mut self, now: InstantMillis, failure: EmbeddedPersistenceFailure) {
928        let retry_at = InstantMillis(now.0.saturating_add(self.policy.retry_interval_millis));
929        self.retry_not_before = Some(retry_at);
930        self.write_failed = true;
931        if self.compaction.is_some() {
932            let target = self
933                .compaction_target
934                .unwrap_or(EmbeddedPersistenceTarget::Routes);
935            if let Some(journal) = self.journal.as_mut() {
936                journal.abort_compaction();
937            }
938            self.compaction = None;
939            self.compaction_target = None;
940            self.require_snapshot(target, now);
941            match target {
942                EmbeddedPersistenceTarget::Routes => {
943                    self.route_dirty_since.get_or_insert(now);
944                }
945                EmbeddedPersistenceTarget::CriticalState => {
946                    self.ratchet_dirty_since.get_or_insert(now);
947                }
948            }
949        }
950        (self.observe_diagnostic)(EmbeddedPersistenceDiagnostic::WriteFailed { failure, retry_at });
951    }
952}
953
954pub(crate) trait ManifoldPersistence<S: StorageLayout> {
955    fn observe(&mut self, journaled: &Journaled<'_>, now: InstantMillis);
956    fn deadline(&self, now: InstantMillis) -> Option<InstantMillis>;
957    async fn progress(&mut self, engine: &mut EngineState<S>, now: InstantMillis);
958}
959
960impl<S, F, Keys, Observe, const PENDING: usize> ManifoldPersistence<S>
961    for EmbeddedFlashPersistence<F, Keys, Observe, PENDING>
962where
963    S: StorageLayout,
964    F: NorFlash,
965    Keys: RouteSnapshotKeys,
966    Observe: FnMut(EmbeddedPersistenceDiagnostic),
967{
968    fn observe(&mut self, journaled: &Journaled<'_>, now: InstantMillis) {
969        self.observe_journaled(journaled, now);
970    }
971
972    fn deadline(&self, now: InstantMillis) -> Option<InstantMillis> {
973        self.next_deadline(now)
974    }
975
976    async fn progress(&mut self, engine: &mut EngineState<S>, now: InstantMillis) {
977        self.progress(engine, now).await;
978    }
979}
980
981pub(crate) struct NoManifoldPersistence;
982
983impl<S: StorageLayout> ManifoldPersistence<S> for NoManifoldPersistence {
984    fn observe(&mut self, _journaled: &Journaled<'_>, _now: InstantMillis) {}
985
986    fn deadline(&self, _now: InstantMillis) -> Option<InstantMillis> {
987        None
988    }
989
990    async fn progress(&mut self, _engine: &mut EngineState<S>, _now: InstantMillis) {}
991}
992
993fn encode_route_delta<S: StorageLayout>(
994    engine: &EngineState<S>,
995    delta: PendingRouteDelta,
996) -> Result<EncodedDelta, ()> {
997    match delta {
998        PendingRouteDelta::RouteUpsert(destination) => {
999            let Some(row) = engine.persisted_route_row(&destination) else {
1000                return encode_tombstone(destination);
1001            };
1002            let mut durable = row.clone();
1003            durable.announce_id_ring = AnnounceIdRing::Table(&[]);
1004            let required = routing_table_snapshot_len(core::iter::once(durable.clone()));
1005            if required > RECORD_SCRATCH_LEN {
1006                return Err(());
1007            }
1008            let mut scratch = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
1009            let written =
1010                write_routing_table_snapshot(core::iter::once(durable), &mut scratch[..required])
1011                    .map_err(|_| ())?;
1012            Ok(EncodedDelta {
1013                kind: FlashJournalRecordKind::RouteUpsert,
1014                payload: scratch,
1015                len: written,
1016            })
1017        }
1018        PendingRouteDelta::RouteRemoval(destination) => encode_tombstone(destination),
1019    }
1020}
1021
1022fn encode_ratchet<S: StorageLayout>(
1023    engine: &EngineState<S>,
1024    destination: DestinationHash,
1025) -> Result<EncodedDelta, ()> {
1026    let Some((last_rotated, secrets)) = engine.persisted_self_ratchet_row(&destination) else {
1027        return Err(());
1028    };
1029    let required = self_ratchets_snapshot_len(secrets.len());
1030    if TRUNCATED_HASH_BYTE_LEN + required > RECORD_SCRATCH_LEN {
1031        return Err(());
1032    }
1033    let mut scratch = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
1034    scratch[..TRUNCATED_HASH_BYTE_LEN].copy_from_slice(destination.as_bytes());
1035    let written = write_self_ratchets_snapshot(
1036        last_rotated,
1037        secrets,
1038        &mut scratch[TRUNCATED_HASH_BYTE_LEN..TRUNCATED_HASH_BYTE_LEN + required],
1039    )
1040    .map_err(|_| ())?;
1041    Ok(EncodedDelta {
1042        kind: FlashJournalRecordKind::SelfRatchet,
1043        payload: scratch,
1044        len: TRUNCATED_HASH_BYTE_LEN + written,
1045    })
1046}
1047
1048fn encode_tombstone(destination: DestinationHash) -> Result<EncodedDelta, ()> {
1049    let mut payload = Zeroizing::new([0u8; RECORD_SCRATCH_LEN]);
1050    payload[..TRUNCATED_HASH_BYTE_LEN].copy_from_slice(destination.as_bytes());
1051    Ok(EncodedDelta {
1052        kind: FlashJournalRecordKind::RouteRemoval,
1053        payload,
1054        len: TRUNCATED_HASH_BYTE_LEN,
1055    })
1056}
1057
1058fn apply_record<S: StorageLayout>(
1059    engine: &mut EngineState<S>,
1060    now: InstantMillis,
1061    record: FlashJournalRecord<'_>,
1062    report: &mut EmbeddedPersistenceRestoreReport,
1063) {
1064    match record.kind {
1065        FlashJournalRecordKind::ArenaCommit => {}
1066        FlashJournalRecordKind::RouteUpsert => {
1067            let Ok(mut rows) = read_routing_table_snapshot(record.payload) else {
1068                report.route_refused_count = report.route_refused_count.saturating_add(1);
1069                return;
1070            };
1071            let Some(Ok(row)) = rows.next() else {
1072                report.route_refused_count = report.route_refused_count.saturating_add(1);
1073                return;
1074            };
1075            if rows.next().is_some() {
1076                report.route_refused_count = report.route_refused_count.saturating_add(1);
1077                return;
1078            }
1079            let destination = row.destination;
1080            let Ok(pending) = engine.prepare_persisted_route(row) else {
1081                report.route_refused_count = report.route_refused_count.saturating_add(1);
1082                return;
1083            };
1084            let Ok(verified) = pending.verify() else {
1085                report.route_refused_count = report.route_refused_count.saturating_add(1);
1086                return;
1087            };
1088            let _ = engine.drop_route(&destination, AttachedInterfaces::new(&[]));
1089            match engine.seed_verified_route(verified, now) {
1090                RouteSeedOutcome::Seeded => {
1091                    report.route_seeded_count = report.route_seeded_count.saturating_add(1);
1092                }
1093                RouteSeedOutcome::RefusedDestinationMismatch
1094                | RouteSeedOutcome::RefusedBlackholedIdentity
1095                | RouteSeedOutcome::RefusedInvalidSignature => {
1096                    report.route_refused_count = report.route_refused_count.saturating_add(1);
1097                }
1098                RouteSeedOutcome::AlreadyPresent
1099                | RouteSeedOutcome::TableFull
1100                | RouteSeedOutcome::AppDataArenaFull => {
1101                    report.route_dropped_count = report.route_dropped_count.saturating_add(1);
1102                }
1103            }
1104        }
1105        FlashJournalRecordKind::RouteRemoval => {
1106            let Ok(bytes) = <[u8; TRUNCATED_HASH_BYTE_LEN]>::try_from(record.payload) else {
1107                report.route_refused_count = report.route_refused_count.saturating_add(1);
1108                return;
1109            };
1110            let destination = DestinationHash::new(bytes);
1111            let _ = engine.drop_route(&destination, AttachedInterfaces::new(&[]));
1112        }
1113        FlashJournalRecordKind::SelfRatchet => {
1114            let Some((destination, sealed)) = record
1115                .payload
1116                .split_first_chunk::<TRUNCATED_HASH_BYTE_LEN>()
1117            else {
1118                report.ratchet_refused_count = report.ratchet_refused_count.saturating_add(1);
1119                return;
1120            };
1121            let Ok(restored) = read_self_ratchets_snapshot(sealed) else {
1122                report.ratchet_refused_count = report.ratchet_refused_count.saturating_add(1);
1123                return;
1124            };
1125            match engine.replace_persisted_self_ratchets(
1126                &DestinationHash::new(*destination),
1127                restored.last_rotated,
1128                restored.secrets_newest_first(),
1129            ) {
1130                SeedSelfRatchetsOutcome::Seeded => {
1131                    report.ratchet_seeded_count = report.ratchet_seeded_count.saturating_add(1);
1132                }
1133                SeedSelfRatchetsOutcome::AlreadyMinted | SeedSelfRatchetsOutcome::Untracked => {
1134                    report.ratchet_refused_count = report.ratchet_refused_count.saturating_add(1);
1135                }
1136            }
1137        }
1138    }
1139}
1140
1141fn failure_from_journal<E>(error: FlashJournalError<E>) -> EmbeddedPersistenceFailure {
1142    match error {
1143        FlashJournalError::Flash(_) => EmbeddedPersistenceFailure::Flash,
1144        FlashJournalError::ArenaFull
1145        | FlashJournalError::OutOfBounds
1146        | FlashJournalError::Misaligned
1147        | FlashJournalError::Uninitialized
1148        | FlashJournalError::CompactionInProgress
1149        | FlashJournalError::NoCompaction
1150        | FlashJournalError::PayloadTooLarge
1151        | FlashJournalError::ScratchTooShort => EmbeddedPersistenceFailure::Capacity,
1152    }
1153}
1154
1155fn earlier(first: Option<InstantMillis>, second: Option<InstantMillis>) -> Option<InstantMillis> {
1156    match (first, second) {
1157        (Some(first), Some(second)) => Some(InstantMillis(first.0.min(second.0))),
1158        (Some(value), None) | (None, Some(value)) => Some(value),
1159        (None, None) => None,
1160    }
1161}
1162
1163#[cfg(test)]
1164mod tests {
1165    use super::*;
1166    use embedded_storage::nor_flash::{ErrorType, NorFlashError, NorFlashErrorKind};
1167    use embedded_storage_async::nor_flash::ReadNorFlash;
1168    use std::cell::{Cell, RefCell};
1169    use std::rc::Rc;
1170    use std::vec::Vec;
1171
1172    const ERASE: usize = 512;
1173    const CAPACITY: usize = ERASE * 6;
1174    const LAYOUT: FlashJournalLayout = FlashJournalLayout::new(
1175        [0, ERASE as u32],
1176        [
1177            crate::persistence::FlashArenaRange::new((ERASE * 2) as u32, (ERASE * 4) as u32),
1178            crate::persistence::FlashArenaRange::new((ERASE * 4) as u32, (ERASE * 6) as u32),
1179        ],
1180    );
1181
1182    #[derive(Debug)]
1183    struct TestFlash {
1184        bytes: [u8; CAPACITY],
1185        sector_erases: [u32; CAPACITY / ERASE],
1186        fail_next_write: Rc<Cell<bool>>,
1187    }
1188
1189    impl TestFlash {
1190        fn new() -> Self {
1191            Self {
1192                bytes: [0xFF; CAPACITY],
1193                sector_erases: [0; CAPACITY / ERASE],
1194                fail_next_write: Rc::new(Cell::new(false)),
1195            }
1196        }
1197
1198        fn controlled() -> (Self, Rc<Cell<bool>>) {
1199            let flash = Self::new();
1200            let control = Rc::clone(&flash.fail_next_write);
1201            (flash, control)
1202        }
1203    }
1204
1205    #[derive(Debug)]
1206    struct TestFlashError;
1207
1208    impl NorFlashError for TestFlashError {
1209        fn kind(&self) -> NorFlashErrorKind {
1210            NorFlashErrorKind::Other
1211        }
1212    }
1213
1214    impl ErrorType for TestFlash {
1215        type Error = TestFlashError;
1216    }
1217
1218    impl ReadNorFlash for TestFlash {
1219        const READ_SIZE: usize = 4;
1220
1221        async fn read(&mut self, offset: u32, bytes: &mut [u8]) -> Result<(), Self::Error> {
1222            let start = offset as usize;
1223            let end = start + bytes.len();
1224            bytes.copy_from_slice(&self.bytes[start..end]);
1225            Ok(())
1226        }
1227
1228        fn capacity(&self) -> usize {
1229            CAPACITY
1230        }
1231    }
1232
1233    impl NorFlash for TestFlash {
1234        const WRITE_SIZE: usize = 4;
1235        const ERASE_SIZE: usize = ERASE;
1236
1237        async fn write(&mut self, offset: u32, bytes: &[u8]) -> Result<(), Self::Error> {
1238            if self.fail_next_write.replace(false) {
1239                return Err(TestFlashError);
1240            }
1241            let start = offset as usize;
1242            for (stored, written) in self.bytes[start..start + bytes.len()].iter_mut().zip(bytes) {
1243                *stored &= *written;
1244            }
1245            Ok(())
1246        }
1247
1248        async fn erase(&mut self, from: u32, to: u32) -> Result<(), Self::Error> {
1249            self.bytes[from as usize..to as usize].fill(0xFF);
1250            for sector in from as usize / ERASE..to as usize / ERASE {
1251                self.sector_erases[sector] = self.sector_erases[sector].saturating_add(1);
1252            }
1253            Ok(())
1254        }
1255    }
1256
1257    fn ready_with_observer<Observe>(
1258        observe: Observe,
1259    ) -> EmbeddedFlashPersistence<TestFlash, FixedRouteSnapshotKeys<8>, Observe, 4>
1260    where
1261        Observe: FnMut(EmbeddedPersistenceDiagnostic),
1262    {
1263        embassy_futures::block_on(async {
1264            let flash = TestFlash::new();
1265            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1266            let (mut journal, _) = FlashJournal::open(flash, LAYOUT, &mut scratch, |_| {})
1267                .await
1268                .unwrap();
1269            journal.initialize_empty().await.unwrap();
1270            journal
1271                .record_compaction_budget(InstantMillis(0))
1272                .await
1273                .unwrap();
1274            let mut persistence = EmbeddedFlashPersistence::new(
1275                TestFlash::new(),
1276                LAYOUT,
1277                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1278                FixedRouteSnapshotKeys::new(),
1279                observe,
1280            );
1281            persistence.flash = None;
1282            persistence.journal = Some(journal);
1283            persistence.next_compaction_not_before =
1284                Some(InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS));
1285            persistence
1286        })
1287    }
1288
1289    fn ready() -> EmbeddedFlashPersistence<
1290        TestFlash,
1291        FixedRouteSnapshotKeys<8>,
1292        fn(EmbeddedPersistenceDiagnostic),
1293        4,
1294    > {
1295        ready_with_observer((|_| {}) as fn(EmbeddedPersistenceDiagnostic))
1296    }
1297
1298    fn signed_route(secret: u8, app_data: &[u8]) -> crate::routing::PersistedRouteRow<'_> {
1299        use crate::identity::in_memory::InMemoryNodeIdentity;
1300        use crate::interfaces::InterfaceId;
1301        use crate::routing::announce::{Announce, AnnounceId, DottedNameHash};
1302        use crate::routing::routes::RouteEntry;
1303        use crate::routing::{AnnounceIdRing, NextHop, RouteResponsiveness};
1304
1305        let signer = InMemoryNodeIdentity::from_secret_key_bytes(&[secret; 64]);
1306        let announce = Announce::build_signed(
1307            &signer,
1308            DottedNameHash::new([secret; 10]),
1309            AnnounceId::from_wire([secret.wrapping_add(1); 10]),
1310            None,
1311            app_data,
1312        )
1313        .unwrap();
1314        crate::routing::PersistedRouteRow {
1315            destination: announce.destination,
1316            entry: RouteEntry {
1317                hops: secret,
1318                learned_at: InstantMillis(500),
1319                last_relayed_at: InstantMillis(700),
1320                responsiveness: RouteResponsiveness::Responsive,
1321                receiving_interface: InterfaceId::new([secret; 8]),
1322                next_hop: NextHop::Direct,
1323            },
1324            public_keys: announce.public_keys,
1325            dotted_name_hash: announce.dotted_name_hash,
1326            announce_id: announce.announce_id,
1327            ratchet: announce.ratchet,
1328            signature: announce.signature,
1329            app_data,
1330            announce_id_ring: AnnounceIdRing::Wire(
1331                &[0; crate::routing::announce::ANNOUNCE_ID_WIRE_LEN],
1332            ),
1333        }
1334    }
1335
1336    #[test]
1337    fn exact_route_and_ratchet_deadlines_are_distinct() {
1338        let policy =
1339            EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(64));
1340        assert_eq!(policy.first_route_commit_delay_millis, 2_000);
1341        assert_eq!(policy.minimum_route_commit_interval_millis, 300_000);
1342        assert_eq!(policy.ratchet_batch_delay_millis, 2_000);
1343        assert_eq!(policy.retry_interval_millis, 300_000);
1344        assert_eq!(
1345            policy.timebase_record_interval_millis,
1346            TIMEBASE_RECORD_INTERVAL_MILLIS
1347        );
1348        assert_eq!(
1349            policy.compaction.minimum_interval_millis,
1350            HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS
1351        );
1352        assert_eq!(policy.compaction.critical_reserve_bytes, 64);
1353    }
1354
1355    #[test]
1356    fn deadline_formula_batches_first_write_then_honors_five_minutes() {
1357        let mut persistence = ready();
1358        persistence.route_dirty_since = Some(InstantMillis(1_000));
1359        assert_eq!(
1360            persistence.next_deadline(InstantMillis(1_500)),
1361            Some(InstantMillis(3_000))
1362        );
1363        persistence.last_route_success = Some(InstantMillis(10_000));
1364        persistence.route_dirty_since = Some(InstantMillis(11_000));
1365        assert_eq!(
1366            persistence.next_deadline(InstantMillis(12_000)),
1367            Some(InstantMillis(310_000))
1368        );
1369        persistence.ratchet_dirty_since = Some(InstantMillis(12_000));
1370        assert_eq!(
1371            persistence.next_deadline(InstantMillis(12_000)),
1372            Some(InstantMillis(14_000))
1373        );
1374        persistence.retry_not_before = Some(InstantMillis(400_000));
1375        assert_eq!(
1376            persistence.next_deadline(InstantMillis(12_000)),
1377            Some(InstantMillis(400_000))
1378        );
1379    }
1380
1381    #[test]
1382    fn repeated_route_and_ratchet_changes_coalesce_by_destination() {
1383        let mut persistence = ready();
1384        let destination = DestinationHash::new([0x11; TRUNCATED_HASH_BYTE_LEN]);
1385        persistence.queue_route(
1386            PendingRouteDelta::RouteUpsert(destination),
1387            InstantMillis(0),
1388        );
1389        persistence.queue_route(
1390            PendingRouteDelta::RouteUpsert(destination),
1391            InstantMillis(0),
1392        );
1393        persistence.queue_route(
1394            PendingRouteDelta::RouteRemoval(destination),
1395            InstantMillis(0),
1396        );
1397        persistence.queue_ratchet(destination, InstantMillis(0));
1398        persistence.queue_ratchet(destination, InstantMillis(0));
1399        assert_eq!(
1400            persistence.pending_routes.as_slice(),
1401            &[PendingRouteDelta::RouteRemoval(destination)]
1402        );
1403        assert_eq!(persistence.pending_ratchets.as_slice(), &[destination]);
1404    }
1405
1406    #[test]
1407    fn pending_overflow_waits_for_the_batch_deadline_before_compacting() {
1408        let mut persistence = ready();
1409        for byte in 0..5 {
1410            persistence.queue_route(
1411                PendingRouteDelta::RouteUpsert(DestinationHash::new(
1412                    [byte; TRUNCATED_HASH_BYTE_LEN],
1413                )),
1414                InstantMillis(1_000),
1415            );
1416        }
1417        persistence.route_dirty_since = Some(InstantMillis(1_000));
1418        assert!(persistence.snapshot_required);
1419        assert_eq!(
1420            persistence.deferred_target,
1421            Some(EmbeddedPersistenceTarget::Routes)
1422        );
1423        assert_eq!(
1424            persistence.next_deadline(InstantMillis(1_500)),
1425            Some(InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS))
1426        );
1427
1428        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1429        embassy_futures::block_on(persistence.progress(&mut engine, InstantMillis(2_999)));
1430        assert_eq!(persistence.compaction, None);
1431        assert!(persistence.snapshot_required);
1432
1433        embassy_futures::block_on(persistence.progress(
1434            &mut engine,
1435            InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS),
1436        ));
1437        assert_eq!(
1438            persistence.compaction,
1439            Some(CompactionPhase::RecordBudget {
1440                at: InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS)
1441            })
1442        );
1443    }
1444
1445    #[test]
1446    fn failures_keep_dirty_state_and_raise_the_notice() {
1447        let mut persistence = ready();
1448        let destination = DestinationHash::new([0x22; TRUNCATED_HASH_BYTE_LEN]);
1449        persistence.queue_route(
1450            PendingRouteDelta::RouteUpsert(destination),
1451            InstantMillis(0),
1452        );
1453        persistence.note_codec_failure(InstantMillis(1_000));
1454        assert_eq!(persistence.pending_routes.len(), 1);
1455        assert_eq!(persistence.retry_not_before, Some(InstantMillis(301_000)));
1456        assert!(persistence.state_not_saved());
1457        assert!(!persistence.snapshot_required);
1458
1459        persistence.retry_not_before = None;
1460        persistence.compaction = Some(CompactionPhase::Commit);
1461        persistence.compaction_target = Some(EmbeddedPersistenceTarget::Routes);
1462        persistence.note_write_failure(InstantMillis(2_000), EmbeddedPersistenceFailure::Flash);
1463        assert_eq!(persistence.pending_routes.len(), 0);
1464        assert_eq!(persistence.retry_not_before, Some(InstantMillis(302_000)));
1465        assert!(persistence.state_not_saved());
1466        assert!(persistence.snapshot_required);
1467        assert_eq!(persistence.compaction, None);
1468    }
1469
1470    #[test]
1471    fn legacy_timebase_allows_one_needed_compaction_then_adopts_the_budget_marker() {
1472        embassy_futures::block_on(async {
1473            let (mut journal, _) = {
1474                let flash = TestFlash::new();
1475                let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1476                FlashJournal::open(flash, LAYOUT, &mut scratch, |_| {})
1477                    .await
1478                    .unwrap()
1479            };
1480            journal.initialize_empty().await.unwrap();
1481            journal
1482                .record_timebase(InstantMillis(10_000))
1483                .await
1484                .unwrap();
1485            let mut persistence =
1486                EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1487                    journal.release(),
1488                    LAYOUT,
1489                    EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(
1490                        0,
1491                    )),
1492                    FixedRouteSnapshotKeys::new(),
1493                    (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1494                );
1495            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1496            let report = persistence.restore(&mut engine, InstantMillis(0)).await;
1497            assert_eq!(persistence.next_compaction_not_before, None);
1498            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, report.logical_start);
1499            persistence.route_dirty_since = Some(InstantMillis(report.logical_start.0 - 2_000));
1500            persistence
1501                .progress(&mut engine, report.logical_start)
1502                .await;
1503            persistence
1504                .progress(&mut engine, report.logical_start)
1505                .await;
1506            assert!(matches!(
1507                persistence.compaction,
1508                Some(CompactionPhase::Erase { sector: 0 })
1509            ));
1510
1511            let flash = persistence.journal.take().unwrap().release();
1512            let mut restored = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1513                flash,
1514                LAYOUT,
1515                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1516                FixedRouteSnapshotKeys::new(),
1517                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1518            );
1519            restored.restore(&mut engine, InstantMillis(0)).await;
1520            assert!(restored.next_compaction_not_before.is_some());
1521        });
1522    }
1523
1524    #[test]
1525    fn failed_compaction_attempt_consumes_the_daily_budget() {
1526        let diagnostics = Rc::new(RefCell::new(Vec::new()));
1527        let observed = Rc::clone(&diagnostics);
1528        let mut persistence = ready_with_observer(move |diagnostic| {
1529            observed.borrow_mut().push(diagnostic);
1530        });
1531        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1532        let first = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1533        persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, first);
1534        persistence.route_dirty_since = Some(InstantMillis(first.0 - 2_000));
1535        embassy_futures::block_on(persistence.progress(&mut engine, first));
1536        let second = InstantMillis(first.0 + HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1537        assert_eq!(persistence.next_compaction_not_before, Some(first));
1538        embassy_futures::block_on(persistence.progress(&mut engine, first));
1539        assert_eq!(persistence.next_compaction_not_before, Some(second));
1540        assert_eq!(
1541            persistence.compaction,
1542            Some(CompactionPhase::Erase { sector: 0 })
1543        );
1544
1545        persistence.note_write_failure(first, EmbeddedPersistenceFailure::Flash);
1546        assert_eq!(persistence.compaction, None);
1547        assert!(persistence.snapshot_required);
1548        assert_eq!(persistence.next_deadline(first), Some(second));
1549        embassy_futures::block_on(persistence.progress(&mut engine, InstantMillis(second.0 - 1)));
1550        assert_eq!(persistence.compaction, None);
1551        embassy_futures::block_on(persistence.progress(&mut engine, second));
1552        assert_eq!(
1553            persistence.compaction,
1554            Some(CompactionPhase::RecordBudget { at: second })
1555        );
1556        embassy_futures::block_on(persistence.progress(&mut engine, second));
1557        assert_eq!(
1558            persistence.compaction,
1559            Some(CompactionPhase::Erase { sector: 0 })
1560        );
1561
1562        let starts = diagnostics
1563            .borrow()
1564            .iter()
1565            .filter(|diagnostic| {
1566                matches!(
1567                    diagnostic,
1568                    EmbeddedPersistenceDiagnostic::CompactionStarted { .. }
1569                )
1570            })
1571            .count();
1572        assert_eq!(starts, 2);
1573    }
1574
1575    #[test]
1576    fn recorded_compaction_budget_survives_reboot() {
1577        embassy_futures::block_on(async {
1578            let mut persistence = ready();
1579            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1580            let attempt = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1581            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, attempt);
1582            persistence.route_dirty_since = Some(InstantMillis(attempt.0 - 2_000));
1583            persistence.progress(&mut engine, attempt).await;
1584            persistence.progress(&mut engine, attempt).await;
1585            assert_eq!(
1586                persistence.compaction,
1587                Some(CompactionPhase::Erase { sector: 0 })
1588            );
1589
1590            let flash = persistence.journal.take().unwrap().release();
1591            let mut restored = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1592                flash,
1593                LAYOUT,
1594                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1595                FixedRouteSnapshotKeys::new(),
1596                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1597            );
1598            let report = restored.restore(&mut engine, InstantMillis(0)).await;
1599            assert!(report.logical_start.0 >= attempt.0);
1600            assert_eq!(
1601                restored.next_compaction_not_before,
1602                Some(InstantMillis(
1603                    attempt
1604                        .0
1605                        .saturating_add(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS)
1606                ))
1607            );
1608        });
1609    }
1610
1611    #[test]
1612    fn marker_write_failure_does_not_consume_the_compaction_budget() {
1613        embassy_futures::block_on(async {
1614            let (flash, fail_next_write) = TestFlash::controlled();
1615            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1616            let (mut journal, _) = FlashJournal::open(flash, LAYOUT, &mut scratch, |_| {})
1617                .await
1618                .unwrap();
1619            journal.initialize_empty().await.unwrap();
1620            let mut persistence =
1621                EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1622                    TestFlash::new(),
1623                    LAYOUT,
1624                    EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(
1625                        0,
1626                    )),
1627                    FixedRouteSnapshotKeys::new(),
1628                    (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1629                );
1630            persistence.flash = None;
1631            persistence.journal = Some(journal);
1632            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1633            let attempt = InstantMillis(2_000);
1634            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, attempt);
1635            persistence.route_dirty_since = Some(InstantMillis(0));
1636            persistence.progress(&mut engine, attempt).await;
1637            fail_next_write.set(true);
1638            persistence.progress(&mut engine, attempt).await;
1639            assert_eq!(persistence.next_compaction_not_before, None);
1640            assert_eq!(persistence.compaction, None);
1641
1642            let flash = persistence.journal.take().unwrap().release();
1643            let mut restored = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1644                flash,
1645                LAYOUT,
1646                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1647                FixedRouteSnapshotKeys::new(),
1648                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1649            );
1650            restored.restore(&mut engine, InstantMillis(0)).await;
1651            assert_eq!(restored.next_compaction_not_before, None);
1652        });
1653    }
1654
1655    #[test]
1656    fn timebase_writes_and_repeated_reboots_do_not_move_the_compaction_deadline() {
1657        embassy_futures::block_on(async {
1658            let mut persistence = ready();
1659            let deadline = persistence.next_compaction_not_before;
1660            persistence
1661                .journal
1662                .as_mut()
1663                .unwrap()
1664                .record_timebase(InstantMillis(3 * 60 * 60 * 1_000))
1665                .await
1666                .unwrap();
1667            let flash = persistence.journal.take().unwrap().release();
1668            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1669
1670            let mut first = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1671                flash,
1672                LAYOUT,
1673                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1674                FixedRouteSnapshotKeys::new(),
1675                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1676            );
1677            first.restore(&mut engine, InstantMillis(0)).await;
1678            assert_eq!(first.next_compaction_not_before, deadline);
1679
1680            let flash = first.journal.take().unwrap().release();
1681            let mut second = EmbeddedFlashPersistence::<_, FixedRouteSnapshotKeys<8>, _, 4>::new(
1682                flash,
1683                LAYOUT,
1684                EmbeddedPersistencePolicy::hopspot_default(EmbeddedCompactionPolicy::hopspot(0)),
1685                FixedRouteSnapshotKeys::new(),
1686                (|_| {}) as fn(EmbeddedPersistenceDiagnostic),
1687            );
1688            second.restore(&mut engine, InstantMillis(0)).await;
1689            assert_eq!(second.next_compaction_not_before, deadline);
1690        });
1691    }
1692
1693    #[test]
1694    fn overflow_during_compaction_commits_once_and_defers_the_next_snapshot() {
1695        let diagnostics = Rc::new(RefCell::new(Vec::new()));
1696        let observed = Rc::clone(&diagnostics);
1697        let mut persistence = ready_with_observer(move |diagnostic| {
1698            observed.borrow_mut().push(diagnostic);
1699        });
1700        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1701        let now = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1702        persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
1703        persistence.route_dirty_since = Some(InstantMillis(now.0 - 2_000));
1704        embassy_futures::block_on(persistence.progress(&mut engine, now));
1705
1706        for byte in 0..6 {
1707            persistence.queue_route(
1708                PendingRouteDelta::RouteUpsert(DestinationHash::new(
1709                    [byte; TRUNCATED_HASH_BYTE_LEN],
1710                )),
1711                now,
1712            );
1713        }
1714        persistence.route_dirty_since = Some(now);
1715        for _ in 0..8 {
1716            embassy_futures::block_on(persistence.progress(&mut engine, now));
1717        }
1718
1719        assert_eq!(persistence.compaction, None);
1720        assert!(persistence.snapshot_required);
1721        assert!(persistence.state_not_saved());
1722        assert_eq!(
1723            persistence.next_deadline(now),
1724            Some(InstantMillis(
1725                now.0 + HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS
1726            ))
1727        );
1728        let diagnostics = diagnostics.borrow();
1729        assert_eq!(
1730            diagnostics
1731                .iter()
1732                .filter(|diagnostic| matches!(
1733                    diagnostic,
1734                    EmbeddedPersistenceDiagnostic::CompactionStarted { .. }
1735                ))
1736                .count(),
1737            1
1738        );
1739        assert_eq!(
1740            diagnostics
1741                .iter()
1742                .filter(|diagnostic| matches!(
1743                    diagnostic,
1744                    EmbeddedPersistenceDiagnostic::CompactionCompleted { .. }
1745                ))
1746                .count(),
1747            1
1748        );
1749        assert_eq!(
1750            diagnostics
1751                .iter()
1752                .filter(|diagnostic| matches!(
1753                    diagnostic,
1754                    EmbeddedPersistenceDiagnostic::DurabilityDeferred { .. }
1755                ))
1756                .count(),
1757            1
1758        );
1759    }
1760
1761    #[test]
1762    fn captured_route_keys_survive_slot_shifts_and_new_routes_land_after_compaction() {
1763        embassy_futures::block_on(async {
1764            let mut persistence = ready();
1765            let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1766            let rows = [signed_route(0x31, &[0xA1]), signed_route(0x32, &[0xA2])];
1767            for row in &rows {
1768                assert_eq!(
1769                    engine.seed_route(row, InstantMillis(1_000)),
1770                    RouteSeedOutcome::Seeded
1771                );
1772            }
1773            let now = InstantMillis(HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1774            persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
1775            persistence.route_dirty_since = Some(InstantMillis(now.0 - 2_000));
1776            persistence.progress(&mut engine, now).await;
1777            persistence.progress(&mut engine, now).await;
1778            persistence.progress(&mut engine, now).await;
1779
1780            let removed = rows[0].destination;
1781            let retained = rows[1].destination;
1782            let _ = engine.drop_route(&removed, AttachedInterfaces::new(&[]));
1783            let added = signed_route(0x34, &[0xA4]);
1784            assert_eq!(engine.seed_route(&added, now), RouteSeedOutcome::Seeded);
1785            persistence.queue_route(PendingRouteDelta::RouteUpsert(added.destination), now);
1786            persistence.route_dirty_since = Some(now);
1787
1788            for _ in 0..8 {
1789                persistence.progress(&mut engine, now).await;
1790            }
1791            assert_eq!(persistence.compaction, None);
1792            assert_eq!(
1793                (
1794                    persistence.pending_routes.len(),
1795                    persistence.snapshot_required,
1796                    persistence.write_failed,
1797                    persistence.route_dirty_since,
1798                    persistence.landing_batch,
1799                ),
1800                (1, false, false, Some(now), None)
1801            );
1802            let correction_at = InstantMillis(now.0 + 2_000);
1803            persistence.progress(&mut engine, correction_at).await;
1804            persistence.progress(&mut engine, correction_at).await;
1805
1806            let flash = persistence.journal.take().unwrap().release();
1807            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1808            let mut restored = Vec::new();
1809            let _ = FlashJournal::open(flash, LAYOUT, &mut scratch, |record| {
1810                if record.kind != FlashJournalRecordKind::RouteUpsert {
1811                    return;
1812                }
1813                let mut rows = read_routing_table_snapshot(record.payload).unwrap();
1814                restored.push(rows.next().unwrap().unwrap().destination);
1815            })
1816            .await
1817            .unwrap();
1818            assert_eq!(restored, vec![retained, added.destination]);
1819        });
1820    }
1821
1822    #[test]
1823    fn sixteen_route_records_restore_eight_and_report_capacity_drops() {
1824        type EightRouteStorage =
1825            crate::storage::TestFixedStorage<8, 8, 256, 2, 2, 16, 4, 4, 4, 4, 4, 16>;
1826        let now = InstantMillis(1_000);
1827        let mut engine = EngineState::<EightRouteStorage>::default();
1828        let mut report = EmbeddedPersistenceRestoreReport {
1829            logical_start: now,
1830            route_seeded_count: 0,
1831            route_refused_count: 0,
1832            route_dropped_count: 0,
1833            ratchet_seeded_count: 0,
1834            ratchet_refused_count: 0,
1835            warning: None,
1836        };
1837
1838        for secret in 1..=16 {
1839            let row = signed_route(secret, &[]);
1840            let required = routing_table_snapshot_len(core::iter::once(row.clone()));
1841            let mut scratch = [0u8; RECORD_SCRATCH_LEN];
1842            let written =
1843                write_routing_table_snapshot(core::iter::once(row), &mut scratch[..required])
1844                    .unwrap();
1845            apply_record(
1846                &mut engine,
1847                now,
1848                FlashJournalRecord {
1849                    epoch: 1,
1850                    kind: FlashJournalRecordKind::RouteUpsert,
1851                    payload: &scratch[..written],
1852                },
1853                &mut report,
1854            );
1855        }
1856
1857        assert_eq!(
1858            report,
1859            EmbeddedPersistenceRestoreReport {
1860                logical_start: now,
1861                route_seeded_count: 8,
1862                route_refused_count: 0,
1863                route_dropped_count: 8,
1864                ratchet_seeded_count: 0,
1865                ratchet_refused_count: 0,
1866                warning: None,
1867            }
1868        );
1869    }
1870
1871    #[test]
1872    fn thirty_days_of_pressure_erase_each_arena_sector_at_most_fifteen_times() {
1873        let diagnostics = Rc::new(RefCell::new(Vec::new()));
1874        let observed = Rc::clone(&diagnostics);
1875        let mut persistence = ready_with_observer(move |diagnostic| {
1876            observed.borrow_mut().push(diagnostic);
1877        });
1878        let mut engine = EngineState::<crate::storage::GrowableHeap>::default();
1879        embassy_futures::block_on(async {
1880            for day in 1..=30 {
1881                let now = InstantMillis(day * HOPSPOT_MINIMUM_COMPACTION_INTERVAL_MILLIS);
1882                persistence.require_snapshot(EmbeddedPersistenceTarget::Routes, now);
1883                persistence.route_dirty_since = Some(InstantMillis(now.0 - 2_000));
1884                for _ in 0..8 {
1885                    persistence.progress(&mut engine, now).await;
1886                }
1887                assert_eq!(persistence.compaction, None);
1888            }
1889        });
1890        assert_eq!(
1891            diagnostics
1892                .borrow()
1893                .iter()
1894                .filter(|diagnostic| matches!(
1895                    diagnostic,
1896                    EmbeddedPersistenceDiagnostic::CompactionStarted { .. }
1897                ))
1898                .count(),
1899            30
1900        );
1901        let flash = persistence.journal.take().unwrap().release();
1902        assert_eq!(
1903            [
1904                flash.sector_erases[2] - 1,
1905                flash.sector_erases[3] - 1,
1906                flash.sector_erases[4],
1907                flash.sector_erases[5],
1908            ],
1909            [15, 15, 15, 15]
1910        );
1911    }
1912}