Skip to main content

rightkit_process/
restart.rs

1use std::{collections::VecDeque, error::Error, fmt, time::Duration};
2
3#[derive(Clone, Copy, Debug, Eq, PartialEq)]
4pub enum RestartMode {
5    Automatic,
6    OnDemand,
7    Disabled,
8}
9
10#[derive(Clone, Copy, Debug, Eq, PartialEq)]
11pub struct RestartPolicy {
12    mode: RestartMode,
13    max_restarts: u32,
14    window: Duration,
15    base_backoff: Duration,
16    max_backoff: Duration,
17}
18
19impl RestartPolicy {
20    pub fn mode(&self) -> RestartMode {
21        self.mode
22    }
23    pub fn max_restarts(&self) -> u32 {
24        self.max_restarts
25    }
26    pub fn window(&self) -> Duration {
27        self.window
28    }
29    pub fn base_backoff(&self) -> Duration {
30        self.base_backoff
31    }
32    pub fn max_backoff(&self) -> Duration {
33        self.max_backoff
34    }
35}
36
37impl Default for RestartPolicy {
38    /// ScrapeRight's crash-recovery plan: restart automatically, at most 3
39    /// unexpected restarts per rolling 60 s window, backoff 250 ms doubling to 5 s.
40    fn default() -> Self {
41        Self {
42            mode: RestartMode::Automatic,
43            max_restarts: 3,
44            window: Duration::from_secs(60),
45            base_backoff: Duration::from_millis(250),
46            max_backoff: Duration::from_secs(5),
47        }
48    }
49}
50
51#[derive(Clone, Copy, Debug, Eq, PartialEq)]
52pub struct InvalidRestartPolicy;
53
54impl fmt::Display for InvalidRestartPolicy {
55    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
56        f.write_str("restart window must be non-zero and max backoff must cover base backoff")
57    }
58}
59
60impl Error for InvalidRestartPolicy {}
61
62impl RestartPolicy {
63    pub fn new(
64        mode: RestartMode,
65        max_restarts: u32,
66        window: Duration,
67        base_backoff: Duration,
68        max_backoff: Duration,
69    ) -> Result<Self, InvalidRestartPolicy> {
70        if window.is_zero() || max_backoff < base_backoff {
71            return Err(InvalidRestartPolicy);
72        }
73        Ok(Self {
74            mode,
75            max_restarts,
76            window,
77            base_backoff,
78            max_backoff,
79        })
80    }
81}
82
83#[derive(Clone, Copy, Debug, Eq, PartialEq)]
84pub struct Generation(u64);
85
86#[derive(Clone, Copy, Debug, Eq, PartialEq)]
87pub struct RestartPermit {
88    epoch: u64,
89    delay: Duration,
90    not_before: Duration,
91}
92
93impl RestartPermit {
94    pub fn delay(&self) -> Duration {
95        self.delay
96    }
97}
98
99#[derive(Clone, Copy, Debug, Eq, PartialEq)]
100pub enum RestartOutcome {
101    Restart(RestartPermit),
102    AwaitingDemand,
103    Exhausted,
104    IgnoredStale,
105    Stopped,
106}
107
108#[derive(Clone, Debug)]
109pub struct RestartTracker {
110    policy: RestartPolicy,
111    active: Option<Generation>,
112    next_generation: u64,
113    restart_times: VecDeque<Duration>,
114    permit_epoch: u64,
115    awaiting_demand: bool,
116    shutdown: bool,
117}
118
119impl RestartTracker {
120    pub fn new(policy: RestartPolicy) -> Self {
121        Self {
122            policy,
123            active: None,
124            next_generation: 0,
125            restart_times: VecDeque::new(),
126            permit_epoch: 0,
127            awaiting_demand: false,
128            shutdown: false,
129        }
130    }
131
132    pub fn started(&mut self, _at: Duration) -> Generation {
133        self.next_generation = self.next_generation.saturating_add(1);
134        self.permit_epoch = self.permit_epoch.saturating_add(1);
135        self.awaiting_demand = false;
136        let generation = Generation(self.next_generation);
137        self.active = Some(generation);
138        generation
139    }
140
141    pub fn exited(&mut self, generation: Generation, at: Duration) -> RestartOutcome {
142        if self.shutdown {
143            return RestartOutcome::Stopped;
144        }
145        if self.active != Some(generation) {
146            return RestartOutcome::IgnoredStale;
147        }
148        self.active = None;
149        match self.policy.mode {
150            RestartMode::Automatic => self.schedule(at),
151            RestartMode::OnDemand => {
152                self.awaiting_demand = true;
153                RestartOutcome::AwaitingDemand
154            }
155            RestartMode::Disabled => RestartOutcome::Stopped,
156        }
157    }
158
159    pub fn request_restart(&mut self, at: Duration) -> RestartOutcome {
160        if self.shutdown || self.policy.mode == RestartMode::Disabled {
161            return RestartOutcome::Stopped;
162        }
163        if self.policy.mode == RestartMode::OnDemand && !self.awaiting_demand {
164            return RestartOutcome::AwaitingDemand;
165        }
166        self.awaiting_demand = false;
167        self.schedule(at)
168    }
169
170    /// Restarts granted inside the rolling window ending at `at` (diagnostics).
171    pub fn restarts_in_window(&self, at: Duration) -> u32 {
172        let window_start = at.saturating_sub(self.policy.window);
173        self.restart_times
174            .iter()
175            .filter(|t| **t >= window_start)
176            .count() as u32
177    }
178
179    /// Forget the restart budget (a user-initiated start after a terminal stop).
180    pub fn reset_budget(&mut self) {
181        self.restart_times.clear();
182        self.awaiting_demand = false;
183        self.shutdown = false;
184    }
185
186    pub fn permit_valid(&self, permit: &RestartPermit, at: Duration) -> bool {
187        !self.shutdown && permit.epoch == self.permit_epoch && at >= permit.not_before
188    }
189
190    pub fn shutdown(&mut self) {
191        self.shutdown = true;
192        self.active = None;
193        self.awaiting_demand = false;
194        self.permit_epoch = self.permit_epoch.saturating_add(1);
195    }
196
197    fn schedule(&mut self, at: Duration) -> RestartOutcome {
198        let window_start = at.saturating_sub(self.policy.window);
199        while self
200            .restart_times
201            .front()
202            .is_some_and(|restart| *restart < window_start)
203        {
204            self.restart_times.pop_front();
205        }
206        if self.restart_times.len() >= self.policy.max_restarts as usize {
207            return RestartOutcome::Exhausted;
208        }
209
210        let exponent = u32::try_from(self.restart_times.len()).unwrap_or(u32::MAX);
211        let multiplier = 1_u32.checked_shl(exponent.min(31)).unwrap_or(u32::MAX);
212        let delay = self
213            .policy
214            .base_backoff
215            .checked_mul(multiplier)
216            .unwrap_or(self.policy.max_backoff)
217            .min(self.policy.max_backoff);
218        self.restart_times.push_back(at);
219        self.permit_epoch = self.permit_epoch.saturating_add(1);
220        let permit = RestartPermit {
221            epoch: self.permit_epoch,
222            delay,
223            not_before: at.checked_add(delay).unwrap_or(Duration::MAX),
224        };
225        RestartOutcome::Restart(permit)
226    }
227}