use std::sync::Arc;
use std::time::Duration;
use super::storage_repair::StorageRepairOutcome;
use super::Stabilizer;
use super::STABILIZATION_STEP_TIMEOUT;
use super::STABILIZATION_STOP_POLL_INTERVAL;
use crate::lifecycle::StopToken;
use crate::swarm::transport::DATA_CHANNEL_SEND_ACCEPT_BUDGET;
use crate::utils::try_sleep;
use crate::utils::Instant;
const STORAGE_REPAIR_PHASE_OFFSET: Duration = Duration::from_secs(5);
const STORAGE_REPAIR_ADMISSION_BUDGET: Duration = DATA_CHANNEL_SEND_ACCEPT_BUDGET;
const MAINTENANCE_QUIET_GAP: Duration = STABILIZATION_STOP_POLL_INTERVAL;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum MaintenanceTask {
Stabilize,
Repair,
}
#[cfg(all(test, target_family = "wasm"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum MaintenancePhaseKind {
Stabilize,
Repair,
}
#[cfg(all(test, target_family = "wasm"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) struct MaintenancePhaseEvent {
pub(crate) local: crate::dht::Did,
pub(crate) kind: MaintenancePhaseKind,
pub(crate) started_at_ms: u64,
}
#[cfg(all(test, target_family = "wasm"))]
thread_local! {
static MAINTENANCE_PHASE_TRACE: std::cell::RefCell<Vec<MaintenancePhaseEvent>> = const {
std::cell::RefCell::new(Vec::new())
};
}
#[cfg(all(test, target_family = "wasm"))]
pub(crate) fn reset_maintenance_phase_trace_for_test() {
MAINTENANCE_PHASE_TRACE.with(|trace| trace.borrow_mut().clear());
}
#[cfg(all(test, target_family = "wasm"))]
pub(crate) fn maintenance_phase_trace_for_test(
local: crate::dht::Did,
) -> Vec<MaintenancePhaseEvent> {
MAINTENANCE_PHASE_TRACE.with(|trace| {
trace
.borrow()
.iter()
.copied()
.filter(|event| event.local == local)
.collect()
})
}
#[cfg(all(test, target_family = "wasm"))]
fn record_maintenance_phase_for_test(
local: crate::dht::Did,
task: MaintenanceTask,
started_at_ms: u64,
) {
let kind = match task {
MaintenanceTask::Stabilize => MaintenancePhaseKind::Stabilize,
MaintenanceTask::Repair => MaintenancePhaseKind::Repair,
};
MAINTENANCE_PHASE_TRACE.with(|trace| {
trace.borrow_mut().push(MaintenancePhaseEvent {
local,
kind,
started_at_ms,
});
});
}
#[cfg(not(all(test, target_family = "wasm")))]
fn record_maintenance_phase_for_test(
_local: crate::dht::Did,
_task: MaintenanceTask,
_started_at_ms: u64,
) {
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct MaintenanceDecision {
task: Option<MaintenanceTask>,
periodic_repair_due: bool,
repair_deferred_for_window: bool,
}
struct MaintenanceSchedule {
period_ms: u64,
next_stabilize_ms: u64,
next_repair_ms: u64,
repair_not_before_ms: u64,
repair_admission_budget_ms: u64,
repair_turn_reserved: bool,
}
impl MaintenanceSchedule {
fn new(now_ms: u64, interval: Duration) -> Self {
let period_ms = duration_ms(interval).max(2);
let offset_ms = duration_ms(STORAGE_REPAIR_PHASE_OFFSET)
.min(period_ms / 2)
.max(1);
let next_stabilize_ms = now_ms.saturating_add(period_ms);
Self {
period_ms,
next_stabilize_ms,
next_repair_ms: next_stabilize_ms.saturating_add(offset_ms),
repair_not_before_ms: now_ms,
repair_admission_budget_ms: duration_ms(STORAGE_REPAIR_ADMISSION_BUDGET),
repair_turn_reserved: false,
}
}
fn poll(&mut self, now_ms: u64, repair_pending: bool) -> MaintenanceDecision {
let periodic_repair_due = self.advance_repair_deadline_if_due(now_ms);
let effective_repair_pending = repair_pending || periodic_repair_due;
if !effective_repair_pending {
self.repair_turn_reserved = false;
}
let reserved_repair_ready = self.reserved_repair_ready(now_ms, effective_repair_pending);
let stabilization_due = now_ms >= self.next_stabilize_ms;
let repair_has_window = effective_repair_pending && self.can_start_storage_repair(now_ms);
let task = if reserved_repair_ready {
Some(MaintenanceTask::Repair)
} else if stabilization_due {
Some(MaintenanceTask::Stabilize)
} else if repair_has_window {
Some(MaintenanceTask::Repair)
} else {
None
};
MaintenanceDecision {
task,
periodic_repair_due,
repair_deferred_for_window: effective_repair_pending
&& !self.repair_turn_reserved
&& !stabilization_due
&& !self.has_storage_repair_window(now_ms),
}
}
fn complete_stabilization(&mut self, completed_at_ms: u64, repair_pending: bool) -> bool {
let periodic_repair_due = self.advance_repair_deadline_if_due(completed_at_ms);
self.next_stabilize_ms =
next_deadline_after(self.next_stabilize_ms, self.period_ms, completed_at_ms);
self.repair_not_before_ms =
completed_at_ms.saturating_add(duration_ms(MAINTENANCE_QUIET_GAP));
self.repair_turn_reserved = repair_pending || periodic_repair_due;
if self.repair_turn_reserved {
let reserved_deadline = self
.repair_not_before_ms
.saturating_add(self.required_repair_window_ms());
self.next_stabilize_ms = self.next_stabilize_ms.max(reserved_deadline);
}
periodic_repair_due
}
fn complete_repair(&mut self, completed_at_ms: u64, succeeded: bool) {
self.repair_turn_reserved = false;
let post_repair_deadline =
completed_at_ms.saturating_add(duration_ms(MAINTENANCE_QUIET_GAP));
self.next_stabilize_ms = self.next_stabilize_ms.max(post_repair_deadline);
self.repair_not_before_ms = if succeeded {
post_repair_deadline
} else {
self.next_repair_ms
};
}
fn advance_repair_deadline_if_due(&mut self, now_ms: u64) -> bool {
if now_ms < self.next_repair_ms {
return false;
}
self.next_repair_ms = next_deadline_after(self.next_repair_ms, self.period_ms, now_ms);
true
}
fn can_start_storage_repair(&self, now_ms: u64) -> bool {
now_ms >= self.repair_not_before_ms && self.has_storage_repair_window(now_ms)
}
fn reserved_repair_ready(&self, now_ms: u64, repair_pending: bool) -> bool {
self.repair_turn_reserved && repair_pending && now_ms >= self.repair_not_before_ms
}
fn has_storage_repair_window(&self, now_ms: u64) -> bool {
self.storage_repair_window_ms(now_ms) >= self.required_repair_window_ms()
}
fn required_repair_window_ms(&self) -> u64 {
self.repair_admission_budget_ms
.saturating_add(duration_ms(MAINTENANCE_QUIET_GAP))
}
fn storage_repair_window_ms(&self, now_ms: u64) -> u64 {
self.next_stabilize_ms.saturating_sub(now_ms)
}
fn next_wake_ms(&self, now_ms: u64, repair_pending: bool) -> u64 {
let mut next_ms = self.next_stabilize_ms.min(self.next_repair_ms);
if !repair_pending {
return next_ms;
}
if self.can_start_storage_repair(now_ms) {
return now_ms;
}
if now_ms < self.repair_not_before_ms
&& self.has_storage_repair_window(self.repair_not_before_ms)
{
next_ms = next_ms.min(self.repair_not_before_ms);
}
next_ms
}
}
fn duration_ms(duration: Duration) -> u64 {
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
}
fn next_deadline_after(deadline_ms: u64, period_ms: u64, now_ms: u64) -> u64 {
let elapsed_ms = now_ms.saturating_sub(deadline_ms);
let periods = elapsed_ms
.checked_div(period_ms)
.unwrap_or(0)
.saturating_add(1);
let next = deadline_ms.saturating_add(periods.saturating_mul(period_ms));
if next <= now_ms {
u64::MAX
} else {
next
}
}
impl Stabilizer {
pub async fn wait(self: Arc<Self>, interval: Duration) {
self.wait_with(interval, StopToken::never()).await;
}
pub async fn wait_with(self: Arc<Self>, interval: Duration, stop: StopToken) {
let origin = Instant::now();
let mut schedule = MaintenanceSchedule::new(0, interval);
loop {
if stop.should_stop() {
return;
}
let now_ms = monotonic_elapsed_ms(&origin);
let decision = schedule.poll(now_ms, self.transport.storage_repair_requested());
if decision.periodic_repair_due {
self.transport.request_storage_repair();
}
if decision.repair_deferred_for_window {
tracing::debug!(
target: "rings_core::dht::stabilization",
local = %self.dht.did,
available_ms = schedule.storage_repair_window_ms(now_ms),
required_ms = schedule.required_repair_window_ms(),
"STABILIZATION deferred storage repair for an admission window"
);
}
match decision.task {
Some(MaintenanceTask::Stabilize) => {
record_maintenance_phase_for_test(
self.dht.did,
MaintenanceTask::Stabilize,
now_ms,
);
self.stabilize_topology_with_step_timeout(STABILIZATION_STEP_TIMEOUT)
.await;
let periodic_repair_due = schedule.complete_stabilization(
monotonic_elapsed_ms(&origin),
self.transport.storage_repair_requested(),
);
if periodic_repair_due {
self.transport.request_storage_repair();
}
}
Some(MaintenanceTask::Repair) => {
record_maintenance_phase_for_test(
self.dht.did,
MaintenanceTask::Repair,
now_ms,
);
if let Some(outcome) = self.run_requested_storage_repair().await {
schedule
.complete_repair(monotonic_elapsed_ms(&origin), outcome.is_complete());
}
}
None => {
let deadline_ms =
schedule.next_wake_ms(now_ms, self.transport.storage_repair_requested());
if !sleep_until_or_stop(&origin, deadline_ms, &stop).await {
return;
}
}
}
}
}
pub(crate) async fn run_requested_storage_repair(&self) -> Option<StorageRepairOutcome> {
if !self.transport.claim_storage_repair() {
return None;
}
let outcome = self
.run_step(
"repair_storage",
STABILIZATION_STEP_TIMEOUT,
self.repair_storage(),
)
.await
.unwrap_or(StorageRepairOutcome::Deferred);
if !outcome.is_complete() {
self.transport.request_storage_repair();
}
Some(outcome)
}
}
fn monotonic_elapsed_ms(origin: &Instant) -> u64 {
duration_ms(origin.elapsed())
}
async fn sleep_until_or_stop(origin: &Instant, deadline_ms: u64, stop: &StopToken) -> bool {
loop {
if stop.should_stop() {
return false;
}
let delay = remaining_delay(deadline_ms, monotonic_elapsed_ms(origin));
if delay.is_zero() {
return !stop.should_stop();
}
if !try_sleep(delay.min(STABILIZATION_STOP_POLL_INTERVAL)).await {
tracing::error!("stopping stabilization maintenance after timer scheduling failed");
return false;
}
}
}
fn remaining_delay(deadline_ms: u64, now_ms: u64) -> Duration {
Duration::from_millis(deadline_ms.saturating_sub(now_ms))
}
#[cfg(test)]
mod tests {
use super::*;
const PERIOD: Duration = Duration::from_secs(15);
#[test]
fn test_maintenance_phases_are_staggered_within_each_period() {
let mut schedule = MaintenanceSchedule::new(0, PERIOD);
assert_eq!(schedule.poll(14_999, false).task, None);
assert_eq!(
schedule.poll(15_000, false).task,
Some(MaintenanceTask::Stabilize)
);
assert!(!schedule.complete_stabilization(15_000, false));
let repair = schedule.poll(20_000, false);
assert!(repair.periodic_repair_due);
assert_eq!(repair.task, Some(MaintenanceTask::Repair));
schedule.complete_repair(20_000, true);
assert_eq!(
schedule.poll(30_000, false).task,
Some(MaintenanceTask::Stabilize)
);
}
#[test]
fn test_repeated_stabilization_overruns_preserve_repair_intent() {
let mut schedule = MaintenanceSchedule::new(0, PERIOD);
assert_eq!(
schedule.poll(15_000, false).task,
Some(MaintenanceTask::Stabilize)
);
assert!(schedule.complete_stabilization(21_000, false));
let missed_phase = schedule.poll(21_000, true);
assert!(!missed_phase.periodic_repair_due);
assert_eq!(missed_phase.task, None);
assert_eq!(
schedule.poll(21_050, true).task,
Some(MaintenanceTask::Repair)
);
}
#[test]
fn test_long_stabilization_skips_missed_stabilization_deadlines() {
let mut schedule = MaintenanceSchedule::new(0, PERIOD);
assert_eq!(
schedule.poll(15_000, false).task,
Some(MaintenanceTask::Stabilize)
);
assert!(schedule.complete_stabilization(46_000, false));
assert_eq!(schedule.next_stabilize_ms, 60_000);
assert_ne!(
schedule.poll(46_000, false).task,
Some(MaintenanceTask::Stabilize)
);
}
#[test]
fn test_stabilization_reserves_a_window_for_pending_repair() {
let mut schedule = MaintenanceSchedule::new(0, PERIOD);
assert_eq!(
schedule.poll(15_000, false).task,
Some(MaintenanceTask::Stabilize)
);
assert!(schedule.complete_stabilization(26_000, false));
let reserved_stabilization_ms = schedule.next_stabilize_ms;
assert!(schedule.has_storage_repair_window(schedule.repair_not_before_ms));
assert_eq!(schedule.poll(26_000, true).task, None);
assert_eq!(
schedule.poll(26_050, true).task,
Some(MaintenanceTask::Repair)
);
let quiet_gap_ms = duration_ms(MAINTENANCE_QUIET_GAP);
let repair_completed_ms = reserved_stabilization_ms.saturating_sub(quiet_gap_ms);
schedule.complete_repair(repair_completed_ms, true);
assert_eq!(schedule.poll(repair_completed_ms, false).task, None);
assert_eq!(
schedule.poll(reserved_stabilization_ms, false).task,
Some(MaintenanceTask::Stabilize)
);
}
#[test]
fn test_reserved_repair_turn_survives_timer_overshoot() {
let mut schedule = MaintenanceSchedule::new(0, Duration::from_millis(100));
assert_eq!(
schedule.poll(100, false).task,
Some(MaintenanceTask::Stabilize)
);
schedule.complete_stabilization(100, true);
let first_missed_deadline = schedule
.repair_not_before_ms
.saturating_add(schedule.required_repair_window_ms())
.saturating_add(1);
assert_eq!(first_missed_deadline, schedule.next_stabilize_ms + 1);
assert_eq!(
schedule.poll(first_missed_deadline, true).task,
Some(MaintenanceTask::Repair)
);
}
#[test]
fn test_repeated_timer_overshoots_preserve_repair_and_stabilization_fairness() {
let mut schedule = MaintenanceSchedule::new(0, Duration::from_millis(500));
let mut stabilization_start_ms = 500;
for _ in 0..3 {
assert_eq!(
schedule.poll(stabilization_start_ms, true).task,
Some(MaintenanceTask::Stabilize)
);
schedule.complete_stabilization(stabilization_start_ms, true);
let late_wake_ms = schedule.next_stabilize_ms.saturating_add(1);
assert_eq!(
schedule.poll(late_wake_ms, true).task,
Some(MaintenanceTask::Repair)
);
schedule.complete_repair(late_wake_ms.saturating_add(1), false);
stabilization_start_ms = schedule.next_stabilize_ms;
}
}
#[test]
fn test_repair_overrun_reconciles_stabilization_with_actual_completion() {
let mut schedule = MaintenanceSchedule::new(0, PERIOD);
assert_eq!(
schedule.poll(15_000, false).task,
Some(MaintenanceTask::Stabilize)
);
assert!(!schedule.complete_stabilization(15_000, false));
assert_eq!(
schedule.poll(20_000, false).task,
Some(MaintenanceTask::Repair)
);
schedule.complete_repair(31_000, true);
assert_eq!(schedule.poll(31_000, false).task, None);
assert_eq!(
schedule.poll(31_050, false).task,
Some(MaintenanceTask::Stabilize)
);
}
#[test]
fn test_failed_repair_waits_for_the_next_topology_phase() {
let mut schedule = MaintenanceSchedule::new(0, PERIOD);
assert_eq!(
schedule.poll(15_000, false).task,
Some(MaintenanceTask::Stabilize)
);
assert!(!schedule.complete_stabilization(15_000, false));
assert_eq!(
schedule.poll(20_000, false).task,
Some(MaintenanceTask::Repair)
);
schedule.complete_repair(20_001, false);
assert_eq!(schedule.next_wake_ms(20_001, true), 30_000);
assert_eq!(
schedule.poll(30_000, true).task,
Some(MaintenanceTask::Stabilize)
);
assert!(!schedule.complete_stabilization(30_000, true));
assert_eq!(
schedule.poll(30_050, true).task,
Some(MaintenanceTask::Repair)
);
}
#[test]
fn test_late_timer_wake_recomputes_from_absolute_deadline() {
assert_eq!(remaining_delay(100, 90), Duration::from_millis(10));
assert_eq!(remaining_delay(100, 125), Duration::ZERO);
}
}