1use core::mem;
4use core::ops::ControlFlow;
5
6use behavior::{
7 Actions, ActiveTurn, Address, Behavior, BehaviorActed, BehaviorAddr, BehaviorBase, BirthMode,
8 Births, ChildCreationProduct, ChildHead, CreateChild, CreationSequence, Creations,
9 CreationsSettled, EndpointAddress, EstablishedDelivery, EstablishedRecipient, Here,
10 InjectEvent, InterpreterRequests, Never, Protocol, ReportToParent, SendEffects,
11 SendSettlements, User,
12};
13
14use crate::{
15 ChildStopped, EstablishedShutdownResolved, ObserveChild, ShutdownEstablished,
16 ShutdownRequested, StopOnShutdown,
17};
18
19mod effects;
20mod operation;
21mod protocol;
22mod state;
23mod worker;
24
25use super::worker::WorkerActivationGrant;
26pub use super::worker::WorkerAttempt;
27pub use super::worker::{
28 ActivationPermit, InitializationAttempt, InitializeWorker, WorkerInitializationFailure,
29 WorkerInitializationReport,
30};
31pub use super::worker::{
32 ActivationPlan, ActivationStartRejection, BeginActivation, ImmediateActivation,
33 WorkerActivation,
34};
35pub use effects::ProxyEffects;
36pub(crate) use operation::ProxyOperationId;
37pub(crate) use operation::ProxyOperationWitness;
38pub use operation::{ProxyControlAdmission, ProxyInputReceipt, ProxyInputResult, ProxyOperation};
39pub use protocol::{
40 InitialWorkerOutcome, ProxyControl, ProxyDiagnostic, ProxyDrain, ProxyOutcome, ProxyPhase,
41 ReplacementOutcome,
42};
43use protocol::{ProxyCommand, ProxyEvent};
44use state::{
45 ActivationDuringDeparture, ActivationProgress, PreReadyFailure, PredecessorReturn,
46 PredecessorShutdown, ProxyReplacement, ProxyRetirement, ProxyShutdown, ProxyState,
47 ReplacementCompletion, WorkerActivationRetirement, WorkerActivationShutdown,
48 WorkerInitializationRetirement, WorkerInitializationShutdown, WorkerStart, WorkerStartKind,
49 WorkerStartPhase, WorkerStartRetirement,
50};
51pub use worker::WorkerStartResult;
52use worker::{CurrentWorker, PendingWorker, StoppedWorker, WorkerCreation, WorkerStopping};
53
54pub struct StableProxy<W, P>
56where
57 W: Behavior,
58 P: ActivationPlan,
59 BehaviorAddr<W>: EndpointAddress,
60{
61 state: ProxyState<W, P>,
62 creations: CreationSequence,
63}
64
65impl<W> StableProxy<W, ImmediateActivation>
66where
67 W: Behavior,
68 BehaviorAddr<W>: EndpointAddress,
69{
70 #[must_use]
72 pub const fn immediate() -> Self {
73 Self::dormant()
74 }
75}
76
77impl<W, P> StableProxy<W, P>
78where
79 W: Behavior,
80 P: ActivationPlan,
81 BehaviorAddr<W>: EndpointAddress,
82{
83 const fn dormant() -> Self {
84 Self {
85 state: ProxyState::Dormant,
86 creations: CreationSequence::new(),
87 }
88 }
89
90 #[must_use]
92 pub const fn activated() -> Self {
93 Self::dormant()
94 }
95
96 #[must_use]
98 pub const fn phase(&self) -> ProxyPhase {
99 self.state.phase()
100 }
101}
102
103#[cfg(test)]
104mod shutdown_ownership_tests {
105 use std::sync::{Arc, Mutex, mpsc};
106 use std::time::{Duration, Instant};
107
108 use core::ops::ControlFlow;
109
110 use behavior::{
111 Actions, ActiveTurn, Address, Behavior, BehaviorActed, CreationSequence, EndpointAddress,
112 EstablishedActor, Never, NoBirths, Protocol, User,
113 };
114
115 use crate::{ChildStopped, Exit, StopOnShutdown};
116
117 use super::{
118 CurrentWorker, ImmediateActivation, InitializationAttempt, StableProxy, WorkerAttempt,
119 WorkerInitializationRetirement, WorkerInitializationShutdown, WorkerStopping,
120 };
121
122 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
123 struct SearchAddress;
124
125 impl Address for SearchAddress {
126 type Nonce = u64;
127 }
128
129 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
130 struct SearchEndpoint;
131
132 struct SearchInstalled<B: Behavior> {
133 endpoint: SearchEndpoint,
134 control: mpsc::Sender<B::Event>,
135 inbox: Arc<Mutex<mpsc::Receiver<B::Event>>>,
136 }
137
138 impl<B: Behavior> Clone for SearchInstalled<B> {
139 fn clone(&self) -> Self {
140 Self {
141 endpoint: self.endpoint,
142 control: self.control.clone(),
143 inbox: Arc::clone(&self.inbox),
144 }
145 }
146 }
147
148 impl<B: Behavior> SearchInstalled<B> {
149 fn new(endpoint: SearchEndpoint) -> Self {
150 let (control, inbox) = mpsc::channel();
151 Self {
152 endpoint,
153 control,
154 inbox: Arc::new(Mutex::new(inbox)),
155 }
156 }
157 }
158
159 impl EndpointAddress for SearchAddress {
160 type Established<P>
161 = SearchEndpoint
162 where
163 P: Protocol<Addr = Self>;
164
165 type Installed<B>
166 = SearchInstalled<B>
167 where
168 B: Behavior<Protocol: Protocol<Addr = Self>>;
169
170 fn recipient<B>(installed: &Self::Installed<B>) -> SearchEndpoint
171 where
172 B: Behavior<Protocol: Protocol<Addr = Self>>,
173 {
174 installed.endpoint
175 }
176 }
177
178 struct SearchWorker;
179
180 impl Protocol for SearchWorker {
181 type Addr = SearchAddress;
182 type Msg = ();
183 }
184
185 impl Behavior for SearchWorker {
186 type Protocol = Self;
187 type Event = User<SearchAddress, ()>;
188 type Sends = Vec<Never>;
189 type Ph = Never;
190 type Error = Never;
191 type Birth = NoBirths;
192
193 fn transition(&mut self, _: ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
194 Ok(Actions::cont())
195 }
196 }
197
198 #[test]
199 fn initialization_and_proxy_stops_remain_separately_owned_during_shutdown() {
200 let mut creations = CreationSequence::new();
201 let creation = creations
202 .issue()
203 .unwrap_or_else(|| panic!("the first worker creation is available"));
204 let worker = WorkerAttempt::issued(creation);
205 let initialization = InitializationAttempt::issued(&worker);
206 let current = CurrentWorker {
207 attempt: worker.clone(),
208 initialization: initialization.clone(),
209 actor: EstablishedActor::<StopOnShutdown<SearchWorker>>::issued(SearchInstalled::new(
210 SearchEndpoint,
211 )),
212 };
213 let (departing, _shutdown) = WorkerStopping::begin(current);
214 let proxy_stop = ChildStopped::new(creation, Ok(Exit::Normal), Instant::now());
215 let initialization_stop = ChildStopped::new(
216 creation,
217 Ok(Exit::Normal),
218 proxy_stop.at + Duration::from_nanos(1),
219 );
220 let departing = match departing.worker_stopped(proxy_stop) {
221 Ok(ControlFlow::Continue(departing)) => departing,
222 Ok(ControlFlow::Break(_)) => panic!("shutdown has not settled"),
223 Err(_) => panic!("the proxy stop belongs to the current worker"),
224 };
225
226 let retained = StableProxy::<SearchWorker, ImmediateActivation>::retain_initialization_stop(
227 departing,
228 worker,
229 initialization,
230 ImmediateActivation,
231 initialization_stop,
232 );
233 let shutdown = match retained {
234 Ok(ControlFlow::Continue(shutdown)) => shutdown,
235 Ok(ControlFlow::Break(_)) => panic!("shutdown still awaits its settlement"),
236 Err(_) => panic!("the initialization stop belongs to the current worker"),
237 };
238 let WorkerInitializationShutdown::Departing {
239 initialization:
240 Some(WorkerInitializationRetirement::Stopped {
241 initialization_stop: Some(retained_stop),
242 ..
243 }),
244 departure,
245 } = shutdown
246 else {
247 panic!("both stop values remain in the current shutdown state");
248 };
249
250 assert_eq!(departure.stopped(), Some(&proxy_stop));
251 assert_eq!(retained_stop, initialization_stop);
252 }
253
254 #[test]
255 fn initialization_stop_occupies_the_worker_stop_slot_once() {
256 let mut creations = CreationSequence::new();
257 let creation = creations
258 .issue()
259 .unwrap_or_else(|| panic!("the first worker creation is available"));
260 let worker = WorkerAttempt::issued(creation);
261 let initialization = InitializationAttempt::issued(&worker);
262 let current = CurrentWorker {
263 attempt: worker.clone(),
264 initialization: initialization.clone(),
265 actor: EstablishedActor::<StopOnShutdown<SearchWorker>>::issued(SearchInstalled::new(
266 SearchEndpoint,
267 )),
268 };
269 let (departing, _shutdown) = WorkerStopping::begin(current);
270 let initialization_stop = ChildStopped::new(creation, Ok(Exit::Normal), Instant::now());
271
272 let retained = StableProxy::<SearchWorker, ImmediateActivation>::retain_initialization_stop(
273 departing,
274 worker,
275 initialization,
276 ImmediateActivation,
277 initialization_stop,
278 );
279 let shutdown = match retained {
280 Ok(ControlFlow::Continue(shutdown)) => shutdown,
281 Ok(ControlFlow::Break(_)) => panic!("shutdown still awaits its settlement"),
282 Err(_) => panic!("the initialization stop belongs to the current worker"),
283 };
284 let WorkerInitializationShutdown::Departing {
285 initialization:
286 Some(WorkerInitializationRetirement::Stopped {
287 initialization_stop: None,
288 ..
289 }),
290 departure,
291 } = shutdown
292 else {
293 panic!("the initialization stop occupies one ownership slot");
294 };
295
296 assert_eq!(departure.stopped(), Some(&initialization_stop));
297 }
298}
299
300impl<W, P> BehaviorBase for StableProxy<W, P>
301where
302 W: Behavior,
303 P: ActivationPlan,
304 BehaviorAddr<W>: EndpointAddress,
305{
306 type Base = Self;
307
308 fn base(&self) -> &Self::Base {
309 self
310 }
311}
312
313impl<W, P> StableProxy<W, P>
314where
315 W: Behavior,
316 P: ActivationPlan,
317 BehaviorAddr<W>: EndpointAddress,
318 <BehaviorAddr<W> as Address>::Nonce: Copy + Eq,
319 StopOnShutdown<W>:
320 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
321 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
322 <StopOnShutdown<W> as Behavior>::Sends: SendSettlements,
323 <W::Birth as BirthMode>::Child: ChildCreationProduct<BehaviorAddr<W>, ChildHead>,
324 EstablishedRecipient<W::Protocol>: Send,
325{
326 fn accept(
327 current: StableProxy<W, P>,
328 event: ProxyEvent<W, P>,
329 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
330 match event {
331 ProxyEvent::Service(service) => Self::service(current, service),
332 ProxyEvent::Owner(control) => Self::owner(current, control),
333 ProxyEvent::WorkerCreationsSettled(workers) => Self::workers_created(current, workers),
334 ProxyEvent::WorkerInitialization(initialization) => {
335 Self::worker_initialized(current, initialization)
336 }
337 ProxyEvent::WorkerActivationReported(activation) => {
338 Self::worker_activated(current, activation)
339 }
340 ProxyEvent::WorkerShutdownResolved(shutdown) => {
341 Self::worker_shutdown_resolved(current, shutdown)
342 }
343 ProxyEvent::WorkerStopped(stopped) => Self::worker_stopped(current, stopped),
344 }
345 }
346
347 fn admit_activation(
348 progress: &ActivationProgress<W, P>,
349 input: WorkerActivation<W, P>,
350 ) -> Result<WorkerActivation<W, P>, WorkerActivation<W, P>> {
351 let expected = match progress {
352 ActivationProgress::WaitingForStart(attempt) | ActivationProgress::Running(attempt) => {
353 attempt
354 }
355 };
356 if input.attempt() == expected {
357 Ok(input)
358 } else {
359 Err(input)
360 }
361 }
362
363 fn service(
364 current: StableProxy<W, P>,
365 service: User<BehaviorAddr<W>, <W::Protocol as Protocol>::Msg>,
366 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
367 match current.state {
368 ProxyState::Ready { worker } => {
369 let delivery = EstablishedDelivery::new(worker.actor.recipient(), service.message);
370 (
371 StableProxy {
372 state: ProxyState::Ready { worker },
373 creations: current.creations,
374 },
375 Actions::send(ProxyEffects {
376 worker_observations: InterpreterRequests::empty(),
377 worker_initializations: InterpreterRequests::empty(),
378 worker_activations: InterpreterRequests::empty(),
379 worker_shutdowns: InterpreterRequests::empty(),
380 worker_deliveries: vec![delivery],
381 owner_outcomes: InterpreterRequests::empty(),
382 diagnostics: InterpreterRequests::empty(),
383 }),
384 )
385 }
386 state => {
387 let outcome = ProxyOutcome::Unavailable {
388 sender: service.from,
389 phase: state.phase(),
390 command: service.message,
391 };
392 (
393 StableProxy {
394 state,
395 creations: current.creations,
396 },
397 Self::report(outcome),
398 )
399 }
400 }
401 }
402
403 fn owner(
404 current: StableProxy<W, P>,
405 control: ProxyControl<W, P>,
406 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
407 match control.command {
408 ProxyCommand::Start(submission) => {
409 Self::start_initial(current, submission.worker, submission.activation)
410 }
411 ProxyCommand::Replace(submission) => {
412 Self::replace_worker(current, submission.worker, submission.activation)
413 }
414 ProxyCommand::Shutdown => Self::owner_shutdown(current),
415 }
416 }
417
418 fn start_initial(
419 current: StableProxy<W, P>,
420 worker: W,
421 activation: P,
422 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
423 match current.state {
424 ProxyState::Dormant => Self::begin_initial(current.creations, worker, activation),
425 state => {
426 let phase = state.phase();
427 let outcome = ProxyOutcome::Initial {
428 outcome: InitialWorkerOutcome::Overlap {
429 worker,
430 activation,
431 phase,
432 },
433 };
434 (
435 StableProxy {
436 state,
437 creations: current.creations,
438 },
439 Self::report(outcome),
440 )
441 }
442 }
443 }
444
445 fn begin_initial(
446 mut creations: CreationSequence,
447 worker: W,
448 activation: P,
449 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
450 match PendingWorker::birth(&mut creations, worker, activation) {
451 Ok((pending_worker, creation)) => {
452 let observation = ObserveChild::new(pending_worker.creation());
453 (
454 StableProxy {
455 state: ProxyState::Starting(WorkerStart {
456 kind: WorkerStartKind::Initial,
457 phase: WorkerStartPhase::Creating {
458 worker: pending_worker,
459 stopped: None,
460 },
461 }),
462 creations,
463 },
464 Actions::new(
465 ProxyEffects {
466 worker_observations: InterpreterRequests::one(observation),
467 worker_initializations: InterpreterRequests::empty(),
468 worker_activations: InterpreterRequests::empty(),
469 worker_shutdowns: InterpreterRequests::empty(),
470 worker_deliveries: Vec::new(),
471 owner_outcomes: InterpreterRequests::empty(),
472 diagnostics: InterpreterRequests::empty(),
473 },
474 Creations::one(creation),
475 behavior::Step::Continue,
476 ),
477 )
478 }
479 Err((worker, activation)) => (
480 StableProxy {
481 state: ProxyState::Dormant,
482 creations,
483 },
484 Self::report(ProxyOutcome::Initial {
485 outcome: InitialWorkerOutcome::WorkerAttemptsExhausted { worker, activation },
486 }),
487 ),
488 }
489 }
490
491 fn workers_created(
492 current: StableProxy<W, P>,
493 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
494 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
495 match current.state {
496 ProxyState::Starting(WorkerStart {
497 kind,
498 phase: WorkerStartPhase::Creating { worker, stopped },
499 }) => match worker.created(workers, stopped) {
500 WorkerCreation::Initializing {
501 worker,
502 activation,
503 stopped,
504 } => {
505 Self::await_initialization(current.creations, kind, worker, activation, stopped)
506 }
507 WorkerCreation::Rejected {
508 rejection,
509 activation,
510 stopped,
511 } => Self::worker_creation_rejected(
512 current.creations,
513 kind,
514 WorkerStartResult::CreationRejected {
515 rejection,
516 activation,
517 stopped,
518 },
519 ),
520 WorkerCreation::Unexpected {
521 worker,
522 stopped,
523 workers,
524 } => {
525 let phase = ProxyPhase::Creating;
526 (
527 StableProxy {
528 state: ProxyState::Starting(WorkerStart {
529 kind,
530 phase: WorkerStartPhase::Creating { worker, stopped },
531 }),
532 creations: current.creations,
533 },
534 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerStart { phase, workers }),
535 )
536 }
537 },
538 ProxyState::ShuttingDown(shutdown) => {
539 Self::shutdown_workers_created(current.creations, shutdown, workers)
540 }
541 state => {
542 let phase = state.phase();
543 let diagnostic = ProxyDiagnostic::UnexpectedWorkerStart { phase, workers };
544 (
545 StableProxy {
546 state,
547 creations: current.creations,
548 },
549 Self::diagnose(diagnostic),
550 )
551 }
552 }
553 }
554
555 fn worker_stopped(
556 current: StableProxy<W, P>,
557 stopped: ChildStopped<BehaviorAddr<W>>,
558 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
559 match current.state {
560 ProxyState::Starting(start) => {
561 Self::worker_start_stopped(current.creations, start, stopped)
562 }
563 ProxyState::Replacing(replacement) => {
564 Self::replacement_worker_stopped(current.creations, replacement, stopped)
565 }
566 ProxyState::ShuttingDown(shutdown) => {
567 Self::shutdown_worker_stopped(current.creations, shutdown, stopped)
568 }
569 ProxyState::Ready { worker } => match worker.admit_stop(stopped) {
570 Ok(stopped) => (
571 StableProxy {
572 state: ProxyState::EmptyAfter {
573 previous: worker.attempt.clone(),
574 },
575 creations: current.creations,
576 },
577 Self::report(ProxyOutcome::WorkerStopped {
578 worker: worker.attempt,
579 stopped,
580 }),
581 ),
582 Err(stopped) => {
583 Self::unexpected_stop(current.creations, ProxyState::Ready { worker }, stopped)
584 }
585 },
586 state => {
587 let phase = state.phase();
588 (
589 StableProxy {
590 state,
591 creations: current.creations,
592 },
593 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerStop { phase, stopped }),
594 )
595 }
596 }
597 }
598
599 fn worker_start_stopped(
600 creations: CreationSequence,
601 start: WorkerStart<W, P>,
602 stopped: ChildStopped<BehaviorAddr<W>>,
603 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
604 let WorkerStart { kind, phase } = start;
605 match phase {
606 WorkerStartPhase::ReturningWorker { departure, failure } => {
607 match departure.worker_stopped(stopped) {
608 Ok(ControlFlow::Continue(departure)) => (
609 StableProxy {
610 state: ProxyState::Starting(WorkerStart {
611 kind,
612 phase: WorkerStartPhase::ReturningWorker { departure, failure },
613 }),
614 creations,
615 },
616 Actions::cont(),
617 ),
618 Ok(ControlFlow::Break(StoppedWorker {
619 worker,
620 shutdown,
621 stopped,
622 })) => {
623 Self::worker_returned(creations, kind, worker, failure, shutdown, stopped)
624 }
625 Err((departure, input)) => Self::unexpected_stop(
626 creations,
627 ProxyState::Starting(WorkerStart {
628 kind,
629 phase: WorkerStartPhase::ReturningWorker { departure, failure },
630 }),
631 input,
632 ),
633 }
634 }
635 phase => match Self::admit_pre_ready_stop(phase, stopped) {
636 Ok(phase) => (
637 StableProxy {
638 state: ProxyState::Starting(WorkerStart { kind, phase }),
639 creations,
640 },
641 Actions::cont(),
642 ),
643 Err((phase, stopped)) => Self::unexpected_stop(
644 creations,
645 ProxyState::Starting(WorkerStart { kind, phase }),
646 stopped,
647 ),
648 },
649 }
650 }
651
652 fn admit_pre_ready_stop(
653 phase: WorkerStartPhase<W, P>,
654 stopped: ChildStopped<BehaviorAddr<W>>,
655 ) -> Result<WorkerStartPhase<W, P>, (WorkerStartPhase<W, P>, ChildStopped<BehaviorAddr<W>>)>
656 {
657 match phase {
658 WorkerStartPhase::Creating {
659 worker,
660 stopped: None,
661 } => match worker.admit_stop(stopped) {
662 Ok(stopped) => Ok(WorkerStartPhase::Creating {
663 worker,
664 stopped: Some(stopped),
665 }),
666 Err(stopped) => Err((
667 WorkerStartPhase::Creating {
668 worker,
669 stopped: None,
670 },
671 stopped,
672 )),
673 },
674 WorkerStartPhase::Initializing {
675 worker,
676 stopped: None,
677 } => match worker.admit_stop(stopped) {
678 Ok(stopped) => Ok(WorkerStartPhase::Initializing {
679 worker,
680 stopped: Some(stopped),
681 }),
682 Err(stopped) => Err((
683 WorkerStartPhase::Initializing {
684 worker,
685 stopped: None,
686 },
687 stopped,
688 )),
689 },
690 WorkerStartPhase::Activating {
691 worker,
692 progress,
693 stopped: None,
694 } => match worker.admit_stop(stopped) {
695 Ok(stopped) => Ok(WorkerStartPhase::Activating {
696 worker,
697 progress,
698 stopped: Some(stopped),
699 }),
700 Err(stopped) => Err((
701 WorkerStartPhase::Activating {
702 worker,
703 progress,
704 stopped: None,
705 },
706 stopped,
707 )),
708 },
709 phase => Err((phase, stopped)),
710 }
711 }
712
713 fn unexpected_stop(
714 creations: CreationSequence,
715 state: ProxyState<W, P>,
716 stopped: ChildStopped<BehaviorAddr<W>>,
717 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
718 let phase = state.phase();
719 (
720 StableProxy { state, creations },
721 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerStop { phase, stopped }),
722 )
723 }
724
725 fn await_initialization(
726 creations: CreationSequence,
727 kind: WorkerStartKind,
728 worker: CurrentWorker<W>,
729 activation: P,
730 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
731 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
732 let request = InitializeWorker::<W, P>::new(
733 worker.attempt.clone(),
734 worker.initialization.clone(),
735 worker.actor.recipient(),
736 activation,
737 );
738 (
739 StableProxy {
740 state: ProxyState::Starting(WorkerStart {
741 kind,
742 phase: WorkerStartPhase::Initializing { worker, stopped },
743 }),
744 creations,
745 },
746 Actions::send(ProxyEffects {
747 worker_observations: InterpreterRequests::empty(),
748 worker_initializations: InterpreterRequests::one(request),
749 worker_activations: InterpreterRequests::empty(),
750 worker_shutdowns: InterpreterRequests::empty(),
751 worker_deliveries: Vec::new(),
752 owner_outcomes: InterpreterRequests::empty(),
753 diagnostics: InterpreterRequests::empty(),
754 }),
755 )
756 }
757
758 fn worker_creation_rejected(
759 creations: CreationSequence,
760 kind: WorkerStartKind,
761 result: WorkerStartResult<W, P>,
762 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
763 match kind {
764 WorkerStartKind::Initial => (
765 StableProxy {
766 state: ProxyState::EmptyInitial,
767 creations,
768 },
769 Self::report(ProxyOutcome::Initial {
770 outcome: InitialWorkerOutcome::Resolved { result },
771 }),
772 ),
773 WorkerStartKind::Replacement {
774 replaces,
775 predecessor_shutdown,
776 } => Self::finish_replacement(
777 creations,
778 replaces.clone(),
779 predecessor_shutdown,
780 ReplacementCompletion::Empty {
781 worker: replaces,
782 result,
783 },
784 ),
785 }
786 }
787
788 fn worker_start_ready(
789 creations: CreationSequence,
790 kind: WorkerStartKind,
791 worker: CurrentWorker<W>,
792 result: WorkerStartResult<W, P>,
793 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
794 match kind {
795 WorkerStartKind::Initial => (
796 StableProxy {
797 state: ProxyState::Ready { worker },
798 creations,
799 },
800 Self::report(ProxyOutcome::Initial {
801 outcome: InitialWorkerOutcome::Resolved { result },
802 }),
803 ),
804 WorkerStartKind::Replacement {
805 replaces,
806 predecessor_shutdown,
807 } => Self::finish_replacement(
808 creations,
809 replaces,
810 predecessor_shutdown,
811 ReplacementCompletion::Ready { worker, result },
812 ),
813 }
814 }
815
816 fn worker_start_empty(
817 creations: CreationSequence,
818 kind: WorkerStartKind,
819 previous: WorkerAttempt,
820 result: WorkerStartResult<W, P>,
821 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
822 match kind {
823 WorkerStartKind::Initial => (
824 StableProxy {
825 state: ProxyState::EmptyAfter { previous },
826 creations,
827 },
828 Self::report(ProxyOutcome::Initial {
829 outcome: InitialWorkerOutcome::Resolved { result },
830 }),
831 ),
832 WorkerStartKind::Replacement {
833 replaces,
834 predecessor_shutdown,
835 } => Self::finish_replacement(
836 creations,
837 replaces,
838 predecessor_shutdown,
839 ReplacementCompletion::Empty {
840 worker: previous,
841 result,
842 },
843 ),
844 }
845 }
846
847 fn create_worker(
848 observation: ObserveChild<W::Protocol, ChildHead>,
849 creation: CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
850 ) -> ProxyActions<W, P> {
851 Actions::new(
852 ProxyEffects {
853 worker_observations: InterpreterRequests::one(observation),
854 worker_initializations: InterpreterRequests::empty(),
855 worker_activations: InterpreterRequests::empty(),
856 worker_shutdowns: InterpreterRequests::empty(),
857 worker_deliveries: Vec::new(),
858 owner_outcomes: InterpreterRequests::empty(),
859 diagnostics: InterpreterRequests::empty(),
860 },
861 Creations::one(creation),
862 behavior::Step::Continue,
863 )
864 }
865
866 fn report(outcome: ProxyOutcome<W, P>) -> ProxyActions<W, P> {
867 Actions::send(ProxyEffects {
868 worker_observations: InterpreterRequests::empty(),
869 worker_initializations: InterpreterRequests::empty(),
870 worker_activations: InterpreterRequests::empty(),
871 worker_shutdowns: InterpreterRequests::empty(),
872 worker_deliveries: Vec::new(),
873 owner_outcomes: InterpreterRequests::one(ReportToParent::new(outcome)),
874 diagnostics: InterpreterRequests::empty(),
875 })
876 }
877
878 fn diagnose(diagnostic: ProxyDiagnostic<W, P>) -> ProxyActions<W, P> {
879 Actions::send(ProxyEffects {
880 worker_observations: InterpreterRequests::empty(),
881 worker_initializations: InterpreterRequests::empty(),
882 worker_activations: InterpreterRequests::empty(),
883 worker_shutdowns: InterpreterRequests::empty(),
884 worker_deliveries: Vec::new(),
885 owner_outcomes: InterpreterRequests::empty(),
886 diagnostics: InterpreterRequests::one(ReportToParent::new(diagnostic)),
887 })
888 }
889}
890
891impl<W, P> Behavior for StableProxy<W, P>
892where
893 W: Behavior,
894 P: ActivationPlan,
895 BehaviorAddr<W>: EndpointAddress,
896 <BehaviorAddr<W> as Address>::Nonce: Copy + Eq + Send,
897 StopOnShutdown<W>:
898 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
899 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
900 <StopOnShutdown<W> as Behavior>::Sends: SendSettlements,
901 <W::Birth as BirthMode>::Child: ChildCreationProduct<BehaviorAddr<W>, ChildHead>,
902 EstablishedRecipient<W::Protocol>: Send,
903{
904 type Protocol = W::Protocol;
905 type Event = ProxyEvent<W, P>;
906 type Sends = ProxyEffects<
907 InterpreterRequests<ObserveChild<W::Protocol, ChildHead>>,
908 InterpreterRequests<InitializeWorker<W, P>>,
909 InterpreterRequests<BeginActivation<W, P>>,
910 InterpreterRequests<ShutdownEstablished<StopOnShutdown<W>, Here>>,
911 Vec<EstablishedDelivery<W::Protocol>>,
912 InterpreterRequests<ReportToParent<ProxyOutcome<W, P>>>,
913 InterpreterRequests<ReportToParent<ProxyDiagnostic<W, P>>>,
914 >;
915 type Ph = Never;
916 type Error = Never;
917 type Birth = Births<StopOnShutdown<W>>;
918
919 fn transition(&mut self, _: ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
920 let current = mem::replace(self, Self::dormant());
921 let (next, actions) = Self::accept(current, event);
922 *self = next;
923 Ok(actions)
924 }
925}
926
927type ProxyActions<W, P> = Actions<
928 BehaviorAddr<W>,
929 Never,
930 ProxyEffects<
931 InterpreterRequests<ObserveChild<<W as Behavior>::Protocol, ChildHead>>,
932 InterpreterRequests<InitializeWorker<W, P>>,
933 InterpreterRequests<BeginActivation<W, P>>,
934 InterpreterRequests<ShutdownEstablished<StopOnShutdown<W>, Here>>,
935 Vec<EstablishedDelivery<<W as Behavior>::Protocol>>,
936 InterpreterRequests<ReportToParent<ProxyOutcome<W, P>>>,
937 InterpreterRequests<ReportToParent<ProxyDiagnostic<W, P>>>,
938 >,
939 Births<StopOnShutdown<W>>,
940>;
941
942impl<W, P> StableProxy<W, P>
944where
945 W: Behavior,
946 P: ActivationPlan,
947 BehaviorAddr<W>: EndpointAddress,
948 <BehaviorAddr<W> as Address>::Nonce: Copy + Eq,
949 StopOnShutdown<W>:
950 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
951 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
952 <StopOnShutdown<W> as Behavior>::Sends: SendSettlements,
953 <W::Birth as BirthMode>::Child: ChildCreationProduct<BehaviorAddr<W>, ChildHead>,
954 EstablishedRecipient<W::Protocol>: Send,
955{
956 fn replace_worker(
957 current: StableProxy<W, P>,
958 worker: W,
959 activation: P,
960 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
961 match current.state {
962 ProxyState::Ready {
963 worker: current_worker,
964 } => {
965 let replaces = current_worker.attempt.clone();
966 let mut creations = current.creations;
967 match PendingWorker::replacement(
968 &mut creations,
969 replaces.creation(),
970 worker,
971 activation,
972 ) {
973 Ok((successor, creation)) => {
974 let (departure, shutdown) = WorkerStopping::begin(current_worker);
975 (
976 StableProxy {
977 state: ProxyState::Replacing(
978 ProxyReplacement::ReturningPredecessor {
979 departure,
980 successor,
981 creation,
982 },
983 ),
984 creations,
985 },
986 Self::request_worker_shutdown(shutdown),
987 )
988 }
989 Err((worker, activation)) => (
990 StableProxy {
991 state: ProxyState::Ready {
992 worker: current_worker,
993 },
994 creations,
995 },
996 Self::report(ProxyOutcome::Replacement {
997 outcome: ReplacementOutcome::WorkerAttemptsExhausted {
998 replaces,
999 worker,
1000 activation,
1001 },
1002 }),
1003 ),
1004 }
1005 }
1006 ProxyState::EmptyAfter { previous } => {
1007 let mut creations = current.creations;
1008 match PendingWorker::replacement(
1009 &mut creations,
1010 previous.creation(),
1011 worker,
1012 activation,
1013 ) {
1014 Ok((pending_worker, creation)) => {
1015 let observation = ObserveChild::new(pending_worker.creation());
1016 (
1017 StableProxy {
1018 state: ProxyState::Starting(WorkerStart {
1019 kind: WorkerStartKind::Replacement {
1020 replaces: previous,
1021 predecessor_shutdown: PredecessorShutdown::Settled,
1022 },
1023 phase: WorkerStartPhase::Creating {
1024 worker: pending_worker,
1025 stopped: None,
1026 },
1027 }),
1028 creations,
1029 },
1030 Self::create_worker(observation, creation),
1031 )
1032 }
1033 Err((worker, activation)) => (
1034 StableProxy {
1035 state: ProxyState::EmptyAfter {
1036 previous: previous.clone(),
1037 },
1038 creations,
1039 },
1040 Self::report(ProxyOutcome::Replacement {
1041 outcome: ReplacementOutcome::WorkerAttemptsExhausted {
1042 replaces: previous,
1043 worker,
1044 activation,
1045 },
1046 }),
1047 ),
1048 }
1049 }
1050 state => {
1051 let phase = state.phase();
1052 (
1053 StableProxy {
1054 state,
1055 creations: current.creations,
1056 },
1057 Self::report(ProxyOutcome::Replacement {
1058 outcome: ReplacementOutcome::NotReplaceable {
1059 worker,
1060 activation,
1061 phase,
1062 },
1063 }),
1064 )
1065 }
1066 }
1067 }
1068
1069 fn replacement_worker_stopped(
1070 creations: CreationSequence,
1071 replacement: ProxyReplacement<W, P>,
1072 stopped: ChildStopped<BehaviorAddr<W>>,
1073 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1074 match replacement {
1075 ProxyReplacement::ReturningPredecessor {
1076 departure,
1077 successor,
1078 creation,
1079 } => match departure.worker_stopped(stopped) {
1080 Ok(ControlFlow::Continue(departure)) => {
1081 match departure.into_stopped_before_shutdown() {
1082 Ok((worker, shutdown, stopped)) => Self::start_successor_after_stop(
1083 creations,
1084 worker.attempt,
1085 PredecessorShutdown::Awaiting { id: shutdown },
1086 successor,
1087 creation,
1088 stopped,
1089 ),
1090 Err(departure) => (
1091 StableProxy {
1092 state: ProxyState::Replacing(
1093 ProxyReplacement::ReturningPredecessor {
1094 departure,
1095 successor,
1096 creation,
1097 },
1098 ),
1099 creations,
1100 },
1101 Actions::cont(),
1102 ),
1103 }
1104 }
1105 Ok(ControlFlow::Break(StoppedWorker {
1106 worker,
1107 shutdown: _,
1108 stopped,
1109 })) => Self::start_successor_after_stop(
1110 creations,
1111 worker.attempt,
1112 PredecessorShutdown::Settled,
1113 successor,
1114 creation,
1115 stopped,
1116 ),
1117 Err((departure, input)) => Self::unexpected_stop(
1118 creations,
1119 ProxyState::Replacing(ProxyReplacement::ReturningPredecessor {
1120 departure,
1121 successor,
1122 creation,
1123 }),
1124 input,
1125 ),
1126 },
1127 ProxyReplacement::SuccessorResultAwaitingShutdown {
1128 replaces,
1129 shutdown,
1130 completion: ReplacementCompletion::Ready { worker, result },
1131 } => match worker.admit_stop(stopped) {
1132 Ok(stopped) => (
1133 StableProxy {
1134 state: ProxyState::Replacing(
1135 ProxyReplacement::SuccessorResultAwaitingShutdown {
1136 replaces,
1137 shutdown,
1138 completion: ReplacementCompletion::ReadyAfterStop {
1139 worker: worker.attempt,
1140 result,
1141 stopped,
1142 },
1143 },
1144 ),
1145 creations,
1146 },
1147 Actions::cont(),
1148 ),
1149 Err(stopped) => Self::unexpected_stop(
1150 creations,
1151 ProxyState::Replacing(ProxyReplacement::SuccessorResultAwaitingShutdown {
1152 replaces,
1153 shutdown,
1154 completion: ReplacementCompletion::Ready { worker, result },
1155 }),
1156 stopped,
1157 ),
1158 },
1159 replacement => {
1160 Self::unexpected_stop(creations, ProxyState::Replacing(replacement), stopped)
1161 }
1162 }
1163 }
1164
1165 fn replacement_shutdown_resolved(
1166 creations: CreationSequence,
1167 replacement: ProxyReplacement<W, P>,
1168 shutdown: EstablishedShutdownResolved<W::Protocol>,
1169 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1170 match replacement {
1171 ProxyReplacement::ReturningPredecessor {
1172 departure,
1173 successor,
1174 creation,
1175 } => match departure.shutdown_resolved(shutdown) {
1176 Ok(ControlFlow::Continue(departure)) => (
1177 StableProxy {
1178 state: ProxyState::Replacing(ProxyReplacement::ReturningPredecessor {
1179 departure,
1180 successor,
1181 creation,
1182 }),
1183 creations,
1184 },
1185 Actions::cont(),
1186 ),
1187 Ok(ControlFlow::Break(StoppedWorker {
1188 worker,
1189 shutdown: _,
1190 stopped,
1191 })) => Self::start_successor_after_stop(
1192 creations,
1193 worker.attempt,
1194 PredecessorShutdown::Settled,
1195 successor,
1196 creation,
1197 stopped,
1198 ),
1199 Err((departure, input)) => Self::unexpected_shutdown(
1200 creations,
1201 ProxyState::Replacing(ProxyReplacement::ReturningPredecessor {
1202 departure,
1203 successor,
1204 creation,
1205 }),
1206 input,
1207 ),
1208 },
1209 ProxyReplacement::SuccessorResultAwaitingShutdown {
1210 replaces,
1211 shutdown: expected,
1212 completion,
1213 } => match shutdown.id() {
1214 id if id == expected => Self::publish_replacement(creations, replaces, completion),
1215 _ => Self::unexpected_shutdown(
1216 creations,
1217 ProxyState::Replacing(ProxyReplacement::SuccessorResultAwaitingShutdown {
1218 replaces,
1219 shutdown: expected,
1220 completion,
1221 }),
1222 shutdown,
1223 ),
1224 },
1225 }
1226 }
1227
1228 fn finish_replacement(
1229 creations: CreationSequence,
1230 replaces: WorkerAttempt,
1231 predecessor_shutdown: PredecessorShutdown,
1232 completion: ReplacementCompletion<W, P>,
1233 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1234 match predecessor_shutdown {
1235 PredecessorShutdown::Awaiting { id } => (
1236 StableProxy {
1237 state: ProxyState::Replacing(
1238 ProxyReplacement::SuccessorResultAwaitingShutdown {
1239 replaces,
1240 shutdown: id,
1241 completion,
1242 },
1243 ),
1244 creations,
1245 },
1246 Actions::cont(),
1247 ),
1248 PredecessorShutdown::Settled => {
1249 Self::publish_replacement(creations, replaces, completion)
1250 }
1251 }
1252 }
1253
1254 fn publish_replacement(
1255 creations: CreationSequence,
1256 replaces: WorkerAttempt,
1257 completion: ReplacementCompletion<W, P>,
1258 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1259 match completion {
1260 ReplacementCompletion::Ready { worker, result } => (
1261 StableProxy {
1262 state: ProxyState::Ready { worker },
1263 creations,
1264 },
1265 Self::report(ProxyOutcome::Replacement {
1266 outcome: ReplacementOutcome::Resolved { replaces, result },
1267 }),
1268 ),
1269 ReplacementCompletion::Empty { worker, result } => (
1270 StableProxy {
1271 state: ProxyState::EmptyAfter {
1272 previous: worker.clone(),
1273 },
1274 creations,
1275 },
1276 Self::report(ProxyOutcome::Replacement {
1277 outcome: ReplacementOutcome::Resolved { replaces, result },
1278 }),
1279 ),
1280 ReplacementCompletion::ReadyAfterStop {
1281 worker,
1282 result,
1283 stopped,
1284 } => (
1285 StableProxy {
1286 state: ProxyState::EmptyAfter {
1287 previous: worker.clone(),
1288 },
1289 creations,
1290 },
1291 Self::report_pair(
1292 ProxyOutcome::Replacement {
1293 outcome: ReplacementOutcome::Resolved { replaces, result },
1294 },
1295 ProxyOutcome::WorkerStopped { worker, stopped },
1296 ),
1297 ),
1298 }
1299 }
1300
1301 fn start_successor_after_stop(
1302 creations: CreationSequence,
1303 replaces: WorkerAttempt,
1304 predecessor_shutdown: PredecessorShutdown,
1305 successor: PendingWorker<P>,
1306 creation: CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
1307 stopped: ChildStopped<BehaviorAddr<W>>,
1308 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1309 let observation = ObserveChild::new(successor.creation());
1310 (
1311 StableProxy {
1312 state: ProxyState::Starting(WorkerStart {
1313 kind: WorkerStartKind::Replacement {
1314 replaces: replaces.clone(),
1315 predecessor_shutdown,
1316 },
1317 phase: WorkerStartPhase::Creating {
1318 worker: successor,
1319 stopped: None,
1320 },
1321 }),
1322 creations,
1323 },
1324 Actions::new(
1325 ProxyEffects {
1326 worker_observations: InterpreterRequests::one(observation),
1327 worker_initializations: InterpreterRequests::empty(),
1328 worker_activations: InterpreterRequests::empty(),
1329 worker_shutdowns: InterpreterRequests::empty(),
1330 worker_deliveries: Vec::new(),
1331 owner_outcomes: InterpreterRequests::one(ReportToParent::new(
1332 ProxyOutcome::WorkerStopped {
1333 worker: replaces,
1334 stopped,
1335 },
1336 )),
1337 diagnostics: InterpreterRequests::empty(),
1338 },
1339 Creations::one(creation),
1340 behavior::Step::Continue,
1341 ),
1342 )
1343 }
1344
1345 fn request_worker_shutdown(
1346 shutdown: ShutdownEstablished<StopOnShutdown<W>, Here>,
1347 ) -> ProxyActions<W, P> {
1348 Actions::send(ProxyEffects {
1349 worker_observations: InterpreterRequests::empty(),
1350 worker_initializations: InterpreterRequests::empty(),
1351 worker_activations: InterpreterRequests::empty(),
1352 worker_shutdowns: InterpreterRequests::one(shutdown),
1353 worker_deliveries: Vec::new(),
1354 owner_outcomes: InterpreterRequests::empty(),
1355 diagnostics: InterpreterRequests::empty(),
1356 })
1357 }
1358
1359 fn report_pair(earlier: ProxyOutcome<W, P>, later: ProxyOutcome<W, P>) -> ProxyActions<W, P> {
1360 Actions::send(ProxyEffects {
1361 worker_observations: InterpreterRequests::empty(),
1362 worker_initializations: InterpreterRequests::empty(),
1363 worker_activations: InterpreterRequests::empty(),
1364 worker_shutdowns: InterpreterRequests::empty(),
1365 worker_deliveries: Vec::new(),
1366 owner_outcomes: InterpreterRequests::new(vec![
1367 ReportToParent::new(earlier),
1368 ReportToParent::new(later),
1369 ]),
1370 diagnostics: InterpreterRequests::empty(),
1371 })
1372 }
1373}
1374
1375impl<W, P> StableProxy<W, P>
1377where
1378 W: Behavior,
1379 P: ActivationPlan,
1380 BehaviorAddr<W>: EndpointAddress,
1381 <BehaviorAddr<W> as Address>::Nonce: Copy + Eq,
1382 StopOnShutdown<W>:
1383 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
1384 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
1385 <StopOnShutdown<W> as Behavior>::Sends: SendSettlements,
1386 <W::Birth as BirthMode>::Child: ChildCreationProduct<BehaviorAddr<W>, behavior::ChildHead>,
1387 EstablishedRecipient<W::Protocol>: Send,
1388{
1389 fn owner_shutdown(current: StableProxy<W, P>) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1390 match current.state {
1391 ProxyState::Dormant | ProxyState::EmptyInitial => {
1392 Self::stop_with(current.creations, ProxyRetirement::Empty { previous: None })
1393 }
1394 ProxyState::EmptyAfter { previous } => Self::stop_with(
1395 current.creations,
1396 ProxyRetirement::Empty {
1397 previous: Some(previous),
1398 },
1399 ),
1400 ProxyState::Ready { worker } => {
1401 let (departure, shutdown) = WorkerStopping::begin(worker);
1402 (
1403 StableProxy {
1404 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningWorker(departure)),
1405 creations: current.creations,
1406 },
1407 Self::request_worker_shutdown(shutdown),
1408 )
1409 }
1410 ProxyState::Starting(start) => Self::shutdown_start(current.creations, start),
1411 ProxyState::Replacing(replacement) => {
1412 Self::shutdown_replacement(current.creations, replacement)
1413 }
1414 ProxyState::ShuttingDown(shutdown) => (
1415 StableProxy {
1416 state: ProxyState::ShuttingDown(shutdown),
1417 creations: current.creations,
1418 },
1419 Actions::cont(),
1420 ),
1421 ProxyState::Stopped(retirement) => Self::stop_with(current.creations, retirement),
1422 }
1423 }
1424
1425 fn shutdown_start(
1426 creations: CreationSequence,
1427 start: WorkerStart<W, P>,
1428 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1429 match start {
1430 WorkerStart {
1431 kind,
1432 phase: WorkerStartPhase::Initializing { worker, stopped },
1433 } => match stopped {
1434 Some(stopped) => {
1435 let initialization = WorkerInitializationShutdown::WaitingForInitialization {
1436 worker,
1437 shutdown: None,
1438 stopped,
1439 };
1440 (
1441 StableProxy {
1442 state: ProxyState::ShuttingDown(ProxyShutdown::Initializing {
1443 kind,
1444 initialization,
1445 }),
1446 creations,
1447 },
1448 Actions::cont(),
1449 )
1450 }
1451 None => {
1452 let (departure, shutdown) = WorkerStopping::begin(worker);
1453 let initialization = WorkerInitializationShutdown::Departing {
1454 initialization: None,
1455 departure,
1456 };
1457 (
1458 StableProxy {
1459 state: ProxyState::ShuttingDown(ProxyShutdown::Initializing {
1460 kind,
1461 initialization,
1462 }),
1463 creations,
1464 },
1465 Self::request_worker_shutdown(shutdown),
1466 )
1467 }
1468 },
1469 WorkerStart {
1470 kind,
1471 phase:
1472 WorkerStartPhase::Activating {
1473 worker,
1474 progress,
1475 stopped,
1476 },
1477 } => {
1478 let (activation, shutdown) =
1479 Self::begin_activation_shutdown(worker, progress, stopped);
1480 let actions = match shutdown {
1481 Some(shutdown) => Self::request_worker_shutdown(shutdown),
1482 None => Actions::cont(),
1483 };
1484 (
1485 StableProxy {
1486 state: ProxyState::ShuttingDown(ProxyShutdown::Activating {
1487 kind,
1488 activation,
1489 }),
1490 creations,
1491 },
1492 actions,
1493 )
1494 }
1495 WorkerStart {
1496 kind,
1497 phase: WorkerStartPhase::ReturningWorker { departure, failure },
1498 } => (
1499 StableProxy {
1500 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningWorkerStart {
1501 kind,
1502 departure,
1503 failure,
1504 }),
1505 creations,
1506 },
1507 Actions::cont(),
1508 ),
1509 start => (
1510 StableProxy {
1511 state: ProxyState::ShuttingDown(ProxyShutdown::Starting(start)),
1512 creations,
1513 },
1514 Actions::cont(),
1515 ),
1516 }
1517 }
1518
1519 fn shutdown_replacement(
1520 creations: CreationSequence,
1521 replacement: ProxyReplacement<W, P>,
1522 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1523 match replacement {
1524 ProxyReplacement::ReturningPredecessor {
1525 departure,
1526 successor,
1527 creation,
1528 } => {
1529 let replaces = departure.worker().attempt.clone();
1530 let (successor_attempt, worker, activation) = successor.cancel(creation);
1531 (
1532 StableProxy {
1533 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningPredecessor {
1534 departure,
1535 replaces: replaces.clone(),
1536 successor: successor_attempt,
1537 }),
1538 creations,
1539 },
1540 Self::report(ProxyOutcome::Replacement {
1541 outcome: ReplacementOutcome::CancelledBeforeBirth {
1542 replaces,
1543 worker,
1544 activation,
1545 },
1546 }),
1547 )
1548 }
1549 ProxyReplacement::SuccessorResultAwaitingShutdown {
1550 replaces,
1551 shutdown,
1552 completion: ReplacementCompletion::Ready { worker, result },
1553 } => {
1554 let (departure, successor_shutdown) = WorkerStopping::begin(worker);
1555 (
1556 StableProxy {
1557 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningSuccessor {
1558 replaces,
1559 predecessor: PredecessorReturn::Awaiting { shutdown },
1560 departure,
1561 result,
1562 }),
1563 creations,
1564 },
1565 Self::request_worker_shutdown(successor_shutdown),
1566 )
1567 }
1568 ProxyReplacement::SuccessorResultAwaitingShutdown {
1569 replaces,
1570 shutdown,
1571 completion,
1572 } => (
1573 StableProxy {
1574 state: ProxyState::ShuttingDown(ProxyShutdown::WaitingForPredecessor {
1575 replaces,
1576 shutdown,
1577 retirement: WorkerStartRetirement::ReplacementCompletion(completion),
1578 }),
1579 creations,
1580 },
1581 Actions::cont(),
1582 ),
1583 }
1584 }
1585
1586 fn shutdown_workers_created(
1587 creations: CreationSequence,
1588 shutdown: ProxyShutdown<W, P>,
1589 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
1590 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1591 match shutdown {
1592 ProxyShutdown::Starting(WorkerStart {
1593 kind,
1594 phase: WorkerStartPhase::Creating { worker, stopped },
1595 }) => match worker.created(workers, stopped) {
1596 WorkerCreation::Initializing {
1597 worker,
1598 activation,
1599 stopped,
1600 } => Self::begin_shutdown_initialization(
1601 creations, kind, worker, activation, stopped,
1602 ),
1603 WorkerCreation::Rejected {
1604 rejection,
1605 activation,
1606 stopped,
1607 } => Self::finish_start_retirement(
1608 creations,
1609 kind,
1610 WorkerStartRetirement::Result(WorkerStartResult::CreationRejected {
1611 rejection,
1612 activation,
1613 stopped,
1614 }),
1615 ),
1616 WorkerCreation::Unexpected {
1617 worker,
1618 stopped,
1619 workers,
1620 } => Self::unexpected_worker_start(
1621 creations,
1622 ProxyShutdown::Starting(WorkerStart {
1623 kind,
1624 phase: WorkerStartPhase::Creating { worker, stopped },
1625 }),
1626 workers,
1627 ),
1628 },
1629 shutdown => Self::unexpected_worker_start(creations, shutdown, workers),
1630 }
1631 }
1632
1633 fn begin_shutdown_initialization(
1634 creations: CreationSequence,
1635 kind: WorkerStartKind,
1636 worker: CurrentWorker<W>,
1637 activation: P,
1638 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
1639 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1640 let request = InitializeWorker::new(
1641 worker.attempt.clone(),
1642 worker.initialization.clone(),
1643 worker.actor.recipient(),
1644 activation,
1645 );
1646 match stopped {
1647 Some(stopped) => {
1648 let initialization = WorkerInitializationShutdown::WaitingForInitialization {
1649 worker,
1650 shutdown: None,
1651 stopped,
1652 };
1653 (
1654 StableProxy {
1655 state: ProxyState::ShuttingDown(ProxyShutdown::Initializing {
1656 kind,
1657 initialization,
1658 }),
1659 creations,
1660 },
1661 Self::request_initialization(request),
1662 )
1663 }
1664 None => {
1665 let (departure, shutdown) = WorkerStopping::begin(worker);
1666 let initialization = WorkerInitializationShutdown::Departing {
1667 initialization: None,
1668 departure,
1669 };
1670 (
1671 StableProxy {
1672 state: ProxyState::ShuttingDown(ProxyShutdown::Initializing {
1673 kind,
1674 initialization,
1675 }),
1676 creations,
1677 },
1678 Self::initialize_and_shutdown(request, shutdown),
1679 )
1680 }
1681 }
1682 }
1683
1684 fn shutdown_worker_initialized(
1685 creations: CreationSequence,
1686 shutdown: ProxyShutdown<W, P>,
1687 initialization: WorkerInitializationReport<W, P>,
1688 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1689 match shutdown {
1690 ProxyShutdown::Initializing {
1691 kind,
1692 initialization: current,
1693 } => match Self::continue_initialization_shutdown(
1694 creations,
1695 kind,
1696 Self::admit_worker_initialized_during_shutdown(current, initialization),
1697 ) {
1698 Ok(changed) => changed,
1699 Err((creations, shutdown, input)) => {
1700 Self::unexpected_initialization(creations, shutdown, input)
1701 }
1702 },
1703 shutdown => Self::unexpected_initialization(creations, shutdown, initialization),
1704 }
1705 }
1706
1707 fn unexpected_initialization(
1708 creations: CreationSequence,
1709 shutdown: ProxyShutdown<W, P>,
1710 initialization: WorkerInitializationReport<W, P>,
1711 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
1712 (
1713 StableProxy {
1714 state: ProxyState::ShuttingDown(shutdown),
1715 creations,
1716 },
1717 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerInitialization {
1718 phase: ProxyPhase::ShuttingDown,
1719 initialization,
1720 }),
1721 )
1722 }
1723
1724 fn continue_initialization_shutdown<I>(
1725 creations: CreationSequence,
1726 kind: WorkerStartKind,
1727 changed: Result<
1728 ControlFlow<WorkerStartRetirement<W, P>, WorkerInitializationShutdown<W, P>>,
1729 (WorkerInitializationShutdown<W, P>, I),
1730 >,
1731 ) -> Result<(StableProxy<W, P>, ProxyActions<W, P>), (CreationSequence, ProxyShutdown<W, P>, I)>
1732 {
1733 match changed {
1734 Ok(ControlFlow::Continue(initialization)) => Ok((
1735 StableProxy {
1736 state: ProxyState::ShuttingDown(ProxyShutdown::Initializing {
1737 kind,
1738 initialization,
1739 }),
1740 creations,
1741 },
1742 Actions::cont(),
1743 )),
1744 Ok(ControlFlow::Break(retirement)) => {
1745 Ok(Self::finish_start_retirement(creations, kind, retirement))
1746 }
1747 Err((initialization, input)) => Err((
1748 creations,
1749 ProxyShutdown::Initializing {
1750 kind,
1751 initialization,
1752 },
1753 input,
1754 )),
1755 }
1756 }
1757
1758 fn admit_worker_initialized_during_shutdown(
1759 current: WorkerInitializationShutdown<W, P>,
1760 input: WorkerInitializationReport<W, P>,
1761 ) -> Result<
1762 ControlFlow<WorkerStartRetirement<W, P>, WorkerInitializationShutdown<W, P>>,
1763 (
1764 WorkerInitializationShutdown<W, P>,
1765 WorkerInitializationReport<W, P>,
1766 ),
1767 > {
1768 match current {
1769 WorkerInitializationShutdown::Departing {
1770 initialization: None,
1771 departure,
1772 } => match departure.worker().admit_initialization(input) {
1773 Ok(input) => Self::retain_initialization_while_departing(departure, input),
1774 Err(input) => Err((
1775 WorkerInitializationShutdown::Departing {
1776 initialization: None,
1777 departure,
1778 },
1779 input,
1780 )),
1781 },
1782 WorkerInitializationShutdown::Departing {
1783 initialization,
1784 departure,
1785 } => Err((
1786 WorkerInitializationShutdown::Departing {
1787 initialization,
1788 departure,
1789 },
1790 input,
1791 )),
1792 WorkerInitializationShutdown::WaitingForInitialization {
1793 worker,
1794 shutdown,
1795 stopped,
1796 } => match worker.admit_initialization(input) {
1797 Ok(input) => match Self::classify_initialization_after_stop(&worker, input) {
1798 Ok(initialization) => {
1799 Ok(ControlFlow::Break(WorkerStartRetirement::Initialization {
1800 initialization,
1801 worker,
1802 shutdown,
1803 stopped,
1804 }))
1805 }
1806 Err(input) => Err((
1807 WorkerInitializationShutdown::WaitingForInitialization {
1808 worker,
1809 shutdown,
1810 stopped,
1811 },
1812 input,
1813 )),
1814 },
1815 Err(input) => Err((
1816 WorkerInitializationShutdown::WaitingForInitialization {
1817 worker,
1818 shutdown,
1819 stopped,
1820 },
1821 input,
1822 )),
1823 },
1824 }
1825 }
1826
1827 fn rejoin_initialization_departure<I>(
1828 initialization: Option<WorkerInitializationRetirement<W, P>>,
1829 admitted: Result<ControlFlow<StoppedWorker<W>, WorkerStopping<W>>, (WorkerStopping<W>, I)>,
1830 ) -> Result<
1831 ControlFlow<WorkerStartRetirement<W, P>, WorkerInitializationShutdown<W, P>>,
1832 (WorkerInitializationShutdown<W, P>, I),
1833 > {
1834 match admitted {
1835 Ok(ControlFlow::Continue(departure)) => Ok(ControlFlow::Continue(
1836 WorkerInitializationShutdown::Departing {
1837 initialization,
1838 departure,
1839 },
1840 )),
1841 Ok(ControlFlow::Break(departure)) => match initialization {
1842 None => {
1843 let StoppedWorker {
1844 worker,
1845 shutdown,
1846 stopped,
1847 } = departure;
1848 Ok(ControlFlow::Continue(
1849 WorkerInitializationShutdown::WaitingForInitialization {
1850 worker,
1851 shutdown: Some(shutdown),
1852 stopped,
1853 },
1854 ))
1855 }
1856 Some(initialization) => {
1857 let StoppedWorker {
1858 worker,
1859 shutdown,
1860 stopped,
1861 } = departure;
1862 Ok(ControlFlow::Break(WorkerStartRetirement::Initialization {
1863 initialization,
1864 worker,
1865 shutdown: Some(shutdown),
1866 stopped,
1867 }))
1868 }
1869 },
1870 Err((departure, input)) => Err((
1871 WorkerInitializationShutdown::Departing {
1872 initialization,
1873 departure,
1874 },
1875 input,
1876 )),
1877 }
1878 }
1879
1880 fn retain_initialization_while_departing(
1881 departure: WorkerStopping<W>,
1882 input: WorkerInitializationReport<W, P>,
1883 ) -> Result<
1884 ControlFlow<WorkerStartRetirement<W, P>, WorkerInitializationShutdown<W, P>>,
1885 (
1886 WorkerInitializationShutdown<W, P>,
1887 WorkerInitializationReport<W, P>,
1888 ),
1889 > {
1890 match input {
1891 WorkerInitializationReport::ReadyForActivation {
1892 permit, activation, ..
1893 } => Ok(ControlFlow::Continue(
1894 Self::wait_for_initialization_departure(
1895 departure,
1896 WorkerInitializationRetirement::Initialized { permit, activation },
1897 ),
1898 )),
1899 WorkerInitializationReport::EffectsRejected {
1900 failure,
1901 activation,
1902 ..
1903 } => Ok(ControlFlow::Continue(
1904 Self::wait_for_initialization_departure(
1905 departure,
1906 WorkerInitializationRetirement::EffectsRejected {
1907 failure,
1908 activation,
1909 },
1910 ),
1911 )),
1912 WorkerInitializationReport::Stopped {
1913 worker,
1914 initialization,
1915 activation,
1916 stopped,
1917 } => Self::retain_initialization_stop(
1918 departure,
1919 worker,
1920 initialization,
1921 activation,
1922 stopped,
1923 ),
1924 }
1925 }
1926
1927 fn retain_initialization_stop(
1928 departure: WorkerStopping<W>,
1929 worker: WorkerAttempt,
1930 initialization: InitializationAttempt,
1931 activation: P,
1932 stopped: ChildStopped<BehaviorAddr<W>>,
1933 ) -> Result<
1934 ControlFlow<WorkerStartRetirement<W, P>, WorkerInitializationShutdown<W, P>>,
1935 (
1936 WorkerInitializationShutdown<W, P>,
1937 WorkerInitializationReport<W, P>,
1938 ),
1939 > {
1940 match departure.stopped() {
1941 Some(_) => match departure.worker().admit_stop(stopped) {
1942 Ok(stopped) => Ok(ControlFlow::Continue(
1943 Self::wait_for_initialization_departure(
1944 departure,
1945 WorkerInitializationRetirement::Stopped {
1946 activation,
1947 initialization_stop: Some(stopped),
1948 },
1949 ),
1950 )),
1951 Err(stopped) => Err(Self::return_unrelated_initialization_stop(
1952 departure,
1953 worker,
1954 initialization,
1955 activation,
1956 stopped,
1957 )),
1958 },
1959 None => match departure.worker_stopped(stopped) {
1960 Ok(ControlFlow::Continue(departure)) => Ok(ControlFlow::Continue(
1961 Self::wait_for_initialization_departure(
1962 departure,
1963 WorkerInitializationRetirement::Stopped {
1964 activation,
1965 initialization_stop: None,
1966 },
1967 ),
1968 )),
1969 Ok(ControlFlow::Break(StoppedWorker {
1970 worker,
1971 shutdown,
1972 stopped,
1973 })) => Ok(ControlFlow::Break(WorkerStartRetirement::Initialization {
1974 initialization: WorkerInitializationRetirement::Stopped {
1975 activation,
1976 initialization_stop: None,
1977 },
1978 worker,
1979 shutdown: Some(shutdown),
1980 stopped,
1981 })),
1982 Err((departure, input)) => Err(Self::return_unrelated_initialization_stop(
1983 departure,
1984 worker,
1985 initialization,
1986 activation,
1987 input,
1988 )),
1989 },
1990 }
1991 }
1992
1993 fn classify_initialization_after_stop(
1994 worker: &CurrentWorker<W>,
1995 input: WorkerInitializationReport<W, P>,
1996 ) -> Result<WorkerInitializationRetirement<W, P>, WorkerInitializationReport<W, P>> {
1997 match input {
1998 WorkerInitializationReport::ReadyForActivation {
1999 permit, activation, ..
2000 } => Ok(WorkerInitializationRetirement::Initialized { permit, activation }),
2001 WorkerInitializationReport::EffectsRejected {
2002 failure,
2003 activation,
2004 ..
2005 } => Ok(WorkerInitializationRetirement::EffectsRejected {
2006 failure,
2007 activation,
2008 }),
2009 WorkerInitializationReport::Stopped {
2010 worker: attempt,
2011 initialization,
2012 activation,
2013 stopped,
2014 } => match worker.admit_stop(stopped) {
2015 Ok(stopped) => Ok(WorkerInitializationRetirement::Stopped {
2016 activation,
2017 initialization_stop: Some(stopped),
2018 }),
2019 Err(stopped) => Err(WorkerInitializationReport::Stopped {
2020 worker: attempt,
2021 initialization,
2022 activation,
2023 stopped,
2024 }),
2025 },
2026 }
2027 }
2028
2029 fn wait_for_initialization_departure(
2030 departure: WorkerStopping<W>,
2031 initialization: WorkerInitializationRetirement<W, P>,
2032 ) -> WorkerInitializationShutdown<W, P> {
2033 WorkerInitializationShutdown::Departing {
2034 initialization: Some(initialization),
2035 departure,
2036 }
2037 }
2038
2039 fn return_unrelated_initialization_stop(
2040 departure: WorkerStopping<W>,
2041 worker: WorkerAttempt,
2042 initialization: InitializationAttempt,
2043 activation: P,
2044 stopped: ChildStopped<BehaviorAddr<W>>,
2045 ) -> (
2046 WorkerInitializationShutdown<W, P>,
2047 WorkerInitializationReport<W, P>,
2048 ) {
2049 (
2050 WorkerInitializationShutdown::Departing {
2051 initialization: None,
2052 departure,
2053 },
2054 WorkerInitializationReport::Stopped {
2055 worker,
2056 initialization,
2057 activation,
2058 stopped,
2059 },
2060 )
2061 }
2062
2063 fn begin_activation_shutdown(
2064 worker: CurrentWorker<W>,
2065 activation: ActivationProgress<W, P>,
2066 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
2067 ) -> (
2068 WorkerActivationShutdown<W, P>,
2069 Option<crate::ShutdownEstablished<StopOnShutdown<W>, behavior::Here>>,
2070 ) {
2071 match stopped {
2072 Some(stopped) => (
2073 WorkerActivationShutdown::WaitingForActivation {
2074 activation,
2075 worker,
2076 shutdown: None,
2077 stopped,
2078 },
2079 None,
2080 ),
2081 None => {
2082 let (departure, shutdown) = WorkerStopping::begin(worker);
2083 (
2084 WorkerActivationShutdown::Departing {
2085 activation: ActivationDuringDeparture::Pending(activation),
2086 departure,
2087 },
2088 Some(shutdown),
2089 )
2090 }
2091 }
2092 }
2093
2094 fn continue_activation_shutdown<I>(
2095 creations: CreationSequence,
2096 kind: WorkerStartKind,
2097 changed: Result<
2098 ControlFlow<WorkerStartRetirement<W, P>, WorkerActivationShutdown<W, P>>,
2099 (WorkerActivationShutdown<W, P>, I),
2100 >,
2101 ) -> Result<(StableProxy<W, P>, ProxyActions<W, P>), (CreationSequence, ProxyShutdown<W, P>, I)>
2102 {
2103 match changed {
2104 Ok(ControlFlow::Continue(activation)) => Ok((
2105 StableProxy {
2106 state: ProxyState::ShuttingDown(ProxyShutdown::Activating { kind, activation }),
2107 creations,
2108 },
2109 Actions::cont(),
2110 )),
2111 Ok(ControlFlow::Break(retirement)) => {
2112 Ok(Self::finish_start_retirement(creations, kind, retirement))
2113 }
2114 Err((activation, input)) => Err((
2115 creations,
2116 ProxyShutdown::Activating { kind, activation },
2117 input,
2118 )),
2119 }
2120 }
2121
2122 fn admit_worker_activation_during_shutdown(
2123 current: WorkerActivationShutdown<W, P>,
2124 input: WorkerActivation<W, P>,
2125 ) -> Result<
2126 ControlFlow<WorkerStartRetirement<W, P>, WorkerActivationShutdown<W, P>>,
2127 (WorkerActivationShutdown<W, P>, WorkerActivation<W, P>),
2128 > {
2129 match current {
2130 WorkerActivationShutdown::Departing {
2131 activation: ActivationDuringDeparture::Pending(progress),
2132 departure,
2133 } => match Self::admit_activation(&progress, input) {
2134 Ok(input) => Self::retain_activation_while_departing(progress, departure, input),
2135 Err(input) => Err((
2136 WorkerActivationShutdown::Departing {
2137 activation: ActivationDuringDeparture::Pending(progress),
2138 departure,
2139 },
2140 input,
2141 )),
2142 },
2143 WorkerActivationShutdown::Departing {
2144 activation,
2145 departure,
2146 } => Err((
2147 WorkerActivationShutdown::Departing {
2148 activation,
2149 departure,
2150 },
2151 input,
2152 )),
2153 WorkerActivationShutdown::WaitingForActivation {
2154 activation,
2155 worker,
2156 shutdown,
2157 stopped,
2158 } => match Self::admit_activation(&activation, input) {
2159 Ok(input) => Self::retain_activation_after_worker(
2160 activation, worker, shutdown, stopped, input,
2161 ),
2162 Err(input) => Err((
2163 WorkerActivationShutdown::WaitingForActivation {
2164 activation,
2165 worker,
2166 shutdown,
2167 stopped,
2168 },
2169 input,
2170 )),
2171 },
2172 }
2173 }
2174
2175 fn advance_activation_during_shutdown(
2176 progress: ActivationProgress<W, P>,
2177 input: WorkerActivation<W, P>,
2178 ) -> Result<ActivationDuringDeparture<W, P>, (ActivationProgress<W, P>, WorkerActivation<W, P>)>
2179 {
2180 match progress {
2181 ActivationProgress::WaitingForStart(attempt) => match input.into_started() {
2182 Ok(()) => Ok(ActivationDuringDeparture::Pending(
2183 ActivationProgress::Running(attempt),
2184 )),
2185 Err(input) => match input.into_start_rejection() {
2186 Ok((request, reason)) => Ok(ActivationDuringDeparture::Returned(
2187 WorkerActivationRetirement::StartRejected { request, reason },
2188 )),
2189 Err(input) => Err((ActivationProgress::WaitingForStart(attempt), input)),
2190 },
2191 },
2192 ActivationProgress::Running(attempt) => match input.into_ready() {
2193 Ok(readiness) => Ok(ActivationDuringDeparture::Returned(
2194 WorkerActivationRetirement::Ready { readiness },
2195 )),
2196 Err(input) => match input.into_rejection() {
2197 Ok(rejection) => Ok(ActivationDuringDeparture::Returned(
2198 WorkerActivationRetirement::Rejected { rejection },
2199 )),
2200 Err(input) => Err((ActivationProgress::Running(attempt), input)),
2201 },
2202 },
2203 }
2204 }
2205
2206 fn retain_activation_while_departing(
2207 progress: ActivationProgress<W, P>,
2208 departure: WorkerStopping<W>,
2209 input: WorkerActivation<W, P>,
2210 ) -> Result<
2211 ControlFlow<WorkerStartRetirement<W, P>, WorkerActivationShutdown<W, P>>,
2212 (WorkerActivationShutdown<W, P>, WorkerActivation<W, P>),
2213 > {
2214 match Self::advance_activation_during_shutdown(progress, input) {
2215 Ok(activation) => Ok(ControlFlow::Continue(WorkerActivationShutdown::Departing {
2216 activation,
2217 departure,
2218 })),
2219 Err((progress, input)) => Err((
2220 WorkerActivationShutdown::Departing {
2221 activation: ActivationDuringDeparture::Pending(progress),
2222 departure,
2223 },
2224 input,
2225 )),
2226 }
2227 }
2228
2229 fn retain_activation_after_worker(
2230 progress: ActivationProgress<W, P>,
2231 worker: CurrentWorker<W>,
2232 shutdown: Option<EstablishedShutdownResolved<W::Protocol>>,
2233 stopped: ChildStopped<BehaviorAddr<W>>,
2234 input: WorkerActivation<W, P>,
2235 ) -> Result<
2236 ControlFlow<WorkerStartRetirement<W, P>, WorkerActivationShutdown<W, P>>,
2237 (WorkerActivationShutdown<W, P>, WorkerActivation<W, P>),
2238 > {
2239 match Self::advance_activation_during_shutdown(progress, input) {
2240 Ok(ActivationDuringDeparture::Pending(activation)) => Ok(ControlFlow::Continue(
2241 WorkerActivationShutdown::WaitingForActivation {
2242 activation,
2243 worker,
2244 shutdown,
2245 stopped,
2246 },
2247 )),
2248 Ok(ActivationDuringDeparture::Returned(activation)) => {
2249 Ok(ControlFlow::Break(WorkerStartRetirement::Activation {
2250 activation,
2251 worker,
2252 shutdown,
2253 stopped,
2254 }))
2255 }
2256 Err((progress, input)) => Err((
2257 WorkerActivationShutdown::WaitingForActivation {
2258 activation: progress,
2259 worker,
2260 shutdown,
2261 stopped,
2262 },
2263 input,
2264 )),
2265 }
2266 }
2267
2268 fn rejoin_activation_departure<I>(
2269 activation: ActivationDuringDeparture<W, P>,
2270 admitted: Result<ControlFlow<StoppedWorker<W>, WorkerStopping<W>>, (WorkerStopping<W>, I)>,
2271 ) -> Result<
2272 ControlFlow<WorkerStartRetirement<W, P>, WorkerActivationShutdown<W, P>>,
2273 (WorkerActivationShutdown<W, P>, I),
2274 > {
2275 match admitted {
2276 Ok(ControlFlow::Continue(departure)) => {
2277 Ok(ControlFlow::Continue(WorkerActivationShutdown::Departing {
2278 activation,
2279 departure,
2280 }))
2281 }
2282 Ok(ControlFlow::Break(StoppedWorker {
2283 worker,
2284 shutdown,
2285 stopped,
2286 })) => match activation {
2287 ActivationDuringDeparture::Pending(activation) => Ok(ControlFlow::Continue(
2288 WorkerActivationShutdown::WaitingForActivation {
2289 activation,
2290 worker,
2291 shutdown: Some(shutdown),
2292 stopped,
2293 },
2294 )),
2295 ActivationDuringDeparture::Returned(activation) => {
2296 Ok(ControlFlow::Break(WorkerStartRetirement::Activation {
2297 activation,
2298 worker,
2299 shutdown: Some(shutdown),
2300 stopped,
2301 }))
2302 }
2303 },
2304 Err((departure, input)) => Err((
2305 WorkerActivationShutdown::Departing {
2306 activation,
2307 departure,
2308 },
2309 input,
2310 )),
2311 }
2312 }
2313
2314 fn shutdown_worker_activated(
2315 creations: CreationSequence,
2316 shutdown: ProxyShutdown<W, P>,
2317 input: WorkerActivation<W, P>,
2318 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2319 match shutdown {
2320 ProxyShutdown::Activating { kind, activation } => {
2321 match Self::continue_activation_shutdown(
2322 creations,
2323 kind,
2324 Self::admit_worker_activation_during_shutdown(activation, input),
2325 ) {
2326 Ok(changed) => changed,
2327 Err((creations, shutdown, input)) => {
2328 Self::unexpected_shutdown_activation(creations, shutdown, input)
2329 }
2330 }
2331 }
2332 shutdown => Self::unexpected_shutdown_activation(creations, shutdown, input),
2333 }
2334 }
2335
2336 fn unexpected_shutdown_activation(
2337 creations: CreationSequence,
2338 shutdown: ProxyShutdown<W, P>,
2339 activation: WorkerActivation<W, P>,
2340 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2341 (
2342 StableProxy {
2343 state: ProxyState::ShuttingDown(shutdown),
2344 creations,
2345 },
2346 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerActivation {
2347 phase: ProxyPhase::ShuttingDown,
2348 activation,
2349 }),
2350 )
2351 }
2352
2353 fn shutdown_resolved(
2354 creations: CreationSequence,
2355 shutdown: ProxyShutdown<W, P>,
2356 resolution: EstablishedShutdownResolved<W::Protocol>,
2357 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2358 match shutdown {
2359 ProxyShutdown::ReturningWorker(departure) => {
2360 match departure.shutdown_resolved(resolution) {
2361 Ok(ControlFlow::Continue(departure)) => (
2362 StableProxy {
2363 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningWorker(
2364 departure,
2365 )),
2366 creations,
2367 },
2368 Actions::cont(),
2369 ),
2370 Ok(ControlFlow::Break(StoppedWorker {
2371 worker,
2372 shutdown,
2373 stopped,
2374 })) => Self::finish_shutdown(creations, worker, shutdown, stopped),
2375 Err((departure, input)) => Self::unexpected_shutdown(
2376 creations,
2377 ProxyState::ShuttingDown(ProxyShutdown::ReturningWorker(departure)),
2378 input,
2379 ),
2380 }
2381 }
2382 ProxyShutdown::ReturningWorkerStart {
2383 kind,
2384 departure,
2385 failure,
2386 } => match departure.shutdown_resolved(resolution) {
2387 Ok(ControlFlow::Continue(departure)) => (
2388 StableProxy {
2389 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningWorkerStart {
2390 kind,
2391 departure,
2392 failure,
2393 }),
2394 creations,
2395 },
2396 Actions::cont(),
2397 ),
2398 Ok(ControlFlow::Break(StoppedWorker {
2399 worker,
2400 shutdown,
2401 stopped,
2402 })) => {
2403 let result = Self::worker_return_result(worker, failure, shutdown, stopped);
2404 Self::finish_start_retirement(
2405 creations,
2406 kind,
2407 WorkerStartRetirement::Result(result),
2408 )
2409 }
2410 Err((departure, input)) => Self::unexpected_shutdown(
2411 creations,
2412 ProxyState::ShuttingDown(ProxyShutdown::ReturningWorkerStart {
2413 kind,
2414 departure,
2415 failure,
2416 }),
2417 input,
2418 ),
2419 },
2420 ProxyShutdown::ReturningPredecessor {
2421 departure,
2422 replaces,
2423 successor,
2424 } => match departure.shutdown_resolved(resolution) {
2425 Ok(ControlFlow::Continue(departure)) => (
2426 StableProxy {
2427 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningPredecessor {
2428 departure,
2429 replaces,
2430 successor,
2431 }),
2432 creations,
2433 },
2434 Actions::cont(),
2435 ),
2436 Ok(ControlFlow::Break(StoppedWorker {
2437 worker,
2438 shutdown,
2439 stopped,
2440 })) => Self::finish_cancelled_replacement(
2441 creations, replaces, successor, worker, shutdown, stopped,
2442 ),
2443 Err((departure, input)) => Self::unexpected_shutdown(
2444 creations,
2445 ProxyState::ShuttingDown(ProxyShutdown::ReturningPredecessor {
2446 departure,
2447 replaces,
2448 successor,
2449 }),
2450 input,
2451 ),
2452 },
2453 ProxyShutdown::ReturningSuccessor {
2454 replaces,
2455 predecessor: PredecessorReturn::Awaiting { shutdown: expected },
2456 departure,
2457 result,
2458 } => match resolution.id() {
2459 id if id == expected => (
2460 StableProxy {
2461 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningSuccessor {
2462 replaces,
2463 predecessor: PredecessorReturn::Returned { resolution },
2464 departure,
2465 result,
2466 }),
2467 creations,
2468 },
2469 Actions::cont(),
2470 ),
2471 _ => Self::successor_shutdown_resolved(
2472 creations,
2473 replaces,
2474 PredecessorReturn::Awaiting { shutdown: expected },
2475 departure,
2476 result,
2477 resolution,
2478 ),
2479 },
2480 ProxyShutdown::ReturningSuccessor {
2481 replaces,
2482 predecessor,
2483 departure,
2484 result,
2485 } => Self::successor_shutdown_resolved(
2486 creations,
2487 replaces,
2488 predecessor,
2489 departure,
2490 result,
2491 resolution,
2492 ),
2493 ProxyShutdown::Initializing {
2494 kind,
2495 initialization,
2496 } => match initialization {
2497 WorkerInitializationShutdown::Departing {
2498 initialization,
2499 departure,
2500 } => match Self::continue_initialization_shutdown(
2501 creations,
2502 kind,
2503 Self::rejoin_initialization_departure(
2504 initialization,
2505 departure.shutdown_resolved(resolution),
2506 ),
2507 ) {
2508 Ok(changed) => changed,
2509 Err((creations, shutdown, input)) => Self::unexpected_shutdown(
2510 creations,
2511 ProxyState::ShuttingDown(shutdown),
2512 input,
2513 ),
2514 },
2515 initialization => Self::unexpected_shutdown(
2516 creations,
2517 ProxyState::ShuttingDown(ProxyShutdown::Initializing {
2518 kind,
2519 initialization,
2520 }),
2521 resolution,
2522 ),
2523 },
2524 ProxyShutdown::Activating { kind, activation } => {
2525 let changed = match activation {
2526 WorkerActivationShutdown::Departing {
2527 activation,
2528 departure,
2529 } => Self::rejoin_activation_departure(
2530 activation,
2531 departure.shutdown_resolved(resolution),
2532 ),
2533 activation => Err((activation, resolution)),
2534 };
2535 match Self::continue_activation_shutdown(creations, kind, changed) {
2536 Ok(changed) => changed,
2537 Err((creations, shutdown, input)) => Self::unexpected_shutdown(
2538 creations,
2539 ProxyState::ShuttingDown(shutdown),
2540 input,
2541 ),
2542 }
2543 }
2544 ProxyShutdown::WaitingForPredecessor {
2545 replaces,
2546 shutdown: expected,
2547 retirement,
2548 } => match resolution.id() {
2549 id if id == expected => Self::stop_with(
2550 creations,
2551 ProxyRetirement::ReplacementAfterPredecessor {
2552 replaces,
2553 predecessor: resolution,
2554 retirement,
2555 },
2556 ),
2557 _ => Self::unexpected_shutdown(
2558 creations,
2559 ProxyState::ShuttingDown(ProxyShutdown::WaitingForPredecessor {
2560 replaces,
2561 shutdown: expected,
2562 retirement,
2563 }),
2564 resolution,
2565 ),
2566 },
2567 shutdown => {
2568 Self::unexpected_shutdown(creations, ProxyState::ShuttingDown(shutdown), resolution)
2569 }
2570 }
2571 }
2572
2573 fn successor_shutdown_resolved(
2574 creations: CreationSequence,
2575 replaces: WorkerAttempt,
2576 predecessor: PredecessorReturn<W>,
2577 departure: WorkerStopping<W>,
2578 result: WorkerStartResult<W, P>,
2579 resolution: EstablishedShutdownResolved<W::Protocol>,
2580 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2581 match departure.shutdown_resolved(resolution) {
2582 Ok(ControlFlow::Continue(departure)) => (
2583 StableProxy {
2584 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningSuccessor {
2585 replaces,
2586 predecessor,
2587 departure,
2588 result,
2589 }),
2590 creations,
2591 },
2592 Actions::cont(),
2593 ),
2594 Ok(ControlFlow::Break(departure)) => {
2595 Self::successor_departed(creations, replaces, predecessor, result, departure)
2596 }
2597 Err((departure, input)) => Self::unexpected_shutdown(
2598 creations,
2599 ProxyState::ShuttingDown(ProxyShutdown::ReturningSuccessor {
2600 replaces,
2601 predecessor,
2602 departure,
2603 result,
2604 }),
2605 input,
2606 ),
2607 }
2608 }
2609
2610 fn successor_departed(
2611 creations: CreationSequence,
2612 replaces: WorkerAttempt,
2613 predecessor: PredecessorReturn<W>,
2614 result: WorkerStartResult<W, P>,
2615 departure: StoppedWorker<W>,
2616 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2617 let retirement = WorkerStartRetirement::ReadyWorker { result, departure };
2618 match predecessor {
2619 PredecessorReturn::Awaiting { shutdown } => (
2620 StableProxy {
2621 state: ProxyState::ShuttingDown(ProxyShutdown::WaitingForPredecessor {
2622 replaces,
2623 shutdown,
2624 retirement,
2625 }),
2626 creations,
2627 },
2628 Actions::cont(),
2629 ),
2630 PredecessorReturn::Returned { resolution } => Self::stop_with(
2631 creations,
2632 ProxyRetirement::ReplacementAfterPredecessor {
2633 replaces,
2634 predecessor: resolution,
2635 retirement,
2636 },
2637 ),
2638 }
2639 }
2640
2641 fn shutdown_worker_stopped(
2642 creations: CreationSequence,
2643 shutdown: ProxyShutdown<W, P>,
2644 stopped: ChildStopped<BehaviorAddr<W>>,
2645 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2646 match shutdown {
2647 ProxyShutdown::ReturningWorker(departure) => match departure.worker_stopped(stopped) {
2648 Ok(ControlFlow::Continue(departure)) => (
2649 StableProxy {
2650 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningWorker(departure)),
2651 creations,
2652 },
2653 Actions::cont(),
2654 ),
2655 Ok(ControlFlow::Break(StoppedWorker {
2656 worker,
2657 shutdown,
2658 stopped,
2659 })) => Self::finish_shutdown(creations, worker, shutdown, stopped),
2660 Err((departure, input)) => Self::unexpected_stop(
2661 creations,
2662 ProxyState::ShuttingDown(ProxyShutdown::ReturningWorker(departure)),
2663 input,
2664 ),
2665 },
2666 ProxyShutdown::ReturningWorkerStart {
2667 kind,
2668 departure,
2669 failure,
2670 } => match departure.worker_stopped(stopped) {
2671 Ok(ControlFlow::Continue(departure)) => (
2672 StableProxy {
2673 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningWorkerStart {
2674 kind,
2675 departure,
2676 failure,
2677 }),
2678 creations,
2679 },
2680 Actions::cont(),
2681 ),
2682 Ok(ControlFlow::Break(StoppedWorker {
2683 worker,
2684 shutdown,
2685 stopped,
2686 })) => {
2687 let result = Self::worker_return_result(worker, failure, shutdown, stopped);
2688 Self::finish_start_retirement(
2689 creations,
2690 kind,
2691 WorkerStartRetirement::Result(result),
2692 )
2693 }
2694 Err((departure, input)) => Self::unexpected_stop(
2695 creations,
2696 ProxyState::ShuttingDown(ProxyShutdown::ReturningWorkerStart {
2697 kind,
2698 departure,
2699 failure,
2700 }),
2701 input,
2702 ),
2703 },
2704 ProxyShutdown::ReturningPredecessor {
2705 departure,
2706 replaces,
2707 successor,
2708 } => match departure.worker_stopped(stopped) {
2709 Ok(ControlFlow::Continue(departure)) => (
2710 StableProxy {
2711 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningPredecessor {
2712 departure,
2713 replaces,
2714 successor,
2715 }),
2716 creations,
2717 },
2718 Actions::cont(),
2719 ),
2720 Ok(ControlFlow::Break(StoppedWorker {
2721 worker,
2722 shutdown,
2723 stopped,
2724 })) => Self::finish_cancelled_replacement(
2725 creations, replaces, successor, worker, shutdown, stopped,
2726 ),
2727 Err((departure, input)) => Self::unexpected_stop(
2728 creations,
2729 ProxyState::ShuttingDown(ProxyShutdown::ReturningPredecessor {
2730 departure,
2731 replaces,
2732 successor,
2733 }),
2734 input,
2735 ),
2736 },
2737 ProxyShutdown::ReturningSuccessor {
2738 replaces,
2739 predecessor,
2740 departure,
2741 result,
2742 } => match departure.worker_stopped(stopped) {
2743 Ok(ControlFlow::Continue(departure)) => (
2744 StableProxy {
2745 state: ProxyState::ShuttingDown(ProxyShutdown::ReturningSuccessor {
2746 replaces,
2747 predecessor,
2748 departure,
2749 result,
2750 }),
2751 creations,
2752 },
2753 Actions::cont(),
2754 ),
2755 Ok(ControlFlow::Break(departure)) => {
2756 Self::successor_departed(creations, replaces, predecessor, result, departure)
2757 }
2758 Err((departure, input)) => Self::unexpected_stop(
2759 creations,
2760 ProxyState::ShuttingDown(ProxyShutdown::ReturningSuccessor {
2761 replaces,
2762 predecessor,
2763 departure,
2764 result,
2765 }),
2766 input,
2767 ),
2768 },
2769 ProxyShutdown::Initializing {
2770 kind,
2771 initialization,
2772 } => match initialization {
2773 WorkerInitializationShutdown::Departing {
2774 initialization,
2775 departure,
2776 } => match Self::continue_initialization_shutdown(
2777 creations,
2778 kind,
2779 Self::rejoin_initialization_departure(
2780 initialization,
2781 departure.worker_stopped(stopped),
2782 ),
2783 ) {
2784 Ok(changed) => changed,
2785 Err((creations, shutdown, input)) => {
2786 Self::unexpected_stop(creations, ProxyState::ShuttingDown(shutdown), input)
2787 }
2788 },
2789 initialization => Self::unexpected_stop(
2790 creations,
2791 ProxyState::ShuttingDown(ProxyShutdown::Initializing {
2792 kind,
2793 initialization,
2794 }),
2795 stopped,
2796 ),
2797 },
2798 ProxyShutdown::Activating { kind, activation } => {
2799 let changed = match activation {
2800 WorkerActivationShutdown::Departing {
2801 activation,
2802 departure,
2803 } => Self::rejoin_activation_departure(
2804 activation,
2805 departure.worker_stopped(stopped),
2806 ),
2807 activation => Err((activation, stopped)),
2808 };
2809 match Self::continue_activation_shutdown(creations, kind, changed) {
2810 Ok(changed) => changed,
2811 Err((creations, shutdown, input)) => {
2812 Self::unexpected_stop(creations, ProxyState::ShuttingDown(shutdown), input)
2813 }
2814 }
2815 }
2816 shutdown => {
2817 Self::unexpected_stop(creations, ProxyState::ShuttingDown(shutdown), stopped)
2818 }
2819 }
2820 }
2821
2822 fn finish_shutdown(
2823 creations: CreationSequence,
2824 worker: CurrentWorker<W>,
2825 shutdown: EstablishedShutdownResolved<W::Protocol>,
2826 stopped: ChildStopped<BehaviorAddr<W>>,
2827 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2828 Self::stop_with(
2829 creations,
2830 ProxyRetirement::Worker(StoppedWorker {
2831 worker,
2832 shutdown,
2833 stopped,
2834 }),
2835 )
2836 }
2837
2838 fn finish_start_retirement(
2839 creations: CreationSequence,
2840 kind: WorkerStartKind,
2841 retirement: WorkerStartRetirement<W, P>,
2842 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2843 match kind {
2844 WorkerStartKind::Initial => Self::stop_with(
2845 creations,
2846 ProxyRetirement::WorkerStart {
2847 replaces: None,
2848 retirement,
2849 },
2850 ),
2851 WorkerStartKind::Replacement {
2852 replaces,
2853 predecessor_shutdown: PredecessorShutdown::Settled,
2854 } => Self::stop_with(
2855 creations,
2856 ProxyRetirement::WorkerStart {
2857 replaces: Some(replaces),
2858 retirement,
2859 },
2860 ),
2861 WorkerStartKind::Replacement {
2862 replaces,
2863 predecessor_shutdown: PredecessorShutdown::Awaiting { id },
2864 } => (
2865 StableProxy {
2866 state: ProxyState::ShuttingDown(ProxyShutdown::WaitingForPredecessor {
2867 replaces,
2868 shutdown: id,
2869 retirement,
2870 }),
2871 creations,
2872 },
2873 Actions::cont(),
2874 ),
2875 }
2876 }
2877
2878 fn unexpected_worker_start(
2879 creations: CreationSequence,
2880 shutdown: ProxyShutdown<W, P>,
2881 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
2882 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2883 (
2884 StableProxy {
2885 state: ProxyState::ShuttingDown(shutdown),
2886 creations,
2887 },
2888 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerStart {
2889 phase: ProxyPhase::ShuttingDown,
2890 workers,
2891 }),
2892 )
2893 }
2894
2895 fn request_initialization(initialization: InitializeWorker<W, P>) -> ProxyActions<W, P> {
2896 Actions::send(ProxyEffects {
2897 worker_observations: InterpreterRequests::empty(),
2898 worker_initializations: InterpreterRequests::one(initialization),
2899 worker_activations: InterpreterRequests::empty(),
2900 worker_shutdowns: InterpreterRequests::empty(),
2901 worker_deliveries: Vec::new(),
2902 owner_outcomes: InterpreterRequests::empty(),
2903 diagnostics: InterpreterRequests::empty(),
2904 })
2905 }
2906
2907 fn initialize_and_shutdown(
2908 initialization: InitializeWorker<W, P>,
2909 shutdown: ShutdownEstablished<StopOnShutdown<W>, Here>,
2910 ) -> ProxyActions<W, P> {
2911 Actions::send(ProxyEffects {
2912 worker_observations: InterpreterRequests::empty(),
2913 worker_initializations: InterpreterRequests::one(initialization),
2914 worker_activations: InterpreterRequests::empty(),
2915 worker_shutdowns: InterpreterRequests::one(shutdown),
2916 worker_deliveries: Vec::new(),
2917 owner_outcomes: InterpreterRequests::empty(),
2918 diagnostics: InterpreterRequests::empty(),
2919 })
2920 }
2921
2922 fn finish_cancelled_replacement(
2923 creations: CreationSequence,
2924 replaces: WorkerAttempt,
2925 successor: WorkerAttempt,
2926 predecessor: CurrentWorker<W>,
2927 shutdown: EstablishedShutdownResolved<W::Protocol>,
2928 stopped: ChildStopped<BehaviorAddr<W>>,
2929 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2930 Self::stop_with(
2931 creations,
2932 ProxyRetirement::ReplacementCancelled {
2933 replaces,
2934 successor,
2935 predecessor: StoppedWorker {
2936 worker: predecessor,
2937 shutdown,
2938 stopped,
2939 },
2940 },
2941 )
2942 }
2943
2944 fn stop_with(
2945 creations: CreationSequence,
2946 retirement: ProxyRetirement<W, P>,
2947 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2948 (
2949 StableProxy {
2950 state: ProxyState::Stopped(retirement),
2951 creations,
2952 },
2953 Actions::stop(),
2954 )
2955 }
2956}
2957
2958impl<W, P> StableProxy<W, P>
2960where
2961 W: Behavior,
2962 P: ActivationPlan,
2963 BehaviorAddr<W>: EndpointAddress,
2964 <BehaviorAddr<W> as Address>::Nonce: Copy + Eq,
2965 StopOnShutdown<W>:
2966 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
2967 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
2968 <StopOnShutdown<W> as Behavior>::Sends: SendSettlements,
2969 <W::Birth as BirthMode>::Child: ChildCreationProduct<BehaviorAddr<W>, behavior::ChildHead>,
2970 EstablishedRecipient<W::Protocol>: Send,
2971{
2972 fn worker_activated(
2973 current: StableProxy<W, P>,
2974 input: WorkerActivation<W, P>,
2975 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
2976 match current.state {
2977 ProxyState::Starting(WorkerStart {
2978 kind,
2979 phase:
2980 WorkerStartPhase::Activating {
2981 worker,
2982 progress,
2983 stopped,
2984 },
2985 }) => match Self::admit_activation(&progress, input) {
2986 Ok(input) => match progress {
2987 ActivationProgress::WaitingForStart(activation) => Self::activation_waiting(
2988 current.creations,
2989 kind,
2990 worker,
2991 activation,
2992 stopped,
2993 input,
2994 ),
2995 ActivationProgress::Running(activation) => Self::activation_running(
2996 current.creations,
2997 kind,
2998 worker,
2999 activation,
3000 stopped,
3001 input,
3002 ),
3003 },
3004 Err(input) => Self::unexpected_activation(
3005 current.creations,
3006 ProxyState::Starting(WorkerStart {
3007 kind,
3008 phase: WorkerStartPhase::Activating {
3009 worker,
3010 progress,
3011 stopped,
3012 },
3013 }),
3014 input,
3015 ),
3016 },
3017 ProxyState::ShuttingDown(shutdown) => {
3018 Self::shutdown_worker_activated(current.creations, shutdown, input)
3019 }
3020 state => Self::unexpected_activation(current.creations, state, input),
3021 }
3022 }
3023
3024 fn activation_waiting(
3025 creations: CreationSequence,
3026 kind: WorkerStartKind,
3027 worker: CurrentWorker<W>,
3028 activation: WorkerActivationGrant<W, P>,
3029 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
3030 input: WorkerActivation<W, P>,
3031 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3032 match input.into_started() {
3033 Ok(()) => (
3034 StableProxy {
3035 state: ProxyState::Starting(WorkerStart {
3036 kind,
3037 phase: WorkerStartPhase::Activating {
3038 worker,
3039 progress: ActivationProgress::Running(activation),
3040 stopped,
3041 },
3042 }),
3043 creations,
3044 },
3045 Actions::cont(),
3046 ),
3047 Err(input) => match input.into_start_rejection() {
3048 Ok((request, reason)) => match stopped {
3049 Some(stopped) => Self::worker_start_empty(
3050 creations,
3051 kind,
3052 worker.attempt.clone(),
3053 WorkerStartResult::Unavailable {
3054 attempt: worker.attempt,
3055 drain: ProxyDrain::ActivationStartRejected {
3056 request,
3057 reason,
3058 shutdown: None,
3059 stopped,
3060 },
3061 },
3062 ),
3063 None => Self::return_worker(
3064 creations,
3065 kind,
3066 worker,
3067 PreReadyFailure::ActivationStart { request, reason },
3068 ),
3069 },
3070 Err(input) => Self::unexpected_activation(
3071 creations,
3072 ProxyState::Starting(WorkerStart {
3073 kind,
3074 phase: WorkerStartPhase::Activating {
3075 worker,
3076 progress: ActivationProgress::WaitingForStart(activation),
3077 stopped,
3078 },
3079 }),
3080 input,
3081 ),
3082 },
3083 }
3084 }
3085
3086 fn activation_running(
3087 creations: CreationSequence,
3088 kind: WorkerStartKind,
3089 worker: CurrentWorker<W>,
3090 activation: WorkerActivationGrant<W, P>,
3091 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
3092 input: WorkerActivation<W, P>,
3093 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3094 match input.into_ready() {
3095 Ok(readiness) => match stopped {
3096 Some(stopped) => Self::worker_start_empty(
3097 creations,
3098 kind,
3099 worker.attempt.clone(),
3100 WorkerStartResult::Unavailable {
3101 attempt: worker.attempt,
3102 drain: ProxyDrain::ActivationCompleted { readiness, stopped },
3103 },
3104 ),
3105 None => {
3106 let attempt = worker.attempt.clone();
3107 Self::worker_start_ready(
3108 creations,
3109 kind,
3110 worker,
3111 WorkerStartResult::Ready { attempt, readiness },
3112 )
3113 }
3114 },
3115 Err(input) => match input.into_rejection() {
3116 Ok(rejection) => match stopped {
3117 Some(stopped) => Self::worker_start_empty(
3118 creations,
3119 kind,
3120 worker.attempt.clone(),
3121 WorkerStartResult::Unavailable {
3122 attempt: worker.attempt,
3123 drain: ProxyDrain::ActivationRejected {
3124 rejection,
3125 shutdown: None,
3126 stopped,
3127 },
3128 },
3129 ),
3130 None => Self::return_worker(
3131 creations,
3132 kind,
3133 worker,
3134 PreReadyFailure::Activation { rejection },
3135 ),
3136 },
3137 Err(input) => Self::unexpected_activation(
3138 creations,
3139 ProxyState::Starting(WorkerStart {
3140 kind,
3141 phase: WorkerStartPhase::Activating {
3142 worker,
3143 progress: ActivationProgress::Running(activation),
3144 stopped,
3145 },
3146 }),
3147 input,
3148 ),
3149 },
3150 }
3151 }
3152
3153 fn unexpected_activation(
3154 creations: CreationSequence,
3155 state: ProxyState<W, P>,
3156 activation: WorkerActivation<W, P>,
3157 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3158 let phase = state.phase();
3159 (
3160 StableProxy { state, creations },
3161 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerActivation { phase, activation }),
3162 )
3163 }
3164}
3165
3166impl<W, P> StableProxy<W, P>
3168where
3169 W: Behavior,
3170 P: ActivationPlan,
3171 BehaviorAddr<W>: EndpointAddress,
3172 <BehaviorAddr<W> as Address>::Nonce: Copy + Eq,
3173 StopOnShutdown<W>:
3174 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
3175 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
3176 <StopOnShutdown<W> as Behavior>::Sends: SendSettlements,
3177 <W::Birth as BirthMode>::Child: ChildCreationProduct<BehaviorAddr<W>, behavior::ChildHead>,
3178 EstablishedRecipient<W::Protocol>: Send,
3179{
3180 fn worker_initialized(
3181 current: StableProxy<W, P>,
3182 initialization: WorkerInitializationReport<W, P>,
3183 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3184 match current.state {
3185 ProxyState::Starting(WorkerStart {
3186 kind,
3187 phase: WorkerStartPhase::Initializing { worker, stopped },
3188 }) => match worker.admit_initialization(initialization) {
3189 Ok(initialization) => Self::initialization_resolved(
3190 current.creations,
3191 kind,
3192 worker,
3193 stopped,
3194 initialization,
3195 ),
3196 Err(initialization) => Self::unexpected_start_initialization(
3197 current.creations,
3198 kind,
3199 worker,
3200 stopped,
3201 initialization,
3202 ),
3203 },
3204 ProxyState::ShuttingDown(shutdown) => {
3205 Self::shutdown_worker_initialized(current.creations, shutdown, initialization)
3206 }
3207 state => {
3208 let phase = state.phase();
3209 (
3210 StableProxy {
3211 state,
3212 creations: current.creations,
3213 },
3214 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerInitialization {
3215 phase,
3216 initialization,
3217 }),
3218 )
3219 }
3220 }
3221 }
3222
3223 fn initialization_resolved(
3224 creations: CreationSequence,
3225 kind: WorkerStartKind,
3226 worker: CurrentWorker<W>,
3227 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
3228 initialization: WorkerInitializationReport<W, P>,
3229 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3230 match initialization {
3231 WorkerInitializationReport::ReadyForActivation {
3232 activation, permit, ..
3233 } => match stopped {
3234 Some(stopped) => Self::worker_start_empty(
3235 creations,
3236 kind,
3237 worker.attempt.clone(),
3238 WorkerStartResult::Unavailable {
3239 attempt: worker.attempt,
3240 drain: ProxyDrain::InitializationCompleted {
3241 permit,
3242 activation,
3243 stopped,
3244 },
3245 },
3246 ),
3247 None => {
3248 let request = BeginActivation::new(activation, permit);
3249 let activation = request.attempt();
3250 (
3251 StableProxy {
3252 state: ProxyState::Starting(WorkerStart {
3253 kind,
3254 phase: WorkerStartPhase::Activating {
3255 worker,
3256 progress: ActivationProgress::WaitingForStart(activation),
3257 stopped: None,
3258 },
3259 }),
3260 creations,
3261 },
3262 Self::begin_activation(request),
3263 )
3264 }
3265 },
3266 WorkerInitializationReport::EffectsRejected {
3267 activation,
3268 failure,
3269 ..
3270 } => match stopped {
3271 Some(stopped) => Self::worker_start_empty(
3272 creations,
3273 kind,
3274 worker.attempt.clone(),
3275 WorkerStartResult::Unavailable {
3276 attempt: worker.attempt,
3277 drain: ProxyDrain::InitializationRejected {
3278 failure,
3279 activation,
3280 shutdown: None,
3281 stopped,
3282 },
3283 },
3284 ),
3285 None => Self::return_worker(
3286 creations,
3287 kind,
3288 worker,
3289 PreReadyFailure::InitializationEffects {
3290 failure,
3291 activation,
3292 },
3293 ),
3294 },
3295 WorkerInitializationReport::Stopped {
3296 worker: worker_attempt,
3297 initialization,
3298 activation,
3299 stopped: returned_stop,
3300 } => match worker.admit_stop(returned_stop) {
3301 Ok(returned_stop) => {
3302 let drain = ProxyDrain::InitializationStopped {
3303 activation,
3304 observed: stopped,
3305 returned: returned_stop,
3306 };
3307 Self::worker_start_empty(
3308 creations,
3309 kind,
3310 worker.attempt.clone(),
3311 WorkerStartResult::Unavailable {
3312 attempt: worker.attempt,
3313 drain,
3314 },
3315 )
3316 }
3317 Err(returned_stop) => Self::unexpected_start_initialization(
3318 creations,
3319 kind,
3320 worker,
3321 stopped,
3322 WorkerInitializationReport::Stopped {
3323 worker: worker_attempt,
3324 initialization,
3325 activation,
3326 stopped: returned_stop,
3327 },
3328 ),
3329 },
3330 }
3331 }
3332
3333 fn unexpected_start_initialization(
3334 creations: CreationSequence,
3335 kind: WorkerStartKind,
3336 worker: CurrentWorker<W>,
3337 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
3338 initialization: WorkerInitializationReport<W, P>,
3339 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3340 let phase = ProxyPhase::Initializing;
3341 (
3342 StableProxy {
3343 state: ProxyState::Starting(WorkerStart {
3344 kind,
3345 phase: WorkerStartPhase::Initializing { worker, stopped },
3346 }),
3347 creations,
3348 },
3349 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerInitialization {
3350 phase,
3351 initialization,
3352 }),
3353 )
3354 }
3355
3356 fn begin_activation(request: BeginActivation<W, P>) -> ProxyActions<W, P> {
3357 Actions::send(ProxyEffects {
3358 worker_observations: InterpreterRequests::empty(),
3359 worker_initializations: InterpreterRequests::empty(),
3360 worker_activations: InterpreterRequests::one(request),
3361 worker_shutdowns: InterpreterRequests::empty(),
3362 worker_deliveries: Vec::new(),
3363 owner_outcomes: InterpreterRequests::empty(),
3364 diagnostics: InterpreterRequests::empty(),
3365 })
3366 }
3367
3368 fn return_worker(
3369 creations: CreationSequence,
3370 kind: WorkerStartKind,
3371 worker: CurrentWorker<W>,
3372 failure: PreReadyFailure<W, P>,
3373 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3374 let (departure, shutdown) = WorkerStopping::begin(worker);
3375 (
3376 StableProxy {
3377 state: ProxyState::Starting(WorkerStart {
3378 kind,
3379 phase: WorkerStartPhase::ReturningWorker { departure, failure },
3380 }),
3381 creations,
3382 },
3383 Self::request_worker_shutdown(shutdown),
3384 )
3385 }
3386}
3387
3388impl<W, P> StableProxy<W, P>
3390where
3391 W: Behavior,
3392 P: ActivationPlan,
3393 BehaviorAddr<W>: EndpointAddress,
3394 <BehaviorAddr<W> as Address>::Nonce: Copy + Eq,
3395 StopOnShutdown<W>:
3396 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
3397 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
3398 <StopOnShutdown<W> as Behavior>::Sends: SendSettlements,
3399 <W::Birth as BirthMode>::Child: ChildCreationProduct<BehaviorAddr<W>, behavior::ChildHead>,
3400 EstablishedRecipient<W::Protocol>: Send,
3401{
3402 fn worker_shutdown_resolved(
3403 current: StableProxy<W, P>,
3404 shutdown: EstablishedShutdownResolved<W::Protocol>,
3405 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3406 match current.state {
3407 ProxyState::Starting(WorkerStart {
3408 kind:
3409 WorkerStartKind::Replacement {
3410 replaces,
3411 predecessor_shutdown: PredecessorShutdown::Awaiting { id },
3412 },
3413 phase,
3414 }) if shutdown.id() == id => (
3415 StableProxy {
3416 state: ProxyState::Starting(WorkerStart {
3417 kind: WorkerStartKind::Replacement {
3418 replaces,
3419 predecessor_shutdown: PredecessorShutdown::Settled,
3420 },
3421 phase,
3422 }),
3423 creations: current.creations,
3424 },
3425 Actions::cont(),
3426 ),
3427 ProxyState::Starting(WorkerStart { kind, phase }) => {
3428 Self::worker_start_shutdown_resolved(current.creations, kind, phase, shutdown)
3429 }
3430 ProxyState::Replacing(replacement) => {
3431 Self::replacement_shutdown_resolved(current.creations, replacement, shutdown)
3432 }
3433 ProxyState::ShuttingDown(proxy_shutdown) => {
3434 Self::shutdown_resolved(current.creations, proxy_shutdown, shutdown)
3435 }
3436 state => Self::unexpected_shutdown(current.creations, state, shutdown),
3437 }
3438 }
3439
3440 fn worker_start_shutdown_resolved(
3441 creations: CreationSequence,
3442 kind: WorkerStartKind,
3443 phase: WorkerStartPhase<W, P>,
3444 shutdown: EstablishedShutdownResolved<W::Protocol>,
3445 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3446 match phase {
3447 WorkerStartPhase::ReturningWorker { departure, failure } => {
3448 match departure.shutdown_resolved(shutdown) {
3449 Ok(ControlFlow::Continue(departure)) => (
3450 StableProxy {
3451 state: ProxyState::Starting(WorkerStart {
3452 kind,
3453 phase: WorkerStartPhase::ReturningWorker { departure, failure },
3454 }),
3455 creations,
3456 },
3457 Actions::cont(),
3458 ),
3459 Ok(ControlFlow::Break(StoppedWorker {
3460 worker,
3461 shutdown,
3462 stopped,
3463 })) => {
3464 Self::worker_returned(creations, kind, worker, failure, shutdown, stopped)
3465 }
3466 Err((departure, input)) => Self::unexpected_shutdown(
3467 creations,
3468 ProxyState::Starting(WorkerStart {
3469 kind,
3470 phase: WorkerStartPhase::ReturningWorker { departure, failure },
3471 }),
3472 input,
3473 ),
3474 }
3475 }
3476 phase => Self::unexpected_shutdown(
3477 creations,
3478 ProxyState::Starting(WorkerStart { kind, phase }),
3479 shutdown,
3480 ),
3481 }
3482 }
3483
3484 fn unexpected_shutdown(
3485 creations: CreationSequence,
3486 state: ProxyState<W, P>,
3487 shutdown: EstablishedShutdownResolved<W::Protocol>,
3488 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3489 let phase = state.phase();
3490 (
3491 StableProxy { state, creations },
3492 Self::diagnose(ProxyDiagnostic::UnexpectedWorkerShutdown { phase, shutdown }),
3493 )
3494 }
3495
3496 fn worker_returned(
3497 creations: CreationSequence,
3498 kind: WorkerStartKind,
3499 worker: CurrentWorker<W>,
3500 failure: PreReadyFailure<W, P>,
3501 shutdown: EstablishedShutdownResolved<W::Protocol>,
3502 stopped: ChildStopped<BehaviorAddr<W>>,
3503 ) -> (StableProxy<W, P>, ProxyActions<W, P>) {
3504 let worker_attempt = worker.attempt.clone();
3505 let result = Self::worker_return_result(worker, failure, shutdown, stopped);
3506 Self::worker_start_empty(creations, kind, worker_attempt, result)
3507 }
3508
3509 fn worker_return_result(
3510 worker: CurrentWorker<W>,
3511 failure: PreReadyFailure<W, P>,
3512 shutdown: EstablishedShutdownResolved<W::Protocol>,
3513 stopped: ChildStopped<BehaviorAddr<W>>,
3514 ) -> WorkerStartResult<W, P> {
3515 let worker_attempt = worker.attempt;
3516 let result = match failure {
3517 PreReadyFailure::InitializationEffects {
3518 failure,
3519 activation,
3520 } => WorkerStartResult::Unavailable {
3521 attempt: worker_attempt,
3522 drain: ProxyDrain::InitializationRejected {
3523 failure,
3524 activation,
3525 shutdown: Some(shutdown),
3526 stopped,
3527 },
3528 },
3529 PreReadyFailure::ActivationStart { request, reason } => {
3530 WorkerStartResult::Unavailable {
3531 attempt: worker_attempt,
3532 drain: ProxyDrain::ActivationStartRejected {
3533 request,
3534 reason,
3535 shutdown: Some(shutdown),
3536 stopped,
3537 },
3538 }
3539 }
3540 PreReadyFailure::Activation { rejection } => WorkerStartResult::Unavailable {
3541 attempt: worker_attempt,
3542 drain: ProxyDrain::ActivationRejected {
3543 rejection,
3544 shutdown: Some(shutdown),
3545 stopped,
3546 },
3547 },
3548 };
3549 result
3550 }
3551}