saluki_core/runtime/restart.rs
1use std::{collections::VecDeque, time::Duration};
2
3use serde::{Deserialize, Serialize};
4use tokio::time::Instant;
5use tracing::debug;
6
7/// Restart mode for child processes.
8#[derive(Clone, Copy, Debug, PartialEq, Eq, Deserialize, Serialize)]
9#[serde(rename_all = "snake_case")]
10pub enum RestartMode {
11 /// Restarts the failed child process only.
12 OneForOne,
13
14 /// Restarts all child processes, including the failed one.
15 OneForAll,
16}
17
18/// Restart strategy for a supervisor.
19///
20/// Defaults to one-to-one mode (only restart the failed process) and a restart intensity of 1 over a period of 5
21/// seconds.
22///
23/// # Restarts and permanent failure
24///
25/// A supervisor will allow up to `intensity` process restarts, across all child processes, over a given `period`. When
26/// this limit is exceeded, the supervisor will stop all child processes and return an error itself, indicating that the
27/// supervisor has failed overall.
28///
29/// Permanent failure bubbles up to the parent supervisor, until reaching the root supervisor. Once permanent failure
30/// reaches the root supervisor, and the root supervisor exceeds its own restart limits, the root supervisor will fail
31/// and cease execution.
32#[derive(Clone, Copy, Debug)]
33pub struct RestartStrategy {
34 mode: RestartMode,
35 intensity: usize,
36 period: Duration,
37}
38
39impl RestartStrategy {
40 /// Creates a new `RestartStrategy` with the given mode, intensity, and period.
41 pub const fn new(mode: RestartMode, intensity: usize, period: Duration) -> Self {
42 Self {
43 mode,
44 intensity,
45 period,
46 }
47 }
48
49 /// Creates a new `RestartStrategy` with the one-to-one restart mode, and the default intensity/period.
50 pub fn one_to_one() -> Self {
51 Self {
52 mode: RestartMode::OneForOne,
53 ..Default::default()
54 }
55 }
56
57 /// Creates a new `RestartStrategy` with the one-for-all restart mode, and the default intensity/period.
58 pub fn one_for_all() -> Self {
59 Self {
60 mode: RestartMode::OneForAll,
61 ..Default::default()
62 }
63 }
64
65 /// Sets the restart intensity and period for the strategy.
66 pub const fn with_intensity_and_period(mut self, intensity: usize, period: Duration) -> Self {
67 self.intensity = intensity;
68 self.period = period;
69 self
70 }
71
72 pub(super) fn mode(&self) -> RestartMode {
73 self.mode
74 }
75
76 pub(super) fn intensity(&self) -> usize {
77 self.intensity
78 }
79
80 pub(super) fn period(&self) -> Duration {
81 self.period
82 }
83}
84
85impl Default for RestartStrategy {
86 fn default() -> Self {
87 Self::new(RestartMode::OneForOne, 1, Duration::from_secs(5))
88 }
89}
90
91/// Restart policy for an individual child process.
92///
93/// While [`RestartStrategy`] governs supervisor-wide behavior (which children are restarted together, and how often
94/// before the supervisor gives up), [`RestartType`] governs whether an _individual_ child is eligible for restart at
95/// all, based on how it exited.
96///
97/// Defaults to [`Permanent`][Self::Permanent], which marks a child process to always be restarted.
98#[derive(Clone, Copy, Debug, PartialEq, Eq, Default, Deserialize, Serialize)]
99#[serde(rename_all = "snake_case")]
100pub enum RestartType {
101 /// The child is always restarted, whether it exits normally or abnormally.
102 ///
103 /// This suits long-lived processes that are always expected to be running.
104 #[default]
105 Permanent,
106
107 /// The child is restarted only if it exits abnormally.
108 ///
109 /// An abnormal exit is an error, panic, or forced abort. A normal exit (the child's future resolves with `Ok(())`)
110 /// is treated as intentional, and the child is not restarted. This governs the child's _own_ exit; a transient
111 /// child is still restarted when a sibling triggers a [`RestartMode::OneForAll`] group restart, matching
112 /// Erlang/OTP.
113 Transient,
114
115 /// The child is never restarted, regardless of how it exits.
116 ///
117 /// This suits short-lived, on-demand children -- for example, one task per network connection -- whose termination
118 /// is a normal part of operation. A temporary child is never restarted even when a sibling triggers a
119 /// [`RestartMode::OneForAll`] group restart: it is shut down with the group but not brought back.
120 Temporary,
121}
122
123impl RestartType {
124 /// Returns `true` if a child with this restart policy should be restarted, given how it exited.
125 ///
126 /// `abnormal` indicates the child exited due to an error, panic, or forced abort, rather than completing normally.
127 pub(super) fn should_restart(self, abnormal: bool) -> bool {
128 match self {
129 Self::Permanent => true,
130 Self::Transient => abnormal,
131 Self::Temporary => false,
132 }
133 }
134}
135
136pub(super) enum RestartAction {
137 /// Execute a restart with the given mode.
138 Restart(RestartMode),
139
140 /// Supervisor must shutdown as the maximum number of restarts has been reached.
141 Shutdown,
142}
143
144pub(super) struct RestartState {
145 strategy: RestartStrategy,
146 restart_history: VecDeque<Instant>,
147}
148
149impl RestartState {
150 /// Creates a new `RestartState` with the given strategy.
151 pub fn new(strategy: RestartStrategy) -> Self {
152 Self {
153 strategy,
154 restart_history: VecDeque::with_capacity(strategy.intensity),
155 }
156 }
157
158 /// Evaluates a restart based on the current state and determine the action the supervisor should take in response.
159 pub fn evaluate_restart(&mut self) -> RestartAction {
160 // Short circuit if our intensity is zero.
161 if self.strategy.intensity == 0 {
162 debug!("Restart strategy configured with restart intensity of zero, shutting down.");
163 return RestartAction::Shutdown;
164 }
165
166 // Since we only keep track of the last `intensity` restarts, we simply need to check if the oldest restart
167 // we're tracking is within `period` of the current time, and if the number of tracked restarts is equal to
168 // `intensity`.
169 //
170 // When both of these are true, we have exceeded the restart intensity limit and must shutdown.
171 let now = Instant::now();
172 if self.restart_history.len() == self.strategy.intensity {
173 let oldest = self.restart_history.front().expect("restart history cannot be empty");
174 saluki_antithesis::always_or_unreachable!(
175 now >= *oldest,
176 "restart-intensity window clock did not move backward"
177 );
178 if now.saturating_duration_since(*oldest) < self.strategy.period {
179 debug!(
180 "Restart limit exceeded ({} in {:?}), shutting down.",
181 self.strategy.intensity, self.strategy.period
182 );
183 return RestartAction::Shutdown;
184 }
185
186 // Remove the oldest restart from the history since it is outside the period.
187 self.restart_history.pop_front();
188 }
189
190 // Track this latest restart.
191 self.restart_history.push_back(now);
192
193 debug!("Restart limit not exceeded, restarting worker.");
194 RestartAction::Restart(self.strategy.mode)
195 }
196}