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}