Skip to main content

behavior_actors/atomic/worker/
mod.rs

1//! Worker submissions and shared start contracts used by atomic aggregates.
2
3use 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/// Initialization actions retained together with an uncommitted worker.
53#[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/// An uncommitted worker returned together with initialization actions that
62/// the host rejected before interpreting.
63#[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    /// Transfer both affine values to Bombay's terminal custodian.
85    #[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
103/// Complete rejection returned by one worker-creation attempt.
104pub enum WorkerCreationRejection<W>
105where
106    W: Behavior,
107{
108    /// Bombay could not route the creator's child creation.
109    NamespaceExhausted { worker: W },
110    /// The worker's pure initialization transition rejected.
111    WorkerRejected { worker: W, error: W::Error },
112    /// The worker's pure initialization panicked before host commitment.
113    WorkerPanicked { worker: W },
114    /// Host establishment rejected after initialization produced actions.
115    HostRejected {
116        recovery: WorkerRecovery<W>,
117        reason: CreationRejection,
118    },
119    /// The creation capability lawfully rejected and returned the worker.
120    CreationRejected {
121        worker: W,
122        reason: CreationRejection,
123    },
124    /// The interpreter violated its creation contract and returned the worker.
125    InterpreterCorrupt { worker: W, fault: InterpreterFault },
126    /// Product traversal stopped before attempting this worker creation.
127    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
240/// Concrete work required before a newly created worker may serve commands.
241///
242/// Bombay interprets the owned plan after creation and initialization settle.
243/// Atomic aggregates never invoke application activation work in their pure
244/// transitions.
245pub trait ActivationPlan: Send {
246    /// Evidence returned when activation makes the worker ready.
247    type Ready: Send;
248    /// Exact application rejection returned by activation.
249    type Rejection: Send;
250
251    /// Perform the selected activation work outside the worker behavior.
252    fn activate(
253        self,
254    ) -> impl core::future::Future<Output = Result<Self::Ready, Self::Rejection>> + Send;
255}
256
257/// A worker that needs no application activation work after initialization.
258#[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/// Exact non-forgeable evidence for one worker creation attempt.
273#[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    /// Creator-local correlation for this worker attempt. This number is not
288    /// actor identity or proof that creation committed.
289    #[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/// Complete affine input for one worker start or replacement.
313#[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    /// Prepare a worker that needs no application activation work.
321    #[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    /// Prepare a worker with one concrete owned activation plan.
332    #[must_use]
333    pub const fn activated(worker: W, activation: P) -> Self {
334        Self { worker, activation }
335    }
336}
337
338/// One semantic role and its complete prepared worker submission.
339#[derive(Debug, Eq, PartialEq)]
340pub struct PreparedWorker<Role, W, P> {
341    /// Application role owned by this worker relationship.
342    pub role: Role,
343    /// Worker behavior and activation work prepared for the role.
344    pub submission: WorkerSubmission<W, P>,
345}