Skip to main content

rings_core/dht/stabilization/
maintenance.rs

1use std::sync::Arc;
2use std::time::Duration;
3
4use super::storage_repair::StorageRepairOutcome;
5use super::Stabilizer;
6use super::STABILIZATION_STEP_TIMEOUT;
7use super::STABILIZATION_STOP_POLL_INTERVAL;
8use crate::lifecycle::StopToken;
9use crate::swarm::transport::DATA_CHANNEL_SEND_ACCEPT_BUDGET;
10use crate::utils::try_sleep;
11use crate::utils::Instant;
12
13/// The quiet phase reserved for topology stabilization before periodic repair.
14const STORAGE_REPAIR_PHASE_OFFSET: Duration = Duration::from_secs(5);
15/// The uninterrupted first-frame admission window for one storage repair delivery.
16const STORAGE_REPAIR_ADMISSION_BUDGET: Duration = DATA_CHANNEL_SEND_ACCEPT_BUDGET;
17/// Separate completed maintenance phases by at least one cooperative poll.
18const MAINTENANCE_QUIET_GAP: Duration = STABILIZATION_STOP_POLL_INTERVAL;
19
20#[derive(Clone, Copy, Debug, Eq, PartialEq)]
21enum MaintenanceTask {
22    Stabilize,
23    Repair,
24}
25
26#[cfg(all(test, target_family = "wasm"))]
27#[derive(Clone, Copy, Debug, Eq, PartialEq)]
28pub(crate) enum MaintenancePhaseKind {
29    Stabilize,
30    Repair,
31}
32
33#[cfg(all(test, target_family = "wasm"))]
34#[derive(Clone, Copy, Debug, Eq, PartialEq)]
35pub(crate) struct MaintenancePhaseEvent {
36    pub(crate) local: crate::dht::Did,
37    pub(crate) kind: MaintenancePhaseKind,
38    pub(crate) started_at_ms: u64,
39}
40
41#[cfg(all(test, target_family = "wasm"))]
42thread_local! {
43    static MAINTENANCE_PHASE_TRACE: std::cell::RefCell<Vec<MaintenancePhaseEvent>> = const {
44        std::cell::RefCell::new(Vec::new())
45    };
46}
47
48#[cfg(all(test, target_family = "wasm"))]
49pub(crate) fn reset_maintenance_phase_trace_for_test() {
50    MAINTENANCE_PHASE_TRACE.with(|trace| trace.borrow_mut().clear());
51}
52
53#[cfg(all(test, target_family = "wasm"))]
54pub(crate) fn maintenance_phase_trace_for_test(
55    local: crate::dht::Did,
56) -> Vec<MaintenancePhaseEvent> {
57    MAINTENANCE_PHASE_TRACE.with(|trace| {
58        trace
59            .borrow()
60            .iter()
61            .copied()
62            .filter(|event| event.local == local)
63            .collect()
64    })
65}
66
67#[cfg(all(test, target_family = "wasm"))]
68fn record_maintenance_phase_for_test(
69    local: crate::dht::Did,
70    task: MaintenanceTask,
71    started_at_ms: u64,
72) {
73    let kind = match task {
74        MaintenanceTask::Stabilize => MaintenancePhaseKind::Stabilize,
75        MaintenanceTask::Repair => MaintenancePhaseKind::Repair,
76    };
77    MAINTENANCE_PHASE_TRACE.with(|trace| {
78        trace.borrow_mut().push(MaintenancePhaseEvent {
79            local,
80            kind,
81            started_at_ms,
82        });
83    });
84}
85
86#[cfg(not(all(test, target_family = "wasm")))]
87fn record_maintenance_phase_for_test(
88    _local: crate::dht::Did,
89    _task: MaintenanceTask,
90    _started_at_ms: u64,
91) {
92}
93
94#[derive(Clone, Copy, Debug, Eq, PartialEq)]
95struct MaintenanceDecision {
96    task: Option<MaintenanceTask>,
97    periodic_repair_due: bool,
98    repair_deferred_for_window: bool,
99}
100
101/// Absolute phase schedule for topology and storage maintenance.
102///
103/// State relation:
104/// - `next_stabilize_ms` is advanced only after the selected stabilization run
105///   completes, skipping every deadline at or before its completion time.
106/// - every elapsed `next_repair_ms` submits a persistent repair intent.
107/// - repair may run only after `repair_not_before_ms` and when its first-frame
108///   admission budget and a post-repair quiet gap fit before
109///   `next_stabilize_ms`.
110/// - when repair is pending at stabilization completion, the following
111///   repair turn takes precedence over the following stabilization, even when
112///   timer jitter wakes the loop after the reserved deadline.
113/// - every repair attempt consumes that precedence before stabilization may
114///   reserve another turn, preserving fairness between both maintenance tasks.
115/// - repair is tracked to its final frame. If its tail exceeds the admission
116///   estimate, this serial loop cannot overlap it with stabilization and
117///   reconciles the next deadline from actual completion.
118struct MaintenanceSchedule {
119    period_ms: u64,
120    next_stabilize_ms: u64,
121    next_repair_ms: u64,
122    repair_not_before_ms: u64,
123    repair_admission_budget_ms: u64,
124    repair_turn_reserved: bool,
125}
126
127impl MaintenanceSchedule {
128    fn new(now_ms: u64, interval: Duration) -> Self {
129        let period_ms = duration_ms(interval).max(2);
130        let offset_ms = duration_ms(STORAGE_REPAIR_PHASE_OFFSET)
131            .min(period_ms / 2)
132            .max(1);
133        let next_stabilize_ms = now_ms.saturating_add(period_ms);
134        Self {
135            period_ms,
136            next_stabilize_ms,
137            next_repair_ms: next_stabilize_ms.saturating_add(offset_ms),
138            repair_not_before_ms: now_ms,
139            repair_admission_budget_ms: duration_ms(STORAGE_REPAIR_ADMISSION_BUDGET),
140            repair_turn_reserved: false,
141        }
142    }
143
144    /// Select at most one task. Stabilization wins when both phases are due,
145    /// while the repair phase is preserved as an intent rather than dropped.
146    fn poll(&mut self, now_ms: u64, repair_pending: bool) -> MaintenanceDecision {
147        let periodic_repair_due = self.advance_repair_deadline_if_due(now_ms);
148        let effective_repair_pending = repair_pending || periodic_repair_due;
149        if !effective_repair_pending {
150            self.repair_turn_reserved = false;
151        }
152        let reserved_repair_ready = self.reserved_repair_ready(now_ms, effective_repair_pending);
153        let stabilization_due = now_ms >= self.next_stabilize_ms;
154        let repair_has_window = effective_repair_pending && self.can_start_storage_repair(now_ms);
155        let task = if reserved_repair_ready {
156            Some(MaintenanceTask::Repair)
157        } else if stabilization_due {
158            Some(MaintenanceTask::Stabilize)
159        } else if repair_has_window {
160            Some(MaintenanceTask::Repair)
161        } else {
162            None
163        };
164
165        MaintenanceDecision {
166            task,
167            periodic_repair_due,
168            repair_deferred_for_window: effective_repair_pending
169                && !self.repair_turn_reserved
170                && !stabilization_due
171                && !self.has_storage_repair_window(now_ms),
172        }
173    }
174
175    /// Reconcile deadlines against completion time. This prevents a long run
176    /// from causing immediate catch-up stabilization passes.
177    fn complete_stabilization(&mut self, completed_at_ms: u64, repair_pending: bool) -> bool {
178        let periodic_repair_due = self.advance_repair_deadline_if_due(completed_at_ms);
179        self.next_stabilize_ms =
180            next_deadline_after(self.next_stabilize_ms, self.period_ms, completed_at_ms);
181        self.repair_not_before_ms =
182            completed_at_ms.saturating_add(duration_ms(MAINTENANCE_QUIET_GAP));
183        self.repair_turn_reserved = repair_pending || periodic_repair_due;
184        if self.repair_turn_reserved {
185            let reserved_deadline = self
186                .repair_not_before_ms
187                .saturating_add(self.required_repair_window_ms());
188            self.next_stabilize_ms = self.next_stabilize_ms.max(reserved_deadline);
189        }
190        periodic_repair_due
191    }
192
193    fn complete_repair(&mut self, completed_at_ms: u64, succeeded: bool) {
194        self.repair_turn_reserved = false;
195        let post_repair_deadline =
196            completed_at_ms.saturating_add(duration_ms(MAINTENANCE_QUIET_GAP));
197        self.next_stabilize_ms = self.next_stabilize_ms.max(post_repair_deadline);
198        self.repair_not_before_ms = if succeeded {
199            post_repair_deadline
200        } else {
201            self.next_repair_ms
202        };
203    }
204
205    fn advance_repair_deadline_if_due(&mut self, now_ms: u64) -> bool {
206        if now_ms < self.next_repair_ms {
207            return false;
208        }
209        self.next_repair_ms = next_deadline_after(self.next_repair_ms, self.period_ms, now_ms);
210        true
211    }
212
213    fn can_start_storage_repair(&self, now_ms: u64) -> bool {
214        now_ms >= self.repair_not_before_ms && self.has_storage_repair_window(now_ms)
215    }
216
217    fn reserved_repair_ready(&self, now_ms: u64, repair_pending: bool) -> bool {
218        self.repair_turn_reserved && repair_pending && now_ms >= self.repair_not_before_ms
219    }
220
221    fn has_storage_repair_window(&self, now_ms: u64) -> bool {
222        self.storage_repair_window_ms(now_ms) >= self.required_repair_window_ms()
223    }
224
225    fn required_repair_window_ms(&self) -> u64 {
226        self.repair_admission_budget_ms
227            .saturating_add(duration_ms(MAINTENANCE_QUIET_GAP))
228    }
229
230    fn storage_repair_window_ms(&self, now_ms: u64) -> u64 {
231        self.next_stabilize_ms.saturating_sub(now_ms)
232    }
233
234    fn next_wake_ms(&self, now_ms: u64, repair_pending: bool) -> u64 {
235        let mut next_ms = self.next_stabilize_ms.min(self.next_repair_ms);
236        if !repair_pending {
237            return next_ms;
238        }
239        if self.can_start_storage_repair(now_ms) {
240            return now_ms;
241        }
242        if now_ms < self.repair_not_before_ms
243            && self.has_storage_repair_window(self.repair_not_before_ms)
244        {
245            next_ms = next_ms.min(self.repair_not_before_ms);
246        }
247        next_ms
248    }
249}
250
251fn duration_ms(duration: Duration) -> u64 {
252    u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
253}
254
255fn next_deadline_after(deadline_ms: u64, period_ms: u64, now_ms: u64) -> u64 {
256    let elapsed_ms = now_ms.saturating_sub(deadline_ms);
257    let periods = elapsed_ms
258        .checked_div(period_ms)
259        .unwrap_or(0)
260        .saturating_add(1);
261    let next = deadline_ms.saturating_add(periods.saturating_mul(period_ms));
262    if next <= now_ms {
263        u64::MAX
264    } else {
265        next
266    }
267}
268
269impl Stabilizer {
270    /// Run topology stabilization and storage repair in staggered phases.
271    pub async fn wait(self: Arc<Self>, interval: Duration) {
272        self.wait_with(interval, StopToken::never()).await;
273    }
274
275    /// Run staggered maintenance until `stop` asks this loop to exit.
276    ///
277    /// Repair requests are shared with disconnect handlers and survive missed
278    /// phase deadlines. Cooperative stop is observed between phases; the
279    /// per-step deadline may still cancel a hung network maintenance future.
280    pub async fn wait_with(self: Arc<Self>, interval: Duration, stop: StopToken) {
281        let origin = Instant::now();
282        let mut schedule = MaintenanceSchedule::new(0, interval);
283        loop {
284            if stop.should_stop() {
285                return;
286            }
287
288            let now_ms = monotonic_elapsed_ms(&origin);
289            let decision = schedule.poll(now_ms, self.transport.storage_repair_requested());
290            if decision.periodic_repair_due {
291                self.transport.request_storage_repair();
292            }
293            if decision.repair_deferred_for_window {
294                tracing::debug!(
295                    target: "rings_core::dht::stabilization",
296                    local = %self.dht.did,
297                    available_ms = schedule.storage_repair_window_ms(now_ms),
298                    required_ms = schedule.required_repair_window_ms(),
299                    "STABILIZATION deferred storage repair for an admission window"
300                );
301            }
302
303            match decision.task {
304                Some(MaintenanceTask::Stabilize) => {
305                    record_maintenance_phase_for_test(
306                        self.dht.did,
307                        MaintenanceTask::Stabilize,
308                        now_ms,
309                    );
310                    self.stabilize_topology_with_step_timeout(STABILIZATION_STEP_TIMEOUT)
311                        .await;
312                    let periodic_repair_due = schedule.complete_stabilization(
313                        monotonic_elapsed_ms(&origin),
314                        self.transport.storage_repair_requested(),
315                    );
316                    if periodic_repair_due {
317                        self.transport.request_storage_repair();
318                    }
319                }
320                Some(MaintenanceTask::Repair) => {
321                    record_maintenance_phase_for_test(
322                        self.dht.did,
323                        MaintenanceTask::Repair,
324                        now_ms,
325                    );
326                    if let Some(outcome) = self.run_requested_storage_repair().await {
327                        schedule
328                            .complete_repair(monotonic_elapsed_ms(&origin), outcome.is_complete());
329                    }
330                }
331                None => {
332                    let deadline_ms =
333                        schedule.next_wake_ms(now_ms, self.transport.storage_repair_requested());
334                    if !sleep_until_or_stop(&origin, deadline_ms, &stop).await {
335                        return;
336                    }
337                }
338            }
339        }
340    }
341
342    pub(crate) async fn run_requested_storage_repair(&self) -> Option<StorageRepairOutcome> {
343        if !self.transport.claim_storage_repair() {
344            return None;
345        }
346        let outcome = self
347            .run_step(
348                "repair_storage",
349                STABILIZATION_STEP_TIMEOUT,
350                self.repair_storage(),
351            )
352            .await
353            .unwrap_or(StorageRepairOutcome::Deferred);
354        if !outcome.is_complete() {
355            self.transport.request_storage_repair();
356        }
357        Some(outcome)
358    }
359}
360
361fn monotonic_elapsed_ms(origin: &Instant) -> u64 {
362    duration_ms(origin.elapsed())
363}
364
365async fn sleep_until_or_stop(origin: &Instant, deadline_ms: u64, stop: &StopToken) -> bool {
366    loop {
367        if stop.should_stop() {
368            return false;
369        }
370        let delay = remaining_delay(deadline_ms, monotonic_elapsed_ms(origin));
371        if delay.is_zero() {
372            return !stop.should_stop();
373        }
374        if !try_sleep(delay.min(STABILIZATION_STOP_POLL_INTERVAL)).await {
375            tracing::error!("stopping stabilization maintenance after timer scheduling failed");
376            return false;
377        }
378    }
379}
380
381fn remaining_delay(deadline_ms: u64, now_ms: u64) -> Duration {
382    Duration::from_millis(deadline_ms.saturating_sub(now_ms))
383}
384
385#[cfg(test)]
386mod tests {
387    use super::*;
388
389    const PERIOD: Duration = Duration::from_secs(15);
390
391    #[test]
392    fn test_maintenance_phases_are_staggered_within_each_period() {
393        let mut schedule = MaintenanceSchedule::new(0, PERIOD);
394
395        assert_eq!(schedule.poll(14_999, false).task, None);
396        assert_eq!(
397            schedule.poll(15_000, false).task,
398            Some(MaintenanceTask::Stabilize)
399        );
400        assert!(!schedule.complete_stabilization(15_000, false));
401        let repair = schedule.poll(20_000, false);
402        assert!(repair.periodic_repair_due);
403        assert_eq!(repair.task, Some(MaintenanceTask::Repair));
404        schedule.complete_repair(20_000, true);
405        assert_eq!(
406            schedule.poll(30_000, false).task,
407            Some(MaintenanceTask::Stabilize)
408        );
409    }
410
411    #[test]
412    fn test_repeated_stabilization_overruns_preserve_repair_intent() {
413        let mut schedule = MaintenanceSchedule::new(0, PERIOD);
414
415        assert_eq!(
416            schedule.poll(15_000, false).task,
417            Some(MaintenanceTask::Stabilize)
418        );
419        assert!(schedule.complete_stabilization(21_000, false));
420        let missed_phase = schedule.poll(21_000, true);
421        assert!(!missed_phase.periodic_repair_due);
422        assert_eq!(missed_phase.task, None);
423        assert_eq!(
424            schedule.poll(21_050, true).task,
425            Some(MaintenanceTask::Repair)
426        );
427    }
428
429    #[test]
430    fn test_long_stabilization_skips_missed_stabilization_deadlines() {
431        let mut schedule = MaintenanceSchedule::new(0, PERIOD);
432
433        assert_eq!(
434            schedule.poll(15_000, false).task,
435            Some(MaintenanceTask::Stabilize)
436        );
437        assert!(schedule.complete_stabilization(46_000, false));
438
439        assert_eq!(schedule.next_stabilize_ms, 60_000);
440        assert_ne!(
441            schedule.poll(46_000, false).task,
442            Some(MaintenanceTask::Stabilize)
443        );
444    }
445
446    #[test]
447    fn test_stabilization_reserves_a_window_for_pending_repair() {
448        let mut schedule = MaintenanceSchedule::new(0, PERIOD);
449
450        assert_eq!(
451            schedule.poll(15_000, false).task,
452            Some(MaintenanceTask::Stabilize)
453        );
454        assert!(schedule.complete_stabilization(26_000, false));
455        let reserved_stabilization_ms = schedule.next_stabilize_ms;
456        assert!(schedule.has_storage_repair_window(schedule.repair_not_before_ms));
457        assert_eq!(schedule.poll(26_000, true).task, None);
458        assert_eq!(
459            schedule.poll(26_050, true).task,
460            Some(MaintenanceTask::Repair)
461        );
462        let quiet_gap_ms = duration_ms(MAINTENANCE_QUIET_GAP);
463        let repair_completed_ms = reserved_stabilization_ms.saturating_sub(quiet_gap_ms);
464        schedule.complete_repair(repair_completed_ms, true);
465        assert_eq!(schedule.poll(repair_completed_ms, false).task, None);
466        assert_eq!(
467            schedule.poll(reserved_stabilization_ms, false).task,
468            Some(MaintenanceTask::Stabilize)
469        );
470    }
471
472    #[test]
473    fn test_reserved_repair_turn_survives_timer_overshoot() {
474        let mut schedule = MaintenanceSchedule::new(0, Duration::from_millis(100));
475
476        assert_eq!(
477            schedule.poll(100, false).task,
478            Some(MaintenanceTask::Stabilize)
479        );
480        schedule.complete_stabilization(100, true);
481        let first_missed_deadline = schedule
482            .repair_not_before_ms
483            .saturating_add(schedule.required_repair_window_ms())
484            .saturating_add(1);
485
486        assert_eq!(first_missed_deadline, schedule.next_stabilize_ms + 1);
487        assert_eq!(
488            schedule.poll(first_missed_deadline, true).task,
489            Some(MaintenanceTask::Repair)
490        );
491    }
492
493    #[test]
494    fn test_repeated_timer_overshoots_preserve_repair_and_stabilization_fairness() {
495        let mut schedule = MaintenanceSchedule::new(0, Duration::from_millis(500));
496        let mut stabilization_start_ms = 500;
497
498        for _ in 0..3 {
499            assert_eq!(
500                schedule.poll(stabilization_start_ms, true).task,
501                Some(MaintenanceTask::Stabilize)
502            );
503            schedule.complete_stabilization(stabilization_start_ms, true);
504            let late_wake_ms = schedule.next_stabilize_ms.saturating_add(1);
505            assert_eq!(
506                schedule.poll(late_wake_ms, true).task,
507                Some(MaintenanceTask::Repair)
508            );
509            schedule.complete_repair(late_wake_ms.saturating_add(1), false);
510            stabilization_start_ms = schedule.next_stabilize_ms;
511        }
512    }
513
514    #[test]
515    fn test_repair_overrun_reconciles_stabilization_with_actual_completion() {
516        let mut schedule = MaintenanceSchedule::new(0, PERIOD);
517
518        assert_eq!(
519            schedule.poll(15_000, false).task,
520            Some(MaintenanceTask::Stabilize)
521        );
522        assert!(!schedule.complete_stabilization(15_000, false));
523        assert_eq!(
524            schedule.poll(20_000, false).task,
525            Some(MaintenanceTask::Repair)
526        );
527
528        schedule.complete_repair(31_000, true);
529
530        assert_eq!(schedule.poll(31_000, false).task, None);
531        assert_eq!(
532            schedule.poll(31_050, false).task,
533            Some(MaintenanceTask::Stabilize)
534        );
535    }
536
537    #[test]
538    fn test_failed_repair_waits_for_the_next_topology_phase() {
539        let mut schedule = MaintenanceSchedule::new(0, PERIOD);
540
541        assert_eq!(
542            schedule.poll(15_000, false).task,
543            Some(MaintenanceTask::Stabilize)
544        );
545        assert!(!schedule.complete_stabilization(15_000, false));
546        assert_eq!(
547            schedule.poll(20_000, false).task,
548            Some(MaintenanceTask::Repair)
549        );
550        schedule.complete_repair(20_001, false);
551
552        assert_eq!(schedule.next_wake_ms(20_001, true), 30_000);
553        assert_eq!(
554            schedule.poll(30_000, true).task,
555            Some(MaintenanceTask::Stabilize)
556        );
557        assert!(!schedule.complete_stabilization(30_000, true));
558        assert_eq!(
559            schedule.poll(30_050, true).task,
560            Some(MaintenanceTask::Repair)
561        );
562    }
563
564    #[test]
565    fn test_late_timer_wake_recomputes_from_absolute_deadline() {
566        assert_eq!(remaining_delay(100, 90), Duration::from_millis(10));
567        assert_eq!(remaining_delay(100, 125), Duration::ZERO);
568    }
569}