1use 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
18pub 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 #[must_use]
32 pub const fn start(worker: W) -> Self {
33 Self {
34 command: ProxyCommand::Start(WorkerSubmission::immediate(worker)),
35 }
36 }
37
38 #[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 #[must_use]
50 pub const fn shutdown() -> Self {
51 Self {
52 command: ProxyCommand::Shutdown,
53 }
54 }
55
56 #[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 #[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#[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
91pub 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 Overlap {
101 worker: W,
102 activation: P,
103 phase: ProxyPhase,
104 },
105 WorkerAttemptsExhausted { worker: W, activation: P },
107 Resolved { result: WorkerStartResult<W, P> },
109}
110
111pub 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 NotReplaceable {
121 worker: W,
122 activation: P,
123 phase: ProxyPhase,
124 },
125 WorkerAttemptsExhausted {
127 replaces: WorkerAttempt,
128 worker: W,
129 activation: P,
130 },
131 CancelledBeforeBirth {
133 replaces: WorkerAttempt,
134 worker: W,
135 activation: P,
136 },
137 Resolved {
139 replaces: WorkerAttempt,
140 result: WorkerStartResult<W, P>,
141 },
142}
143
144pub enum ProxyDiagnostic<W, P>
146where
147 W: Behavior,
148 P: ActivationPlan,
149 BehaviorAddr<W>: EndpointAddress,
150{
151 UnexpectedWorkerStart {
153 phase: ProxyPhase,
154 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
155 },
156 UnexpectedWorkerStop {
158 phase: ProxyPhase,
159 stopped: ChildStopped<BehaviorAddr<W>>,
160 },
161 UnexpectedWorkerInitialization {
163 phase: ProxyPhase,
164 initialization: WorkerInitializationReport<W, P>,
165 },
166 UnexpectedWorkerActivation {
168 phase: ProxyPhase,
169 activation: WorkerActivation<W, P>,
170 },
171 UnexpectedWorkerShutdown {
173 phase: ProxyPhase,
174 shutdown: EstablishedShutdownResolved<W::Protocol>,
175 },
176}
177
178pub 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 Initial { outcome: InitialWorkerOutcome<W, P> },
188 Replacement { outcome: ReplacementOutcome<W, P> },
190 WorkerStopped {
192 worker: WorkerAttempt,
193 stopped: ChildStopped<BehaviorAddr<W>>,
194 },
195 Unavailable {
197 sender: BehaviorAddr<W>,
198 phase: ProxyPhase,
199 command: <W::Protocol as Protocol>::Msg,
200 },
201}
202
203#[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 InitializationRejected {
213 failure: WorkerInitializationFailure,
214 activation: P,
215 shutdown: Option<EstablishedShutdownResolved<W::Protocol>>,
216 stopped: ChildStopped<BehaviorAddr<W>>,
217 },
218 InitializationStopped {
220 activation: P,
221 observed: Option<ChildStopped<BehaviorAddr<W>>>,
222 returned: ChildStopped<BehaviorAddr<W>>,
223 },
224 InitializationCompleted {
226 permit: ActivationPermit<W>,
227 activation: P,
228 stopped: ChildStopped<BehaviorAddr<W>>,
229 },
230 ActivationStartRejected {
232 request: BeginActivation<W, P>,
233 reason: ActivationStartRejection,
234 shutdown: Option<EstablishedShutdownResolved<W::Protocol>>,
235 stopped: ChildStopped<BehaviorAddr<W>>,
236 },
237 ActivationRejected {
239 rejection: P::Rejection,
240 shutdown: Option<EstablishedShutdownResolved<W::Protocol>>,
241 stopped: ChildStopped<BehaviorAddr<W>>,
242 },
243 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}