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
13const STORAGE_REPAIR_PHASE_OFFSET: Duration = Duration::from_secs(5);
15const STORAGE_REPAIR_ADMISSION_BUDGET: Duration = DATA_CHANNEL_SEND_ACCEPT_BUDGET;
17const 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
101struct 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 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 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 pub async fn wait(self: Arc<Self>, interval: Duration) {
272 self.wait_with(interval, StopToken::never()).await;
273 }
274
275 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}