Skip to main content

behavior_actors/atomic/pool/
policy.rs

1//! Policies shared by direct-worker pools.
2
3use std::mem;
4
5use behavior::Never;
6
7use super::super::worker::StopKind;
8use super::super::{RestartLimit, RestartRelease};
9
10/// Maximum number of accepted jobs that may wait in one admission queue.
11#[derive(Clone, Copy, Debug, Eq, PartialEq)]
12pub struct BacklogCapacity {
13    maximum: usize,
14}
15
16impl BacklogCapacity {
17    /// Select the waiting-job limit. Zero permits only immediate assignment.
18    #[must_use]
19    pub const fn new(maximum: usize) -> Self {
20        Self { maximum }
21    }
22
23    pub(in crate::atomic) const fn maximum(self) -> usize {
24        self.maximum
25    }
26}
27
28/// Customer disposition when an assigned worker stops.
29#[derive(Clone, Copy, Debug, Eq, PartialEq)]
30pub enum Interruption {
31    /// Return the accepted job to its customer.
32    Fail,
33    /// Reinsert the accepted job at its immutable admission position.
34    Retry,
35}
36
37/// Pool disposition when one worker role cannot recover.
38#[derive(Clone, Copy, Debug, Eq, PartialEq)]
39pub enum PoolFailureReaction {
40    /// Retire only the unavailable role and keep serving through other workers.
41    RetireRole,
42    /// Close the complete pool and begin actor-graph retirement.
43    StopPool,
44}
45
46/// Direct-worker eligibility, source, restart policy, and failure reaction.
47#[derive(Clone, Copy, Debug, Eq, PartialEq)]
48pub enum PoolRecovery<Source> {
49    /// Recover after normal and abnormal worker stops.
50    Permanent {
51        /// Typed capability used by Bombay to prepare a replacement worker.
52        source: Source,
53        /// Sliding-window restart admission limit.
54        limit: RestartLimit,
55        /// Immediate or delayed release policy.
56        release: RestartRelease,
57        /// Disposition when preparation or restart admission fails.
58        failure: PoolFailureReaction,
59    },
60    /// Recover only after abnormal worker stops.
61    Transient {
62        /// Typed capability used by Bombay to prepare a replacement worker.
63        source: Source,
64        /// Sliding-window restart admission limit.
65        limit: RestartLimit,
66        /// Immediate or delayed release policy.
67        release: RestartRelease,
68        /// Disposition when preparation or restart admission fails.
69        failure: PoolFailureReaction,
70    },
71    /// Never create a replacement worker.
72    Temporary {
73        /// Disposition after the role retires.
74        failure: PoolFailureReaction,
75    },
76}
77
78impl<Source> PoolRecovery<Source> {
79    /// Recover after every exact worker stop.
80    #[must_use]
81    pub const fn permanent(
82        source: Source,
83        limit: RestartLimit,
84        release: RestartRelease,
85        failure: PoolFailureReaction,
86    ) -> Self {
87        Self::Permanent {
88            source,
89            limit,
90            release,
91            failure,
92        }
93    }
94
95    /// Recover only after abnormal worker stops.
96    #[must_use]
97    pub const fn transient(
98        source: Source,
99        limit: RestartLimit,
100        release: RestartRelease,
101        failure: PoolFailureReaction,
102    ) -> Self {
103        Self::Transient {
104            source,
105            limit,
106            release,
107            failure,
108        }
109    }
110}
111
112impl PoolRecovery<Never> {
113    /// Retire stopped workers without a source placeholder.
114    #[must_use]
115    pub const fn temporary(failure: PoolFailureReaction) -> Self {
116        Self::Temporary { failure }
117    }
118}
119
120pub(in crate::atomic) enum WorkerRecoverySource<Source> {
121    Available(Source),
122    AwaitingReturn,
123}
124
125pub(in crate::atomic) enum WorkerRecoveryDecision<Source> {
126    PrepareWorker(Source),
127    WaitForSource,
128    RetireRole,
129    StopPool,
130}
131
132pub(in crate::atomic) enum PoolRecoveryState<Source> {
133    Permanent {
134        source: WorkerRecoverySource<Source>,
135        limit: RestartLimit,
136        release: RestartRelease,
137        failure: PoolFailureReaction,
138    },
139    Transient {
140        source: WorkerRecoverySource<Source>,
141        limit: RestartLimit,
142        release: RestartRelease,
143        failure: PoolFailureReaction,
144    },
145    Temporary {
146        failure: PoolFailureReaction,
147    },
148}
149
150impl<Source> From<PoolRecovery<Source>> for PoolRecoveryState<Source> {
151    fn from(recovery: PoolRecovery<Source>) -> Self {
152        match recovery {
153            PoolRecovery::Permanent {
154                source,
155                limit,
156                release,
157                failure,
158            } => Self::Permanent {
159                source: WorkerRecoverySource::Available(source),
160                limit,
161                release,
162                failure,
163            },
164            PoolRecovery::Transient {
165                source,
166                limit,
167                release,
168                failure,
169            } => Self::Transient {
170                source: WorkerRecoverySource::Available(source),
171                limit,
172                release,
173                failure,
174            },
175            PoolRecovery::Temporary { failure } => Self::Temporary { failure },
176        }
177    }
178}
179
180impl<Source> PoolRecoveryState<Source> {
181    pub(in crate::atomic) fn decide(&mut self, stop: StopKind) -> WorkerRecoveryDecision<Source> {
182        match (self, stop) {
183            (Self::Permanent { source, .. }, _)
184            | (Self::Transient { source, .. }, StopKind::Abnormal) => source.claim(),
185            (Self::Transient { failure, .. }, StopKind::Normal)
186            | (Self::Temporary { failure }, _) => match failure {
187                PoolFailureReaction::RetireRole => WorkerRecoveryDecision::RetireRole,
188                PoolFailureReaction::StopPool => WorkerRecoveryDecision::StopPool,
189            },
190        }
191    }
192
193    pub(in crate::atomic) const fn failure(&self) -> PoolFailureReaction {
194        match self {
195            Self::Permanent { failure, .. }
196            | Self::Transient { failure, .. }
197            | Self::Temporary { failure } => *failure,
198        }
199    }
200
201    pub(in crate::atomic) fn restore_source(&mut self, returned: Source) -> Result<(), Source> {
202        let source = match self {
203            Self::Permanent { source, .. } | Self::Transient { source, .. } => source,
204            Self::Temporary { .. } => return Err(returned),
205        };
206        match source {
207            WorkerRecoverySource::AwaitingReturn => {
208                *source = WorkerRecoverySource::Available(returned);
209                Ok(())
210            }
211            WorkerRecoverySource::Available(_) => Err(returned),
212        }
213    }
214
215    pub(in crate::atomic) fn claim_waiting_source(&mut self) -> Option<Source> {
216        match self {
217            Self::Permanent { source, .. } | Self::Transient { source, .. } => {
218                source.claim_source()
219            }
220            Self::Temporary { .. } => None,
221        }
222    }
223}
224
225impl<Source> WorkerRecoverySource<Source> {
226    fn claim(&mut self) -> WorkerRecoveryDecision<Source> {
227        match self.claim_source() {
228            Some(source) => WorkerRecoveryDecision::PrepareWorker(source),
229            None => WorkerRecoveryDecision::WaitForSource,
230        }
231    }
232
233    fn claim_source(&mut self) -> Option<Source> {
234        match mem::replace(self, Self::AwaitingReturn) {
235            Self::Available(source) => Some(source),
236            Self::AwaitingReturn => None,
237        }
238    }
239}