behavior_actors/atomic/pool/
policy.rs1use std::mem;
4
5use behavior::Never;
6
7use super::super::worker::StopKind;
8use super::super::{RestartLimit, RestartRelease};
9
10#[derive(Clone, Copy, Debug, Eq, PartialEq)]
12pub struct BacklogCapacity {
13 maximum: usize,
14}
15
16impl BacklogCapacity {
17 #[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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
30pub enum Interruption {
31 Fail,
33 Retry,
35}
36
37#[derive(Clone, Copy, Debug, Eq, PartialEq)]
39pub enum PoolFailureReaction {
40 RetireRole,
42 StopPool,
44}
45
46#[derive(Clone, Copy, Debug, Eq, PartialEq)]
48pub enum PoolRecovery<Source> {
49 Permanent {
51 source: Source,
53 limit: RestartLimit,
55 release: RestartRelease,
57 failure: PoolFailureReaction,
59 },
60 Transient {
62 source: Source,
64 limit: RestartLimit,
66 release: RestartRelease,
68 failure: PoolFailureReaction,
70 },
71 Temporary {
73 failure: PoolFailureReaction,
75 },
76}
77
78impl<Source> PoolRecovery<Source> {
79 #[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 #[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 #[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}