Skip to main content

behavior_actors/atomic/worker/
activation.rs

1//! Worker activation authorized by successful initialization.
2
3use std::sync::Arc;
4
5use core::marker::PhantomData;
6
7use behavior::{
8    ActionItem, Behavior, BehaviorAddr, EndpointAddress, EstablishedRecipient, Here,
9    InterpretationProgress, InterpreterRequest, ItemSettlement, Never, ReturnsToEmitter,
10    finish_item, prepare_item,
11};
12
13use super::initialization::{ActivationPermit, InitializationAttempt};
14use super::{ActivationPlan, WorkerAttempt};
15
16/// The original worker/plan-typed activation correlation issued with a permit.
17///
18/// Readiness and rejection are absent from this owner. Its original private
19/// issuer remains inside BeginActivation. A caller cannot fabricate a grant,
20/// but into_parts returns the original permit and supports reissuing a request.
21/// This owner does not impose a one-mint-ever policy on that affine permit.
22pub struct WorkerActivationGrant<W, P> {
23    worker: WorkerAttempt,
24    token: Arc<()>,
25    worker_plan: PhantomData<fn() -> (W, P)>,
26}
27
28impl<W, P> Clone for WorkerActivationGrant<W, P> {
29    fn clone(&self) -> Self {
30        Self {
31            worker: self.worker.clone(),
32            token: Arc::clone(&self.token),
33            worker_plan: PhantomData,
34        }
35    }
36}
37
38impl<W, P> WorkerActivationGrant<W, P> {
39    fn issued(worker: &WorkerAttempt) -> Self {
40        Self {
41            worker: worker.clone(),
42            token: Arc::new(()),
43            worker_plan: PhantomData,
44        }
45    }
46
47    /// Original worker correlation carried by this grant.
48    #[must_use]
49    pub fn worker(&self) -> WorkerAttempt {
50        self.worker.clone()
51    }
52}
53
54impl<W, P> PartialEq for WorkerActivationGrant<W, P> {
55    fn eq(&self, other: &Self) -> bool {
56        self.worker == other.worker && Arc::ptr_eq(&self.token, &other.token)
57    }
58}
59
60impl<W, P> Eq for WorkerActivationGrant<W, P> {}
61
62impl<W, P> WorkerActivationGrant<W, P>
63where
64    W: Behavior,
65    BehaviorAddr<W>: EndpointAddress,
66    P: ActivationPlan,
67{
68    /// Form the existing input from the original concrete plan's acquired reply.
69    ///
70    /// The independent correlation stays owned, as when started() produces a
71    /// correlated input. Acceptance remains the actor's existing policy.
72    /// The receiving actor applies its existing correlation policy. Resolving a
73    /// reply does not consume the initialization permit.
74    #[must_use]
75    pub fn resolve(&self, reply: Result<P::Ready, P::Rejection>) -> WorkerActivation<W, P> {
76        WorkerActivation {
77            worker: self.worker(),
78            activation: self.clone(),
79            outcome: match reply {
80                Ok(readiness) => WorkerActivationOutcome::Ready(readiness),
81                Err(rejection) => WorkerActivationOutcome::Rejected(rejection),
82            },
83            worker_type: PhantomData,
84        }
85    }
86}
87
88/// Exact reason activation work was not admitted by its owner.
89#[derive(Clone, Copy, Debug, Eq, PartialEq)]
90pub enum ActivationStartRejection {
91    OwnerStopped,
92}
93
94/// One activation request consuming an initialized worker permit and plan.
95///
96/// ```compile_fail,E0382
97/// fn duplicate<W, P>(request: behavior_actors::atomic::BeginActivation<W, P>)
98/// where
99///     W: behavior::Behavior,
100///     behavior::BehaviorAddr<W>: behavior::EndpointAddress,
101///     P: behavior_actors::atomic::ActivationPlan,
102/// {
103///     let _accepted = request;
104///     let _duplicate = request;
105/// }
106/// ```
107#[must_use = "activation work must be admitted or returned whole"]
108pub struct BeginActivation<W, P>
109where
110    W: Behavior,
111    BehaviorAddr<W>: EndpointAddress,
112    P: ActivationPlan,
113{
114    activation: WorkerActivationGrant<W, P>,
115    permit: ActivationPermit<W>,
116    plan: P,
117}
118
119impl<W, P> BeginActivation<W, P>
120where
121    W: Behavior,
122    BehaviorAddr<W>: EndpointAddress,
123    P: ActivationPlan,
124{
125    /// Consume the successful initialization authority and concrete plan.
126    #[must_use]
127    pub fn new(plan: P, permit: ActivationPermit<W>) -> Self {
128        Self {
129            activation: WorkerActivationGrant::issued(permit.worker_evidence()),
130            permit,
131            plan,
132        }
133    }
134
135    /// Transfer original permit/plan and its exact typed grant independently.
136    #[must_use]
137    pub fn into_parts(self) -> (ActivationPermit<W>, P, WorkerActivationGrant<W, P>) {
138        let Self {
139            activation,
140            permit,
141            plan,
142        } = self;
143        (permit, plan, activation)
144    }
145
146    /// Exact installed worker target carried by the permit.
147    #[must_use]
148    pub fn target(&self) -> EstablishedRecipient<W::Protocol> {
149        self.permit.target()
150    }
151
152    /// Worker creation correlation authorized by the consumed permit.
153    #[must_use]
154    pub fn worker(&self) -> WorkerAttempt {
155        self.permit.worker()
156    }
157
158    /// Initialization correlation authorized by the consumed permit.
159    #[must_use]
160    pub fn initialization(&self) -> InitializationAttempt {
161        self.permit.initialization()
162    }
163
164    pub(in super::super) fn attempt(&self) -> WorkerActivationGrant<W, P> {
165        self.activation.clone()
166    }
167
168    /// Produce the input Bombay must admit before polling the plan.
169    #[must_use]
170    pub fn started(&self) -> WorkerActivation<W, P> {
171        WorkerActivation {
172            worker: self.worker(),
173            activation: self.activation.clone(),
174            outcome: WorkerActivationOutcome::Started,
175            worker_type: PhantomData,
176        }
177    }
178
179    /// Return a request that its owner could not admit.
180    #[must_use]
181    pub fn start_rejected(self, reason: ActivationStartRejection) -> WorkerActivation<W, P> {
182        let worker = self.worker();
183        let activation = self.activation.clone();
184        WorkerActivation {
185            worker,
186            activation,
187            outcome: WorkerActivationOutcome::StartRejected {
188                request: self,
189                reason,
190            },
191            worker_type: PhantomData,
192        }
193    }
194
195    /// Consume the plan and produce its exact ready or rejected input.
196    pub async fn activate(self) -> WorkerActivation<W, P> {
197        let Self {
198            activation,
199            permit,
200            plan,
201        } = self;
202        let worker = permit.worker();
203        match plan.activate().await {
204            Ok(readiness) => WorkerActivation {
205                worker,
206                activation,
207                outcome: WorkerActivationOutcome::Ready(readiness),
208                worker_type: PhantomData,
209            },
210            Err(rejection) => WorkerActivation {
211                worker,
212                activation,
213                outcome: WorkerActivationOutcome::Rejected(rejection),
214                worker_type: PhantomData,
215            },
216        }
217    }
218}
219
220impl<W, P> InterpreterRequest for BeginActivation<W, P>
221where
222    W: Behavior,
223    BehaviorAddr<W>: EndpointAddress,
224    P: ActivationPlan,
225{
226    type ReturnToEmitter = ReturnsToEmitter<WorkerActivation<W, P>, Here>;
227    type LogicalProtocols = behavior::NoBirthProtocols;
228}
229
230impl<W, P> ActionItem for BeginActivation<W, P>
231where
232    W: Behavior,
233    BehaviorAddr<W>: EndpointAddress,
234    EstablishedRecipient<W::Protocol>: Send,
235    P: ActivationPlan,
236{
237    type Custody = (Option<Self>, Option<Self::Reply>);
238    type Input<'a>
239        = &'a mut Option<Self>
240    where
241        Self: 'a;
242    type Reply = ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>;
243
244    fn prepare_interpretation(
245        progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
246    ) {
247        prepare_item::<Self>(progress);
248    }
249
250    fn interpretation_input<'a>(
251        custody: &'a mut Self::Custody,
252    ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
253    where
254        Self: 'a,
255    {
256        match custody {
257            (input @ Some(_), received @ None) => Some((input, received)),
258            _ => None,
259        }
260    }
261
262    fn finish_interpretation(
263        progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
264    ) {
265        finish_item::<Self>(progress);
266    }
267
268    type Accepted = ();
269    type Rejection = ActivationStartRejection;
270    type Prerequisite = Never;
271}
272
273pub(in crate::atomic) enum WorkerActivationOutcome<W, P>
274where
275    W: Behavior,
276    BehaviorAddr<W>: EndpointAddress,
277    P: ActivationPlan,
278{
279    Started,
280    StartRejected {
281        request: BeginActivation<W, P>,
282        reason: ActivationStartRejection,
283    },
284    Ready(P::Ready),
285    Rejected(P::Rejection),
286}
287
288/// Exact input produced by one accepted worker activation request.
289#[must_use = "worker activation must be admitted or transferred outward"]
290pub struct WorkerActivation<W, P>
291where
292    W: Behavior,
293    BehaviorAddr<W>: EndpointAddress,
294    P: ActivationPlan,
295{
296    worker: WorkerAttempt,
297    activation: WorkerActivationGrant<W, P>,
298    outcome: WorkerActivationOutcome<W, P>,
299    worker_type: PhantomData<fn() -> W>,
300}
301
302impl<W, P> WorkerActivation<W, P>
303where
304    W: Behavior,
305    BehaviorAddr<W>: EndpointAddress,
306    P: ActivationPlan,
307{
308    /// Worker creation correlation carried by this input.
309    #[must_use]
310    pub fn worker(&self) -> WorkerAttempt {
311        self.worker.clone()
312    }
313
314    pub(in super::super) const fn attempt(&self) -> &WorkerActivationGrant<W, P> {
315        &self.activation
316    }
317
318    pub(in crate::atomic) fn into_parts(
319        self,
320    ) -> (
321        WorkerAttempt,
322        WorkerActivationGrant<W, P>,
323        WorkerActivationOutcome<W, P>,
324    ) {
325        (self.worker, self.activation, self.outcome)
326    }
327
328    pub(in crate::atomic) fn into_started(self) -> Result<(), Self> {
329        let Self {
330            worker,
331            activation,
332            outcome,
333            worker_type,
334        } = self;
335        match outcome {
336            WorkerActivationOutcome::Started => Ok(()),
337            outcome => Err(Self {
338                worker,
339                activation,
340                outcome,
341                worker_type,
342            }),
343        }
344    }
345
346    pub(in super::super) fn into_start_rejection(
347        self,
348    ) -> Result<(BeginActivation<W, P>, ActivationStartRejection), Self> {
349        let Self {
350            worker,
351            activation,
352            outcome,
353            worker_type,
354        } = self;
355        match outcome {
356            WorkerActivationOutcome::StartRejected { request, reason } => Ok((request, reason)),
357            outcome => Err(Self {
358                worker,
359                activation,
360                outcome,
361                worker_type,
362            }),
363        }
364    }
365
366    /// Consume readiness only when this is the ready alternative.
367    pub fn into_ready(self) -> Result<P::Ready, Self> {
368        let Self {
369            worker,
370            activation,
371            outcome,
372            worker_type,
373        } = self;
374        match outcome {
375            WorkerActivationOutcome::Ready(readiness) => Ok(readiness),
376            outcome => Err(Self {
377                worker,
378                activation,
379                outcome,
380                worker_type,
381            }),
382        }
383    }
384
385    /// Consume the application rejection only when this is the rejected alternative.
386    pub fn into_rejection(self) -> Result<P::Rejection, Self> {
387        let Self {
388            worker,
389            activation,
390            outcome,
391            worker_type,
392        } = self;
393        match outcome {
394            WorkerActivationOutcome::Rejected(rejection) => Ok(rejection),
395            outcome => Err(Self {
396                worker,
397                activation,
398                outcome,
399                worker_type,
400            }),
401        }
402    }
403}