Skip to main content

behavior_actors/atomic/
restart.rs

1//! Checked restart admission shared by fixed supervision and direct pools.
2
3use core::cmp::Ordering;
4use core::num::{NonZeroU32, NonZeroUsize};
5use core::time::Duration;
6use std::time::Instant;
7
8use thiserror::Error;
9
10/// Sliding-window admission limit for automatic worker replacement.
11#[derive(Clone, Copy, Debug, Eq, PartialEq)]
12pub struct RestartLimit {
13    maximum: u32,
14    window: Duration,
15}
16
17impl RestartLimit {
18    /// Define the admitted replacement count and inclusive time window.
19    #[must_use]
20    pub const fn new(maximum: u32, window: Duration) -> Self {
21        Self { maximum, window }
22    }
23}
24
25#[derive(Clone, Copy, Debug, Eq, PartialEq)]
26enum RestartReleaseSchedule {
27    Immediate,
28    Constant {
29        delay: Duration,
30    },
31    Linear {
32        initial: Duration,
33        maximum: Duration,
34    },
35    Exponential {
36        initial: Duration,
37        maximum: Duration,
38    },
39}
40
41/// Checked delay schedule for successive automatic recoveries.
42#[derive(Clone, Copy, Debug, Eq, PartialEq)]
43pub struct RestartRelease {
44    schedule: RestartReleaseSchedule,
45}
46
47/// Invalid delay configuration rejected before an aggregate exists.
48#[derive(Clone, Copy, Debug, Eq, Error, PartialEq)]
49pub enum RestartReleaseError {
50    /// A delayed schedule used a zero duration.
51    #[error("restart release delay must be positive")]
52    ZeroDelay,
53    /// A growing schedule capped its delay below the initial duration.
54    #[error("restart release maximum must not be below its initial delay")]
55    MaximumBelowInitial,
56}
57
58/// Exact automatic-release calculation failure.
59#[derive(Clone, Copy, Debug, Eq, Error, PartialEq)]
60pub enum RestartReleaseFailure {
61    /// Checked delay arithmetic exceeded the representable duration.
62    #[error("restart release delay arithmetic overflowed")]
63    DurationOverflow,
64}
65
66#[derive(Clone, Copy, Debug, Eq, PartialEq)]
67pub(in crate::atomic) enum RecoveryRelease {
68    Immediate,
69    Delayed(Duration),
70}
71
72#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
73pub(in crate::atomic) struct RecoveryCount {
74    admitted: u32,
75}
76
77#[derive(Clone, Copy, Debug, Eq, PartialEq)]
78pub(in crate::atomic) struct RecoveryOrdinal(NonZeroU32);
79
80pub(in crate::atomic) struct ProposedRecovery<'a> {
81    count: &'a mut RecoveryCount,
82    ordinal: RecoveryOrdinal,
83}
84
85pub(in crate::atomic) enum RecoveryProposal<'a> {
86    Available(ProposedRecovery<'a>),
87    Exhausted { admitted: u32 },
88}
89
90#[derive(Clone, Copy, Debug, Eq, PartialEq)]
91struct RestartCharge {
92    observed_at: Instant,
93    replacements: NonZeroU32,
94}
95
96#[derive(Debug, Eq, PartialEq)]
97pub(in crate::atomic) struct RestartBudget {
98    charges: Vec<RestartCharge>,
99}
100
101#[derive(Clone, Copy, Debug, Eq, PartialEq)]
102pub(in crate::atomic) enum RestartDenial {
103    RestartLimitReached {
104        active: u32,
105        requested: NonZeroUsize,
106        maximum: u32,
107    },
108    ClockRegressed {
109        previous: Instant,
110        observed: Instant,
111    },
112    RecoveryCountExhausted {
113        admitted: u32,
114    },
115    ReleaseCalculationFailed(RestartReleaseFailure),
116}
117
118pub(in crate::atomic) enum RestartAdmission<'a> {
119    Proposed(RestartProposal<'a>),
120    Denied {
121        budget: RestartBudget,
122        reason: RestartDenial,
123    },
124}
125
126pub(in crate::atomic) struct RestartProposal<'a> {
127    recovery: ProposedRecovery<'a>,
128    budget: RestartBudget,
129    proposed_budget: RestartBudget,
130    release: RecoveryRelease,
131}
132
133enum BudgetCheck {
134    Admitted(RestartBudget),
135    Denied(RestartDenial),
136}
137
138impl RecoveryCount {
139    pub(in crate::atomic) const fn new() -> Self {
140        Self { admitted: 0 }
141    }
142
143    pub(in crate::atomic) fn propose(&mut self) -> RecoveryProposal<'_> {
144        match self.admitted.checked_add(1).and_then(NonZeroU32::new) {
145            Some(next) => RecoveryProposal::Available(ProposedRecovery {
146                count: self,
147                ordinal: RecoveryOrdinal(next),
148            }),
149            None => RecoveryProposal::Exhausted {
150                admitted: self.admitted,
151            },
152        }
153    }
154
155    #[cfg(test)]
156    const fn from_admitted(admitted: u32) -> Self {
157        Self { admitted }
158    }
159
160    #[cfg(test)]
161    const fn admitted(self) -> u32 {
162        self.admitted
163    }
164}
165
166impl RecoveryOrdinal {
167    pub(in crate::atomic) const fn get(self) -> u32 {
168        self.0.get()
169    }
170}
171
172impl ProposedRecovery<'_> {
173    pub(in crate::atomic) const fn ordinal(&self) -> RecoveryOrdinal {
174        self.ordinal
175    }
176
177    pub(in crate::atomic) fn accept(self) {
178        self.count.admitted = self.ordinal.get();
179    }
180
181    pub(in crate::atomic) fn decline(self) {}
182}
183
184impl RestartBudget {
185    pub(in crate::atomic) const fn empty() -> Self {
186        Self {
187            charges: Vec::new(),
188        }
189    }
190
191    fn preview(
192        &self,
193        limit: RestartLimit,
194        observed_at: Instant,
195        replacements: NonZeroUsize,
196    ) -> BudgetCheck {
197        if let Some(previous) = self.charges.last().map(|charge| charge.observed_at) {
198            match observed_at.cmp(&previous) {
199                Ordering::Less => {
200                    return BudgetCheck::Denied(RestartDenial::ClockRegressed {
201                        previous,
202                        observed: observed_at,
203                    });
204                }
205                Ordering::Equal | Ordering::Greater => {}
206            }
207        }
208
209        let cutoff = observed_at.checked_sub(limit.window);
210        let first_active = cutoff.map_or(0, |cutoff| {
211            self.charges
212                .partition_point(|charge| charge.observed_at < cutoff)
213        });
214        let active_charges = &self.charges[first_active..];
215        let active = active_charges
216            .iter()
217            .try_fold(0_u32, |active, charge| {
218                active.checked_add(charge.replacements.get())
219            })
220            .expect("committed restart charges are bounded by their accepted u32 limit");
221
222        let charge = match u32::try_from(replacements.get())
223            .ok()
224            .and_then(NonZeroU32::new)
225        {
226            Some(charge) => charge,
227            None => {
228                return BudgetCheck::Denied(RestartDenial::RestartLimitReached {
229                    active,
230                    requested: replacements,
231                    maximum: limit.maximum,
232                });
233            }
234        };
235
236        let total = match active.checked_add(charge.get()) {
237            Some(total) => total,
238            None => {
239                return BudgetCheck::Denied(RestartDenial::RestartLimitReached {
240                    active,
241                    requested: replacements,
242                    maximum: limit.maximum,
243                });
244            }
245        };
246        match total.cmp(&limit.maximum) {
247            Ordering::Greater => {
248                return BudgetCheck::Denied(RestartDenial::RestartLimitReached {
249                    active,
250                    requested: replacements,
251                    maximum: limit.maximum,
252                });
253            }
254            Ordering::Less | Ordering::Equal => {}
255        }
256
257        let mut charges = Vec::with_capacity(active_charges.len() + 1);
258        charges.extend_from_slice(active_charges);
259        charges.push(RestartCharge {
260            observed_at,
261            replacements: charge,
262        });
263        BudgetCheck::Admitted(Self { charges })
264    }
265
266    #[cfg(test)]
267    fn latest(&self) -> Option<(Instant, NonZeroU32)> {
268        self.charges
269            .last()
270            .map(|charge| (charge.observed_at, charge.replacements))
271    }
272}
273
274impl RestartProposal<'_> {
275    #[cfg(test)]
276    pub(in crate::atomic) const fn ordinal(&self) -> RecoveryOrdinal {
277        self.recovery.ordinal()
278    }
279
280    pub(in crate::atomic) const fn release(&self) -> RecoveryRelease {
281        self.release
282    }
283
284    pub(in crate::atomic) fn accept(self) -> RestartBudget {
285        self.recovery.accept();
286        self.proposed_budget
287    }
288
289    pub(in crate::atomic) fn decline(self) -> RestartBudget {
290        self.recovery.decline();
291        self.budget
292    }
293}
294
295impl RestartRelease {
296    /// Release every eligible recovery without a timer prerequisite.
297    #[must_use]
298    pub const fn immediate() -> Self {
299        Self {
300            schedule: RestartReleaseSchedule::Immediate,
301        }
302    }
303
304    /// Use one positive delay for every eligible recovery.
305    pub fn constant(delay: Duration) -> Result<Self, RestartReleaseError> {
306        let delay = Self::positive(delay)?;
307        Ok(Self {
308            schedule: RestartReleaseSchedule::Constant { delay },
309        })
310    }
311
312    /// Increase delay linearly from a positive initial value up to a maximum.
313    pub fn linear(initial: Duration, maximum: Duration) -> Result<Self, RestartReleaseError> {
314        let (initial, maximum) = Self::growth(initial, maximum)?;
315        Ok(Self {
316            schedule: RestartReleaseSchedule::Linear { initial, maximum },
317        })
318    }
319
320    /// Increase delay exponentially from a positive initial value up to a maximum.
321    pub fn exponential(initial: Duration, maximum: Duration) -> Result<Self, RestartReleaseError> {
322        let (initial, maximum) = Self::growth(initial, maximum)?;
323        Ok(Self {
324            schedule: RestartReleaseSchedule::Exponential { initial, maximum },
325        })
326    }
327
328    fn calculate(self, ordinal: RecoveryOrdinal) -> Result<RecoveryRelease, RestartReleaseFailure> {
329        let delayed = match self.schedule {
330            RestartReleaseSchedule::Immediate => return Ok(RecoveryRelease::Immediate),
331            RestartReleaseSchedule::Constant { delay } => delay,
332            RestartReleaseSchedule::Linear { initial, maximum } => {
333                multiply(initial, ordinal.get().into())?.min(maximum)
334            }
335            RestartReleaseSchedule::Exponential { initial, maximum } => {
336                let shift = ordinal.get() - 1;
337                let factor = 1_u128
338                    .checked_shl(shift)
339                    .ok_or(RestartReleaseFailure::DurationOverflow)?;
340                multiply(initial, factor)?.min(maximum)
341            }
342        };
343        Ok(RecoveryRelease::Delayed(delayed))
344    }
345
346    fn positive(delay: Duration) -> Result<Duration, RestartReleaseError> {
347        match delay.cmp(&Duration::ZERO) {
348            Ordering::Equal => Err(RestartReleaseError::ZeroDelay),
349            Ordering::Less | Ordering::Greater => Ok(delay),
350        }
351    }
352
353    fn growth(
354        initial: Duration,
355        maximum: Duration,
356    ) -> Result<(Duration, Duration), RestartReleaseError> {
357        let initial = Self::positive(initial)?;
358        match maximum.cmp(&initial) {
359            Ordering::Less => Err(RestartReleaseError::MaximumBelowInitial),
360            Ordering::Equal | Ordering::Greater => Ok((initial, maximum)),
361        }
362    }
363}
364
365pub(in crate::atomic) fn admit_restart(
366    count: &mut RecoveryCount,
367    budget: RestartBudget,
368    limit: RestartLimit,
369    release: RestartRelease,
370    observed_at: Instant,
371    replacements: NonZeroUsize,
372) -> RestartAdmission<'_> {
373    let recovery = match count.propose() {
374        RecoveryProposal::Available(proposal) => proposal,
375        RecoveryProposal::Exhausted { admitted } => {
376            return RestartAdmission::Denied {
377                budget,
378                reason: RestartDenial::RecoveryCountExhausted { admitted },
379            };
380        }
381    };
382    let proposed_budget = match budget.preview(limit, observed_at, replacements) {
383        BudgetCheck::Admitted(budget) => budget,
384        BudgetCheck::Denied(reason) => {
385            recovery.decline();
386            return RestartAdmission::Denied { budget, reason };
387        }
388    };
389    match release.calculate(recovery.ordinal()) {
390        Ok(release) => RestartAdmission::Proposed(RestartProposal {
391            recovery,
392            budget,
393            proposed_budget,
394            release,
395        }),
396        Err(reason) => {
397            recovery.decline();
398            RestartAdmission::Denied {
399                budget,
400                reason: RestartDenial::ReleaseCalculationFailed(reason),
401            }
402        }
403    }
404}
405
406fn multiply(duration: Duration, factor: u128) -> Result<Duration, RestartReleaseFailure> {
407    let nanoseconds = duration
408        .as_nanos()
409        .checked_mul(factor)
410        .ok_or(RestartReleaseFailure::DurationOverflow)?;
411    let seconds = u64::try_from(nanoseconds / 1_000_000_000)
412        .map_err(|_| RestartReleaseFailure::DurationOverflow)?;
413    let subsecond = u32::try_from(nanoseconds % 1_000_000_000)
414        .map_err(|_| RestartReleaseFailure::DurationOverflow)?;
415    Ok(Duration::new(seconds, subsecond))
416}
417
418#[cfg(test)]
419mod tests {
420    use core::num::{NonZeroU32, NonZeroUsize};
421    use std::time::{Duration, Instant};
422
423    use super::{
424        RecoveryCount, RecoveryRelease, RestartAdmission, RestartBudget, RestartDenial,
425        RestartLimit, RestartRelease, RestartReleaseFailure, admit_restart,
426    };
427
428    fn workers(count: usize) -> NonZeroUsize {
429        NonZeroUsize::new(count).expect("test worker count is positive")
430    }
431
432    fn charge(count: u32) -> NonZeroU32 {
433        NonZeroU32::new(count).expect("test charge is positive")
434    }
435
436    #[test]
437    fn count_budget_and_release_commit_together() {
438        let observed = Instant::now();
439        let mut count = RecoveryCount::new();
440        let RestartAdmission::Proposed(proposal) = admit_restart(
441            &mut count,
442            RestartBudget::empty(),
443            RestartLimit::new(2, Duration::from_secs(10)),
444            RestartRelease::immediate(),
445            observed,
446            workers(2),
447        ) else {
448            panic!("valid restart is proposed")
449        };
450        assert_eq!(proposal.ordinal().get(), 1);
451        let release = proposal.release();
452        assert_eq!(release, RecoveryRelease::Immediate);
453        let budget = proposal.accept();
454        assert_eq!(count.admitted(), 1);
455        assert_eq!(budget.latest(), Some((observed, charge(2))));
456    }
457
458    #[test]
459    fn denial_preserves_count_and_budget() {
460        let observed = Instant::now();
461        let mut count = RecoveryCount::new();
462        let budget = RestartBudget::empty();
463        let RestartAdmission::Denied { budget, reason } = admit_restart(
464            &mut count,
465            budget,
466            RestartLimit::new(0, Duration::from_secs(10)),
467            RestartRelease::immediate(),
468            observed,
469            workers(1),
470        ) else {
471            panic!("zero maximum denies the restart")
472        };
473        assert_eq!(count.admitted(), 0);
474        assert_eq!(budget.latest(), None);
475        assert_eq!(
476            reason,
477            RestartDenial::RestartLimitReached {
478                active: 0,
479                requested: NonZeroUsize::MIN,
480                maximum: 0,
481            }
482        );
483    }
484
485    #[test]
486    fn release_failure_does_not_commit_count_or_budget() {
487        let observed = Instant::now();
488        let mut count = RecoveryCount::from_admitted(95);
489        let release = RestartRelease::exponential(Duration::from_nanos(2), Duration::MAX)
490            .expect("valid release policy");
491        let RestartAdmission::Denied { budget, reason } = admit_restart(
492            &mut count,
493            RestartBudget::empty(),
494            RestartLimit::new(1, Duration::from_secs(10)),
495            release,
496            observed,
497            workers(1),
498        ) else {
499            panic!("unrepresentable release is denied")
500        };
501        assert_eq!(count.admitted(), 95);
502        assert_eq!(budget.latest(), None);
503        assert_eq!(
504            reason,
505            RestartDenial::ReleaseCalculationFailed(RestartReleaseFailure::DurationOverflow)
506        );
507    }
508
509    #[test]
510    fn maximum_committed_budget_has_only_proposal_or_denial() {
511        let observed = Instant::now();
512        let mut count = RecoveryCount::new();
513        let budget = match admit_restart(
514            &mut count,
515            RestartBudget::empty(),
516            RestartLimit::new(u32::MAX, Duration::from_secs(10)),
517            RestartRelease::immediate(),
518            observed,
519            workers(u32::MAX as usize),
520        ) {
521            RestartAdmission::Proposed(proposal) => proposal.accept(),
522            RestartAdmission::Denied { .. } => panic!("the maximum charge is representable"),
523        };
524        assert_eq!(count.admitted(), 1);
525        assert_eq!(budget.latest(), Some((observed, charge(u32::MAX))));
526
527        match admit_restart(
528            &mut count,
529            budget,
530            RestartLimit::new(u32::MAX, Duration::from_secs(10)),
531            RestartRelease::immediate(),
532            observed,
533            workers(1),
534        ) {
535            RestartAdmission::Proposed(_) => panic!("the successor exceeds the maximum"),
536            RestartAdmission::Denied { budget, reason } => {
537                assert_eq!(count.admitted(), 1);
538                assert_eq!(budget.latest(), Some((observed, charge(u32::MAX))));
539                assert_eq!(
540                    reason,
541                    RestartDenial::RestartLimitReached {
542                        active: u32::MAX,
543                        requested: workers(1),
544                        maximum: u32::MAX,
545                    }
546                );
547            }
548        }
549    }
550
551    #[cfg(target_pointer_width = "64")]
552    #[test]
553    fn oversized_worker_request_is_an_unchanged_limit_denial() {
554        let observed = Instant::now();
555        let mut count = RecoveryCount::new();
556        let requested = workers(u32::MAX as usize + 1);
557        match admit_restart(
558            &mut count,
559            RestartBudget::empty(),
560            RestartLimit::new(u32::MAX, Duration::from_secs(10)),
561            RestartRelease::immediate(),
562            observed,
563            requested,
564        ) {
565            RestartAdmission::Proposed(_) => panic!("the request exceeds the charge width"),
566            RestartAdmission::Denied { budget, reason } => {
567                assert_eq!(count.admitted(), 0);
568                assert_eq!(budget.latest(), None);
569                assert_eq!(
570                    reason,
571                    RestartDenial::RestartLimitReached {
572                        active: 0,
573                        requested,
574                        maximum: u32::MAX,
575                    }
576                );
577            }
578        }
579    }
580}