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}