1use std::sync::Arc;
4
5use behavior::{
6 Actions, Behavior, BehaviorAddr, ChildCreationOutcome, CreationId, CreationRejection,
7 EndpointAddress, EstablishedActor, InterpreterFault, ItemSettlement, Never, RoutedCreation,
8 SettledItem,
9};
10
11use crate::{Exit, StopOnShutdown, TerminalOutcome};
12
13mod activation;
14mod initialization;
15mod preparation;
16mod roster;
17
18pub(in crate::atomic) use activation::WorkerActivationOutcome;
19pub use activation::{
20 ActivationStartRejection, BeginActivation, WorkerActivation, WorkerActivationGrant,
21};
22pub use initialization::{
23 ActivationPermit, InitializationAttempt, InitializeWorker, WorkerInitializationFailure,
24 WorkerInitializationOutcome, WorkerInitializationReport,
25};
26pub use preparation::{
27 PendingWorkerPreparation, PrepareWorkers, StartingWorkerPreparation, WorkerPreparation,
28 WorkerPreparationStarted, WorkerSource,
29};
30pub(in crate::atomic) use preparation::{
31 PreparationTicket, WorkerPreparationExpectation, WorkerPreparationOutcome,
32};
33pub use roster::InitialWorkerRejection;
34pub(in crate::atomic) use roster::prepare_initial_workers;
35
36#[derive(Clone, Copy, Debug, Eq, PartialEq)]
37pub(in crate::atomic) enum StopKind {
38 Normal,
39 Abnormal,
40}
41
42pub(in crate::atomic) const fn stop_kind<A>(outcome: &TerminalOutcome<A>) -> StopKind
43where
44 A: behavior::Address,
45{
46 match outcome {
47 Ok(Exit::Normal | Exit::Collected) => StopKind::Normal,
48 Ok(Exit::LinkDied(_) | Exit::SupervisionFailed(_)) | Err(_) => StopKind::Abnormal,
49 }
50}
51
52#[doc(hidden)]
54pub type HostedInitialization<W> = Actions<
55 BehaviorAddr<W>,
56 <W as Behavior>::Ph,
57 <StopOnShutdown<W> as Behavior>::Sends,
58 <W as Behavior>::Birth,
59>;
60
61#[must_use = "a rejected worker and its initialization actions require custody"]
64pub struct WorkerRecovery<W>
65where
66 W: Behavior,
67{
68 worker: W,
69 initialization: HostedInitialization<W>,
70}
71
72impl<W> WorkerRecovery<W>
73where
74 W: Behavior,
75 StopOnShutdown<W>: Behavior<Protocol = W::Protocol, Ph = W::Ph, Birth = W::Birth>,
76{
77 pub(in crate::atomic) const fn new(worker: W, initialization: HostedInitialization<W>) -> Self {
78 Self {
79 worker,
80 initialization,
81 }
82 }
83
84 #[must_use]
86 pub fn into_retirement(self) -> (W, HostedInitialization<W>) {
87 (self.worker, self.initialization)
88 }
89}
90
91impl<W> core::fmt::Debug for WorkerRecovery<W>
92where
93 W: Behavior,
94 StopOnShutdown<W>: Behavior<Protocol = W::Protocol, Ph = W::Ph, Birth = W::Birth>,
95{
96 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
97 formatter
98 .debug_struct("WorkerRecovery")
99 .finish_non_exhaustive()
100 }
101}
102
103pub enum WorkerCreationRejection<W>
105where
106 W: Behavior,
107{
108 NamespaceExhausted { worker: W },
110 WorkerRejected { worker: W, error: W::Error },
112 WorkerPanicked { worker: W },
114 HostRejected {
116 recovery: WorkerRecovery<W>,
117 reason: CreationRejection,
118 },
119 CreationRejected {
121 worker: W,
122 reason: CreationRejection,
123 },
124 InterpreterCorrupt { worker: W, fault: InterpreterFault },
126 InterpretationSkipped { worker: W },
128}
129
130pub(in crate::atomic) type WorkerCreationSettlement<W> = SettledItem<
131 RoutedCreation<BehaviorAddr<W>, StopOnShutdown<W>>,
132 ItemSettlement<
133 RoutedCreation<BehaviorAddr<W>, StopOnShutdown<W>>,
134 ChildCreationOutcome<StopOnShutdown<W>, behavior::ChildHead>,
135 CreationRejection,
136 Never,
137 >,
138>;
139
140pub(in crate::atomic) enum WorkerCreationOutcome<W>
141where
142 W: Behavior,
143 BehaviorAddr<W>: EndpointAddress,
144 StopOnShutdown<W>: Behavior<Protocol = W::Protocol>,
145{
146 Established(EstablishedActor<StopOnShutdown<W>>),
147 Rejected(WorkerCreationRejection<W>),
148}
149
150#[cfg(test)]
151mod creation_outcome_tests {
152 use behavior::{Behavior, BehaviorAddr, EndpointAddress};
153
154 use super::WorkerCreationOutcome;
155 use crate::StopOnShutdown;
156
157 #[expect(
158 dead_code,
159 reason = "compile contract for complete worker creation custody"
160 )]
161 fn settled_worker_outcome<W>(outcome: WorkerCreationOutcome<W>)
162 where
163 W: Behavior,
164 BehaviorAddr<W>: EndpointAddress,
165 StopOnShutdown<W>: Behavior<Protocol = W::Protocol>,
166 {
167 match outcome {
168 WorkerCreationOutcome::Established(_) | WorkerCreationOutcome::Rejected(_) => {}
169 }
170 }
171}
172
173pub(in crate::atomic) fn settle_worker_creation<W>(
174 settlement: WorkerCreationSettlement<W>,
175) -> WorkerCreationOutcome<W>
176where
177 W: Behavior,
178 BehaviorAddr<W>: EndpointAddress,
179 StopOnShutdown<W>:
180 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
181{
182 match settlement {
183 SettledItem::Attempted(ItemSettlement::Accepted(created)) => match created {
184 ChildCreationOutcome::Established(child) => {
185 WorkerCreationOutcome::Established(child.into_parts().2)
186 }
187 ChildCreationOutcome::InitializationRejected { creation, error } => {
188 WorkerCreationOutcome::Rejected(WorkerCreationRejection::WorkerRejected {
189 worker: recover_worker(creation),
190 error,
191 })
192 }
193 ChildCreationOutcome::InitializationPanicked { creation } => {
194 WorkerCreationOutcome::Rejected(WorkerCreationRejection::WorkerPanicked {
195 worker: recover_worker(creation),
196 })
197 }
198 ChildCreationOutcome::HostRejected {
199 creation,
200 initialization,
201 reason,
202 } => WorkerCreationOutcome::Rejected(WorkerCreationRejection::HostRejected {
203 recovery: WorkerRecovery::new(recover_worker(creation), initialization),
204 reason,
205 }),
206 },
207 SettledItem::Attempted(ItemSettlement::Rejected { item, reason }) => {
208 WorkerCreationOutcome::Rejected(WorkerCreationRejection::CreationRejected {
209 worker: recover_worker(item),
210 reason,
211 })
212 }
213 SettledItem::Attempted(ItemSettlement::Blocked {
214 item: _,
215 prerequisite,
216 }) => match prerequisite {},
217 SettledItem::Attempted(ItemSettlement::Corrupt { item, fault }) => {
218 WorkerCreationOutcome::Rejected(WorkerCreationRejection::InterpreterCorrupt {
219 worker: recover_worker(item),
220 fault,
221 })
222 }
223 SettledItem::Unattempted(item) => {
224 WorkerCreationOutcome::Rejected(WorkerCreationRejection::InterpretationSkipped {
225 worker: recover_worker(item),
226 })
227 }
228 }
229}
230
231fn recover_worker<W>(creation: RoutedCreation<BehaviorAddr<W>, StopOnShutdown<W>>) -> W
232where
233 W: Behavior,
234{
235 let (creation, _) = creation.into_parts();
236 let (_, worker, _) = creation.into_parts();
237 worker.into_inner()
238}
239
240pub trait ActivationPlan: Send {
246 type Ready: Send;
248 type Rejection: Send;
250
251 fn activate(
253 self,
254 ) -> impl core::future::Future<Output = Result<Self::Ready, Self::Rejection>> + Send;
255}
256
257#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
259pub struct ImmediateActivation;
260
261impl ActivationPlan for ImmediateActivation {
262 type Ready = ();
263 type Rejection = Never;
264
265 fn activate(
266 self,
267 ) -> impl core::future::Future<Output = Result<Self::Ready, Self::Rejection>> + Send {
268 core::future::ready(Ok(()))
269 }
270}
271
272#[derive(Clone)]
274pub struct WorkerAttempt {
275 creation: CreationId,
276 token: Arc<()>,
277}
278
279impl WorkerAttempt {
280 pub(in crate::atomic) fn issued(creation: CreationId) -> Self {
281 Self {
282 creation,
283 token: Arc::new(()),
284 }
285 }
286
287 #[must_use]
290 pub const fn creation(&self) -> CreationId {
291 self.creation
292 }
293}
294
295impl core::fmt::Debug for WorkerAttempt {
296 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
297 formatter
298 .debug_struct("WorkerAttempt")
299 .field("creation", &self.creation)
300 .finish_non_exhaustive()
301 }
302}
303
304impl PartialEq for WorkerAttempt {
305 fn eq(&self, other: &Self) -> bool {
306 self.creation == other.creation && Arc::ptr_eq(&self.token, &other.token)
307 }
308}
309
310impl Eq for WorkerAttempt {}
311
312#[derive(Debug, Eq, PartialEq)]
314pub struct WorkerSubmission<W, P> {
315 pub(crate) worker: W,
316 pub(crate) activation: P,
317}
318
319impl<W> WorkerSubmission<W, ImmediateActivation> {
320 #[must_use]
322 pub const fn immediate(worker: W) -> Self {
323 Self {
324 worker,
325 activation: ImmediateActivation,
326 }
327 }
328}
329
330impl<W, P> WorkerSubmission<W, P> {
331 #[must_use]
333 pub const fn activated(worker: W, activation: P) -> Self {
334 Self { worker, activation }
335 }
336}
337
338#[derive(Debug, Eq, PartialEq)]
340pub struct PreparedWorker<Role, W, P> {
341 pub role: Role,
343 pub submission: WorkerSubmission<W, P>,
345}