Skip to main content

behavior_actors/atomic/stable_proxy/worker/
mod.rs

1//! Exact current worker ownership and correlation.
2
3use core::ops::ControlFlow;
4
5use behavior::{
6    Address, Behavior, BehaviorAddr, ChildCreationOutcome, CreateChild, CreationId, CreationKind,
7    CreationSequence, CreationSettlement, CreationsSettled, EndpointAddress, EstablishedActor,
8    Here, Ingress, InjectEvent, ItemSettlement, SettledItem,
9};
10
11use crate::{
12    ChildStopped, EstablishedShutdownResolved, ShutdownEstablished, ShutdownId, ShutdownRequested,
13    StopOnShutdown,
14};
15
16use crate::atomic::worker::{
17    WorkerCreationOutcome, WorkerCreationSettlement, settle_worker_creation,
18};
19use crate::atomic::{
20    ActivationPlan, InitializationAttempt, WorkerAttempt, WorkerCreationRejection,
21    WorkerInitializationReport,
22};
23
24use super::ProxyDrain;
25
26#[cfg(test)]
27mod direct_worker_admission_contract {
28    use behavior::{Behavior, BehaviorAddr, EndpointAddress};
29
30    use crate::ChildStopped;
31
32    use super::{CurrentWorker, PendingWorker, WorkerCreationSettlement};
33    use crate::atomic::{ActivationPlan, WorkerInitializationReport};
34
35    #[expect(dead_code, reason = "compile contract for exact worker input")]
36    fn current_stop<W>(
37        worker: &CurrentWorker<W>,
38        input: ChildStopped<BehaviorAddr<W>>,
39    ) -> Result<ChildStopped<BehaviorAddr<W>>, ChildStopped<BehaviorAddr<W>>>
40    where
41        W: Behavior,
42        BehaviorAddr<W>: EndpointAddress,
43    {
44        worker.admit_stop(input)
45    }
46
47    #[expect(dead_code, reason = "compile contract for exact worker input")]
48    fn current_initialization<W, P>(
49        worker: &CurrentWorker<W>,
50        input: WorkerInitializationReport<W, P>,
51    ) -> Result<WorkerInitializationReport<W, P>, WorkerInitializationReport<W, P>>
52    where
53        W: Behavior,
54        P: ActivationPlan,
55        BehaviorAddr<W>: EndpointAddress,
56    {
57        worker.admit_initialization(input)
58    }
59
60    #[expect(dead_code, reason = "compile contract for exact worker input")]
61    fn pending_stop<W, P>(
62        worker: &PendingWorker<P>,
63        input: ChildStopped<BehaviorAddr<W>>,
64    ) -> Result<ChildStopped<BehaviorAddr<W>>, ChildStopped<BehaviorAddr<W>>>
65    where
66        W: Behavior,
67        BehaviorAddr<W>: EndpointAddress,
68    {
69        worker.admit_stop(input)
70    }
71
72    #[expect(dead_code, reason = "compile contract for exact worker input")]
73    fn pending_creation<W, P>(
74        worker: &PendingWorker<P>,
75        input: WorkerCreationSettlement<W>,
76    ) -> Result<WorkerCreationSettlement<W>, WorkerCreationSettlement<W>>
77    where
78        W: Behavior,
79        BehaviorAddr<W>: EndpointAddress,
80    {
81        worker.admit_creation(input)
82    }
83}
84
85/// Complete result of one owner-submitted worker start.
86pub enum WorkerStartResult<W, P>
87where
88    W: Behavior,
89    P: ActivationPlan,
90    BehaviorAddr<W>: EndpointAddress,
91    StopOnShutdown<W>: Behavior<Protocol = W::Protocol, Ph = W::Ph, Birth = W::Birth>,
92{
93    /// Worker creation rejected before a worker was committed.
94    CreationRejected {
95        rejection: WorkerCreationRejection<W>,
96        activation: P,
97        stopped: Option<ChildStopped<BehaviorAddr<W>>>,
98    },
99    /// The committed worker became ready and may receive service commands.
100    Ready {
101        attempt: WorkerAttempt,
102        readiness: P::Ready,
103    },
104    /// A committed worker could not become available.
105    Unavailable {
106        attempt: WorkerAttempt,
107        drain: ProxyDrain<W, P>,
108    },
109}
110
111pub(in super::super) struct PendingWorker<P> {
112    attempt: WorkerAttempt,
113    initialization: InitializationAttempt,
114    kind: CreationKind,
115    activation: P,
116}
117
118impl<P> PendingWorker<P> {
119    pub(in super::super) fn birth<W>(
120        creations: &mut CreationSequence,
121        worker: W,
122        activation: P,
123    ) -> Result<
124        (
125            Self,
126            CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
127        ),
128        (W, P),
129    >
130    where
131        W: Behavior,
132    {
133        let Some(creation) = creations.issue() else {
134            return Err((worker, activation));
135        };
136        let request = CreateChild::birth(creation, StopOnShutdown::new(worker));
137        Ok(Self::await_creation(activation, request))
138    }
139
140    pub(in super::super) fn replacement<W>(
141        creations: &mut CreationSequence,
142        previous: CreationId,
143        worker: W,
144        activation: P,
145    ) -> Result<
146        (
147            Self,
148            CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
149        ),
150        (W, P),
151    >
152    where
153        W: Behavior,
154    {
155        let Some(creation) = creations.issue() else {
156            return Err((worker, activation));
157        };
158        let request = CreateChild::replacement(creation, previous, StopOnShutdown::new(worker));
159        Ok(Self::await_creation(activation, request))
160    }
161
162    fn await_creation<W>(
163        activation: P,
164        request: CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
165    ) -> (
166        Self,
167        CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
168    )
169    where
170        W: Behavior,
171    {
172        let creation = request.id();
173        let kind = request.kind();
174        let attempt = WorkerAttempt::issued(creation);
175        let initialization = InitializationAttempt::issued(&attempt);
176        (
177            Self {
178                attempt,
179                initialization,
180                kind,
181                activation,
182            },
183            request,
184        )
185    }
186
187    pub(in super::super) const fn creation(&self) -> CreationId {
188        self.attempt.creation()
189    }
190
191    pub(in super::super) fn cancel<W>(
192        self,
193        creation: CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
194    ) -> (WorkerAttempt, W, P)
195    where
196        W: Behavior,
197    {
198        let (_, worker, _) = creation.into_parts();
199        (self.attempt, worker.into_inner(), self.activation)
200    }
201
202    pub(in super::super) fn admit_stop<A>(
203        &self,
204        stopped: ChildStopped<A>,
205    ) -> Result<ChildStopped<A>, ChildStopped<A>>
206    where
207        A: Address,
208    {
209        if stopped.child == self.attempt.creation() {
210            Ok(stopped)
211        } else {
212            Err(stopped)
213        }
214    }
215
216    fn admit_creation<W>(
217        &self,
218        settlement: WorkerCreationSettlement<W>,
219    ) -> Result<WorkerCreationSettlement<W>, WorkerCreationSettlement<W>>
220    where
221        W: Behavior,
222        BehaviorAddr<W>: EndpointAddress,
223    {
224        let received = match &settlement {
225            SettledItem::Unattempted(creation) => (creation.id(), creation.kind()),
226            SettledItem::Attempted(ItemSettlement::Rejected { item, .. })
227            | SettledItem::Attempted(ItemSettlement::Blocked { item, .. })
228            | SettledItem::Attempted(ItemSettlement::Corrupt { item, .. }) => {
229                (item.id(), item.kind())
230            }
231            SettledItem::Attempted(ItemSettlement::Accepted(creation)) => match creation {
232                ChildCreationOutcome::Established(child) => (child.id(), child.kind()),
233                ChildCreationOutcome::InitializationRejected { creation, .. }
234                | ChildCreationOutcome::InitializationPanicked { creation }
235                | ChildCreationOutcome::HostRejected { creation, .. } => {
236                    (creation.id(), creation.kind())
237                }
238            },
239        };
240        if received == (self.attempt.creation(), self.kind) {
241            Ok(settlement)
242        } else {
243            Err(settlement)
244        }
245    }
246}
247
248enum WorkerShutdown<P>
249where
250    P: behavior::Protocol,
251{
252    Requested(ShutdownId),
253    Settled(EstablishedShutdownResolved<P>),
254}
255
256pub(in super::super) struct WorkerStopping<W>
257where
258    W: Behavior,
259    BehaviorAddr<W>: EndpointAddress,
260{
261    worker: CurrentWorker<W>,
262    shutdown: WorkerShutdown<W::Protocol>,
263    stopped: Option<ChildStopped<BehaviorAddr<W>>>,
264}
265
266pub(in super::super) struct StoppedWorker<W>
267where
268    W: Behavior,
269    BehaviorAddr<W>: EndpointAddress,
270{
271    pub(in super::super) worker: CurrentWorker<W>,
272    pub(in super::super) shutdown: EstablishedShutdownResolved<W::Protocol>,
273    pub(in super::super) stopped: ChildStopped<BehaviorAddr<W>>,
274}
275
276impl<W> WorkerStopping<W>
277where
278    W: Behavior,
279    BehaviorAddr<W>: EndpointAddress,
280{
281    pub(in super::super) const fn worker(&self) -> &CurrentWorker<W> {
282        &self.worker
283    }
284
285    pub(in super::super) const fn stopped(&self) -> Option<&ChildStopped<BehaviorAddr<W>>> {
286        self.stopped.as_ref()
287    }
288
289    pub(in super::super) fn begin(
290        worker: CurrentWorker<W>,
291    ) -> (Self, ShutdownEstablished<StopOnShutdown<W>, Here>)
292    where
293        StopOnShutdown<W>: Behavior<Protocol = W::Protocol>,
294        <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
295    {
296        let shutdown = ShutdownEstablished::new(
297            ShutdownId(worker.attempt.creation().get()),
298            worker.actor.clone(),
299            Ingress::new(),
300        );
301        (
302            Self {
303                worker,
304                shutdown: WorkerShutdown::Requested(shutdown.id),
305                stopped: None,
306            },
307            shutdown,
308        )
309    }
310
311    pub(in super::super) fn shutdown_resolved(
312        self,
313        input: EstablishedShutdownResolved<W::Protocol>,
314    ) -> Result<ControlFlow<StoppedWorker<W>, Self>, (Self, EstablishedShutdownResolved<W::Protocol>)>
315    {
316        let Self {
317            worker,
318            shutdown,
319            stopped,
320        } = self;
321        match shutdown {
322            WorkerShutdown::Requested(expected) if input.id() == expected => Ok(Self {
323                worker,
324                shutdown: WorkerShutdown::Settled(input),
325                stopped,
326            }
327            .complete()),
328            shutdown => Err((
329                Self {
330                    worker,
331                    shutdown,
332                    stopped,
333                },
334                input,
335            )),
336        }
337    }
338
339    pub(in super::super) fn worker_stopped(
340        self,
341        input: ChildStopped<BehaviorAddr<W>>,
342    ) -> Result<ControlFlow<StoppedWorker<W>, Self>, (Self, ChildStopped<BehaviorAddr<W>>)> {
343        let Self {
344            worker,
345            shutdown,
346            stopped,
347        } = self;
348        match stopped {
349            None => match worker.admit_stop(input) {
350                Ok(stopped) => Ok(Self {
351                    worker,
352                    shutdown,
353                    stopped: Some(stopped),
354                }
355                .complete()),
356                Err(input) => Err((
357                    Self {
358                        worker,
359                        shutdown,
360                        stopped: None,
361                    },
362                    input,
363                )),
364            },
365            Some(stopped) => Err((
366                Self {
367                    worker,
368                    shutdown,
369                    stopped: Some(stopped),
370                },
371                input,
372            )),
373        }
374    }
375
376    pub(in super::super) fn into_stopped_before_shutdown(
377        self,
378    ) -> Result<(CurrentWorker<W>, ShutdownId, ChildStopped<BehaviorAddr<W>>), Self> {
379        match (self.shutdown, self.stopped) {
380            (WorkerShutdown::Requested(shutdown), Some(stopped)) => {
381                Ok((self.worker, shutdown, stopped))
382            }
383            (shutdown, stopped) => Err(Self {
384                worker: self.worker,
385                shutdown,
386                stopped,
387            }),
388        }
389    }
390
391    fn complete(self) -> ControlFlow<StoppedWorker<W>, Self> {
392        match (self.shutdown, self.stopped) {
393            (WorkerShutdown::Settled(shutdown), Some(stopped)) => {
394                ControlFlow::Break(StoppedWorker {
395                    worker: self.worker,
396                    shutdown,
397                    stopped,
398                })
399            }
400            (shutdown, stopped) => ControlFlow::Continue(Self {
401                worker: self.worker,
402                shutdown,
403                stopped,
404            }),
405        }
406    }
407}
408
409pub(in super::super) struct CurrentWorker<W>
410where
411    W: Behavior,
412    BehaviorAddr<W>: EndpointAddress,
413{
414    pub(in super::super) attempt: WorkerAttempt,
415    pub(in super::super) initialization: InitializationAttempt,
416    pub(in super::super) actor: EstablishedActor<StopOnShutdown<W>>,
417}
418
419impl<W> CurrentWorker<W>
420where
421    W: Behavior,
422    BehaviorAddr<W>: EndpointAddress,
423{
424    pub(in super::super) fn admit_stop(
425        &self,
426        stopped: ChildStopped<BehaviorAddr<W>>,
427    ) -> Result<ChildStopped<BehaviorAddr<W>>, ChildStopped<BehaviorAddr<W>>> {
428        if stopped.child == self.attempt.creation() {
429            Ok(stopped)
430        } else {
431            Err(stopped)
432        }
433    }
434
435    pub(in super::super) fn admit_initialization<P>(
436        &self,
437        input: WorkerInitializationReport<W, P>,
438    ) -> Result<WorkerInitializationReport<W, P>, WorkerInitializationReport<W, P>>
439    where
440        P: ActivationPlan,
441    {
442        match (input.worker(), input.initialization()) {
443            (worker, initialization)
444                if worker == &self.attempt && initialization == &self.initialization =>
445            {
446                Ok(input)
447            }
448            _ => Err(input),
449        }
450    }
451}
452
453pub(in super::super) enum WorkerCreation<W, P>
454where
455    W: Behavior,
456    BehaviorAddr<W>: EndpointAddress,
457{
458    Initializing {
459        worker: CurrentWorker<W>,
460        activation: P,
461        stopped: Option<ChildStopped<BehaviorAddr<W>>>,
462    },
463    Rejected {
464        rejection: WorkerCreationRejection<W>,
465        activation: P,
466        stopped: Option<ChildStopped<BehaviorAddr<W>>>,
467    },
468    Unexpected {
469        worker: PendingWorker<P>,
470        stopped: Option<ChildStopped<BehaviorAddr<W>>>,
471        workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
472    },
473}
474
475impl<P> PendingWorker<P> {
476    pub(in super::super) fn created<W>(
477        self,
478        workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
479        stopped: Option<ChildStopped<BehaviorAddr<W>>>,
480    ) -> WorkerCreation<W, P>
481    where
482        W: Behavior,
483        BehaviorAddr<W>: EndpointAddress,
484        StopOnShutdown<W>:
485            Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
486    {
487        match workers.into_settlement() {
488            CreationSettlement::Rejected { creations, reason } => {
489                let creation = match creations.into_one() {
490                    Ok(creation) => creation,
491                    Err(creations) => {
492                        return WorkerCreation::Unexpected {
493                            worker: self,
494                            stopped,
495                            workers: CreationsSettled::new(CreationSettlement::Rejected {
496                                creations,
497                                reason,
498                            }),
499                        };
500                    }
501                };
502                if (creation.id(), creation.kind()) != (self.creation(), self.kind) {
503                    return WorkerCreation::Unexpected {
504                        worker: self,
505                        stopped,
506                        workers: CreationsSettled::new(CreationSettlement::Rejected {
507                            creations: behavior::Creations::one(creation),
508                            reason,
509                        }),
510                    };
511                }
512                let (_, worker, _) = creation.into_parts();
513                WorkerCreation::Rejected {
514                    rejection: WorkerCreationRejection::NamespaceExhausted {
515                        worker: worker.into_inner(),
516                    },
517                    activation: self.activation,
518                    stopped,
519                }
520            }
521            CreationSettlement::Settled(settlements) => {
522                let settlement = match settlements.into_one() {
523                    Ok(settlement) => settlement,
524                    Err(settlements) => {
525                        return WorkerCreation::Unexpected {
526                            worker: self,
527                            stopped,
528                            workers: CreationsSettled::new(CreationSettlement::Settled(
529                                settlements,
530                            )),
531                        };
532                    }
533                };
534                self.created_one(settlement, stopped)
535            }
536            settlement @ CreationSettlement::Corrupt { .. } => WorkerCreation::Unexpected {
537                worker: self,
538                stopped,
539                workers: CreationsSettled::new(settlement),
540            },
541        }
542    }
543
544    fn created_one<W>(
545        self,
546        worker: WorkerCreationSettlement<W>,
547        stopped: Option<ChildStopped<BehaviorAddr<W>>>,
548    ) -> WorkerCreation<W, P>
549    where
550        W: Behavior,
551        BehaviorAddr<W>: EndpointAddress,
552        StopOnShutdown<W>:
553            Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
554    {
555        match self.admit_creation(worker) {
556            Ok(settlement) => match settle_worker_creation(settlement) {
557                WorkerCreationOutcome::Established(actor) => WorkerCreation::Initializing {
558                    worker: CurrentWorker {
559                        attempt: self.attempt,
560                        initialization: self.initialization,
561                        actor,
562                    },
563                    activation: self.activation,
564                    stopped,
565                },
566                WorkerCreationOutcome::Rejected(rejection) => WorkerCreation::Rejected {
567                    rejection,
568                    activation: self.activation,
569                    stopped,
570                },
571            },
572            Err(settlement) => WorkerCreation::Unexpected {
573                worker: self,
574                stopped,
575                workers: CreationsSettled::new(CreationSettlement::Settled(
576                    behavior::Creations::one(settlement),
577                )),
578            },
579        }
580    }
581}