Skip to main content

behavior_actors/atomic/stable_proxy/
protocol.rs

1//! StableProxy activation policy, owner control, outcomes, and diagnostics.
2
3use behavior::{
4    Behavior, BehaviorAddr, Births, ChildInputIngress, CreationsSettled, EndpointAddress,
5    EventIngress, Here, InjectEvent, Protocol, RecoverEvent, User, UserEvent,
6};
7
8use crate::atomic::{
9    ActivationPermit, ActivationPlan, ActivationStartRejection, BeginActivation,
10    ImmediateActivation, WorkerInitializationFailure,
11};
12use crate::{ChildStopped, EstablishedShutdownResolved, StopOnShutdown, WorkerSubmission};
13
14use super::{
15    StableProxy, WorkerActivation, WorkerAttempt, WorkerInitializationReport, WorkerStartResult,
16};
17
18/// Owner-only operation accepted by [`super::StableProxy`].
19pub struct ProxyControl<W, P> {
20    pub(super) command: ProxyCommand<W, P>,
21}
22
23pub(super) enum ProxyCommand<W, P> {
24    Start(WorkerSubmission<W, P>),
25    Replace(WorkerSubmission<W, P>),
26    Shutdown,
27}
28
29impl<W> ProxyControl<W, ImmediateActivation> {
30    /// Start an immediately ready worker without an unused plan argument.
31    #[must_use]
32    pub const fn start(worker: W) -> Self {
33        Self {
34            command: ProxyCommand::Start(WorkerSubmission::immediate(worker)),
35        }
36    }
37
38    /// Replace with an immediately ready worker.
39    #[must_use]
40    pub const fn replace(worker: W) -> Self {
41        Self {
42            command: ProxyCommand::Replace(WorkerSubmission::immediate(worker)),
43        }
44    }
45}
46
47impl<W, P> ProxyControl<W, P> {
48    /// Shut down the proxy without an unused worker or activation input.
49    #[must_use]
50    pub const fn shutdown() -> Self {
51        Self {
52            command: ProxyCommand::Shutdown,
53        }
54    }
55
56    /// Start a worker with application-defined activation work.
57    #[must_use]
58    pub const fn start_with(worker: W, activation: P) -> Self {
59        Self {
60            command: ProxyCommand::Start(WorkerSubmission::activated(worker, activation)),
61        }
62    }
63
64    /// Replace with a worker using application-defined activation work.
65    #[must_use]
66    pub const fn replace_with(worker: W, activation: P) -> Self {
67        Self {
68            command: ProxyCommand::Replace(WorkerSubmission::activated(worker, activation)),
69        }
70    }
71}
72
73/// Observable lifecycle phase carried by proxy rejections and diagnostics.
74///
75/// This value grants no transition or routing authority.
76#[derive(Clone, Copy, Debug, Eq, PartialEq)]
77pub enum ProxyPhase {
78    Dormant,
79    Creating,
80    Initializing,
81    Activating,
82    Ready,
83    EmptyInitial,
84    EmptyAfter,
85    Replacing,
86    ReturningWorker,
87    ShuttingDown,
88    Stopped,
89}
90
91/// Complete disposition of one initial worker submission.
92pub enum InitialWorkerOutcome<W, P>
93where
94    W: Behavior,
95    P: ActivationPlan,
96    BehaviorAddr<W>: EndpointAddress,
97    StopOnShutdown<W>: Behavior<Protocol = W::Protocol, Ph = W::Ph, Birth = W::Birth>,
98{
99    /// Another initial submission already owns the proxy.
100    Overlap {
101        worker: W,
102        activation: P,
103        phase: ProxyPhase,
104    },
105    /// The proxy cannot reserve another non-reused worker attempt.
106    WorkerAttemptsExhausted { worker: W, activation: P },
107    /// The worker-start operation reached one complete semantic result.
108    Resolved { result: WorkerStartResult<W, P> },
109}
110
111/// Complete disposition of one submitted replacement worker.
112pub enum ReplacementOutcome<W, P>
113where
114    W: Behavior,
115    P: ActivationPlan,
116    BehaviorAddr<W>: EndpointAddress,
117    StopOnShutdown<W>: Behavior<Protocol = W::Protocol, Ph = W::Ph, Birth = W::Birth>,
118{
119    /// The current proxy phase cannot accept another replacement.
120    NotReplaceable {
121        worker: W,
122        activation: P,
123        phase: ProxyPhase,
124    },
125    /// The proxy cannot reserve another non-reused worker attempt.
126    WorkerAttemptsExhausted {
127        replaces: WorkerAttempt,
128        worker: W,
129        activation: P,
130    },
131    /// Owner shutdown returned a successor before its creation was emitted.
132    CancelledBeforeBirth {
133        replaces: WorkerAttempt,
134        worker: W,
135        activation: P,
136    },
137    /// The successor start produced one complete result.
138    Resolved {
139        replaces: WorkerAttempt,
140        result: WorkerStartResult<W, P>,
141    },
142}
143
144/// Lifecycle input returned without changing the current proxy state.
145pub enum ProxyDiagnostic<W, P>
146where
147    W: Behavior,
148    P: ActivationPlan,
149    BehaviorAddr<W>: EndpointAddress,
150{
151    /// A worker-creation settlement did not belong to the current start.
152    UnexpectedWorkerStart {
153        phase: ProxyPhase,
154        workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
155    },
156    /// A worker stop did not belong to the worker expected by this phase.
157    UnexpectedWorkerStop {
158        phase: ProxyPhase,
159        stopped: ChildStopped<BehaviorAddr<W>>,
160    },
161    /// A child-host result did not belong to the worker expected by this phase.
162    UnexpectedWorkerInitialization {
163        phase: ProxyPhase,
164        initialization: WorkerInitializationReport<W, P>,
165    },
166    /// An activation input did not belong to the worker expected by this phase.
167    UnexpectedWorkerActivation {
168        phase: ProxyPhase,
169        activation: WorkerActivation<W, P>,
170    },
171    /// An exact shutdown resolution did not belong to the current worker drain.
172    UnexpectedWorkerShutdown {
173        phase: ProxyPhase,
174        shutdown: EstablishedShutdownResolved<W::Protocol>,
175    },
176}
177
178/// Complete report emitted to the structural owner of a stable proxy.
179pub enum ProxyOutcome<W, P>
180where
181    W: Behavior,
182    P: ActivationPlan,
183    StopOnShutdown<W>: Behavior<Protocol = W::Protocol, Ph = W::Ph, Birth = W::Birth>,
184    <W::Protocol as Protocol>::Addr: EndpointAddress,
185{
186    /// Result of the first submitted worker.
187    Initial { outcome: InitialWorkerOutcome<W, P> },
188    /// Result of one submitted replacement worker.
189    Replacement { outcome: ReplacementOutcome<W, P> },
190    /// The exact ready worker stopped and left the stable service unavailable.
191    WorkerStopped {
192        worker: WorkerAttempt,
193        stopped: ChildStopped<BehaviorAddr<W>>,
194    },
195    /// A service command arrived while no worker was ready.
196    Unavailable {
197        sender: BehaviorAddr<W>,
198        phase: ProxyPhase,
199        command: <W::Protocol as Protocol>::Msg,
200    },
201}
202
203/// Complete retained values proving how one unavailable worker was returned.
204#[must_use = "a proxy return retains affine worker lifecycle values"]
205pub enum ProxyDrain<W, P>
206where
207    W: Behavior,
208    BehaviorAddr<W>: EndpointAddress,
209    P: ActivationPlan,
210{
211    /// Pure initialization rejected after the worker was committed.
212    InitializationRejected {
213        failure: WorkerInitializationFailure,
214        activation: P,
215        shutdown: Option<EstablishedShutdownResolved<W::Protocol>>,
216        stopped: ChildStopped<BehaviorAddr<W>>,
217    },
218    /// Initialization stopped after the worker was committed.
219    InitializationStopped {
220        activation: P,
221        observed: Option<ChildStopped<BehaviorAddr<W>>>,
222        returned: ChildStopped<BehaviorAddr<W>>,
223    },
224    /// Initialization completed, but an exact worker stop prevented activation.
225    InitializationCompleted {
226        permit: ActivationPermit<W>,
227        activation: P,
228        stopped: ChildStopped<BehaviorAddr<W>>,
229    },
230    /// The activation request was rejected before the plan ran.
231    ActivationStartRejected {
232        request: BeginActivation<W, P>,
233        reason: ActivationStartRejection,
234        shutdown: Option<EstablishedShutdownResolved<W::Protocol>>,
235        stopped: ChildStopped<BehaviorAddr<W>>,
236    },
237    /// The activation plan rejected after its request was accepted.
238    ActivationRejected {
239        rejection: P::Rejection,
240        shutdown: Option<EstablishedShutdownResolved<W::Protocol>>,
241        stopped: ChildStopped<BehaviorAddr<W>>,
242    },
243    /// Activation completed, but an exact worker stop prevented availability.
244    ActivationCompleted {
245        readiness: P::Ready,
246        stopped: ChildStopped<BehaviorAddr<W>>,
247    },
248}
249
250#[doc(hidden)]
251pub enum ProxyEvent<W, P>
252where
253    W: Behavior,
254    P: ActivationPlan,
255    BehaviorAddr<W>: EndpointAddress,
256{
257    Service(User<BehaviorAddr<W>, <W::Protocol as Protocol>::Msg>),
258    Owner(ProxyControl<W, P>),
259    WorkerCreationsSettled(CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>),
260    WorkerInitialization(WorkerInitializationReport<W, P>),
261    WorkerActivationReported(WorkerActivation<W, P>),
262    WorkerShutdownResolved(EstablishedShutdownResolved<W::Protocol>),
263    WorkerStopped(ChildStopped<BehaviorAddr<W>>),
264}
265
266impl<W, P>
267    EventIngress<Births<StopOnShutdown<W>>, CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>>
268    for ProxyEvent<W, P>
269where
270    W: Behavior,
271    P: ActivationPlan,
272    BehaviorAddr<W>: EndpointAddress,
273{
274    fn ingress(input: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>) -> Self {
275        Self::WorkerCreationsSettled(input)
276    }
277}
278
279impl<W, P> InjectEvent<CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>, Here>
280    for ProxyEvent<W, P>
281where
282    W: Behavior,
283    P: ActivationPlan,
284    BehaviorAddr<W>: EndpointAddress,
285{
286    fn inject_at(input: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>) -> Self {
287        Self::WorkerCreationsSettled(input)
288    }
289}
290
291impl<W, P> UserEvent for ProxyEvent<W, P>
292where
293    W: Behavior,
294    P: ActivationPlan,
295    BehaviorAddr<W>: EndpointAddress,
296{
297    type Addr = BehaviorAddr<W>;
298    type Message = <W::Protocol as Protocol>::Msg;
299
300    fn user(from: Self::Addr, message: Self::Message) -> Self {
301        Self::Service(User::new(from, message))
302    }
303
304    fn into_user(self) -> Result<User<Self::Addr, Self::Message>, Self> {
305        match self {
306            Self::Service(service) => Ok(service),
307            other => Err(other),
308        }
309    }
310}
311
312impl<W, P> InjectEvent<ProxyControl<W, P>, Here> for ProxyEvent<W, P>
313where
314    W: Behavior,
315    P: ActivationPlan,
316    BehaviorAddr<W>: EndpointAddress,
317{
318    fn inject_at(input: ProxyControl<W, P>) -> Self {
319        Self::Owner(input)
320    }
321}
322
323impl<W, P> ChildInputIngress<StableProxy<W, P>, ProxyControl<W, P>> for ProxyEvent<W, P>
324where
325    W: Behavior,
326    P: ActivationPlan,
327    BehaviorAddr<W>: EndpointAddress,
328{
329    fn child_input(input: ProxyControl<W, P>) -> Self {
330        Self::Owner(input)
331    }
332}
333
334impl<W, P> RecoverEvent<ProxyControl<W, P>, Here> for ProxyEvent<W, P>
335where
336    W: Behavior,
337    P: ActivationPlan,
338    BehaviorAddr<W>: EndpointAddress,
339{
340    fn recover(event: Self) -> Result<ProxyControl<W, P>, Self> {
341        match event {
342            Self::Owner(control) => Ok(control),
343            other => Err(other),
344        }
345    }
346}
347
348impl<W, P> RecoverEvent<CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>, Here>
349    for ProxyEvent<W, P>
350where
351    W: Behavior,
352    P: ActivationPlan,
353    BehaviorAddr<W>: EndpointAddress,
354{
355    fn recover(event: Self) -> Result<CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>, Self> {
356        match event {
357            Self::WorkerCreationsSettled(workers) => Ok(workers),
358            other => Err(other),
359        }
360    }
361}
362
363impl<W, P> InjectEvent<ChildStopped<BehaviorAddr<W>>, Here> for ProxyEvent<W, P>
364where
365    W: Behavior,
366    P: ActivationPlan,
367    BehaviorAddr<W>: EndpointAddress,
368{
369    fn inject_at(input: ChildStopped<BehaviorAddr<W>>) -> Self {
370        Self::WorkerStopped(input)
371    }
372}
373
374impl<W, P> RecoverEvent<ChildStopped<BehaviorAddr<W>>, Here> for ProxyEvent<W, P>
375where
376    W: Behavior,
377    P: ActivationPlan,
378    BehaviorAddr<W>: EndpointAddress,
379{
380    fn recover(event: Self) -> Result<ChildStopped<BehaviorAddr<W>>, Self> {
381        match event {
382            Self::WorkerStopped(stopped) => Ok(stopped),
383            other => Err(other),
384        }
385    }
386}
387
388impl<W, P> InjectEvent<WorkerInitializationReport<W, P>, Here> for ProxyEvent<W, P>
389where
390    W: Behavior,
391    P: ActivationPlan,
392    BehaviorAddr<W>: EndpointAddress,
393{
394    fn inject_at(input: WorkerInitializationReport<W, P>) -> Self {
395        Self::WorkerInitialization(input)
396    }
397}
398
399impl<W, P> RecoverEvent<WorkerInitializationReport<W, P>, Here> for ProxyEvent<W, P>
400where
401    W: Behavior,
402    P: ActivationPlan,
403    BehaviorAddr<W>: EndpointAddress,
404{
405    fn recover(event: Self) -> Result<WorkerInitializationReport<W, P>, Self> {
406        match event {
407            Self::WorkerInitialization(initialization) => Ok(initialization),
408            other => Err(other),
409        }
410    }
411}
412
413impl<W, P> InjectEvent<WorkerActivation<W, P>, Here> for ProxyEvent<W, P>
414where
415    W: Behavior,
416    P: ActivationPlan,
417    BehaviorAddr<W>: EndpointAddress,
418{
419    fn inject_at(input: WorkerActivation<W, P>) -> Self {
420        Self::WorkerActivationReported(input)
421    }
422}
423
424impl<W, P> RecoverEvent<WorkerActivation<W, P>, Here> for ProxyEvent<W, P>
425where
426    W: Behavior,
427    P: ActivationPlan,
428    BehaviorAddr<W>: EndpointAddress,
429{
430    fn recover(event: Self) -> Result<WorkerActivation<W, P>, Self> {
431        match event {
432            Self::WorkerActivationReported(activation) => Ok(activation),
433            other => Err(other),
434        }
435    }
436}
437
438impl<W, P> InjectEvent<EstablishedShutdownResolved<W::Protocol>, Here> for ProxyEvent<W, P>
439where
440    W: Behavior,
441    P: ActivationPlan,
442    BehaviorAddr<W>: EndpointAddress,
443{
444    fn inject_at(input: EstablishedShutdownResolved<W::Protocol>) -> Self {
445        Self::WorkerShutdownResolved(input)
446    }
447}
448
449impl<W, P> RecoverEvent<EstablishedShutdownResolved<W::Protocol>, Here> for ProxyEvent<W, P>
450where
451    W: Behavior,
452    P: ActivationPlan,
453    BehaviorAddr<W>: EndpointAddress,
454{
455    fn recover(event: Self) -> Result<EstablishedShutdownResolved<W::Protocol>, Self> {
456        match event {
457            Self::WorkerShutdownResolved(shutdown) => Ok(shutdown),
458            other => Err(other),
459        }
460    }
461}