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 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 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 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}