1use core::convert::Infallible;
4use core::mem;
5use core::ops::ControlFlow;
6
7use crate::{
8 DeliveryRoute, DiagnosticAction, DiagnosticDisposition, ObserveChild, ReplyRoute, ScheduleAfter,
9};
10use behavior::{
11 ActionItemResult, Actions, ActiveTurn, Address, Behavior, BehaviorActed, BehaviorAddr,
12 BehaviorBase, Births, ChildHead, CreationId, CreationSequence, Creations, CreationsSettled,
13 EndpointAddress, EstablishedActor, InitializationTurn, InterpreterRequests, MessageProtocol,
14 Never, SendEffects, SourceActions, Step, User,
15};
16
17use super::worker::{InitialWorkerRejection, prepare_initial_workers};
18use super::{
19 ActivationPlan, ActivationPolicy, ActorDrainPolicy, OrderedRoles, PreparedWorker,
20 ProxyOperation, RestartLimit, RestartRelease, RoleName, StableProxy, WorkerSubmission,
21};
22
23mod diagnostic;
24mod event;
25mod lifecycle;
26mod member;
27mod protocol;
28mod proxy;
29mod recovery;
30mod requests;
31mod restart;
32mod role;
33mod shutdown;
34
35pub use super::worker::{
36 PendingWorkerPreparation, PrepareWorkers, WorkerPreparation, WorkerSource,
37};
38pub use diagnostic::{
39 FixedDiagnostic, ProxyInputFailure, ProxyOutcomeFailure, RecoveryDenied,
40 RestartScheduleFailure, WorkerPreparationFailure, WorkerPreparationFailureReason,
41 WorkerUnavailable,
42};
43pub use event::FixedSupervisorEvent;
44pub use lifecycle::{FixedLifecycle, FixedLifecycleEvent, FixedLifecycleRoute};
45pub use protocol::{CapabilityResult, FixedCommand, FixedSnapshot, MemberStatus, UnavailablePhase};
46pub use requests::FixedSupervisorRequests;
47pub use restart::RecoveryDenialReason;
48
49use member::{OnlineMember, RosterOwner};
50use proxy::{FixedRoster, InitialProxyDecision};
51use recovery::{
52 AcceptedWorkerStop, PreparationAcceptance, PreparedRecoveryDecision, RecoveryChoice,
53 RecoveryFailure, RecoveryTicket, ReplacementCompletion, RestartScheduleAdmission,
54 SupervisorRecovery, UnrecoveredMember, WorkerPreparationRetirement, WorkerStopDecision,
55};
56use shutdown::FixedShutdown;
57
58#[derive(Clone, Copy, Debug, Eq, PartialEq)]
60pub enum Strategy {
61 OneForOne,
63 OneForAll,
65 RestForOne,
67}
68
69#[derive(Clone, Copy, Debug, Eq, PartialEq)]
70enum RecoveryDecision<Source> {
71 Permanent {
72 source: Source,
73 strategy: Strategy,
74 limit: RestartLimit,
75 release: RestartRelease,
76 },
77 Transient {
78 source: Source,
79 strategy: Strategy,
80 limit: RestartLimit,
81 release: RestartRelease,
82 },
83 Temporary,
84}
85
86#[derive(Clone, Copy, Debug, Eq, PartialEq)]
88pub struct Recovery<Source> {
89 decision: RecoveryDecision<Source>,
90}
91
92impl<Source> Recovery<Source> {
93 #[must_use]
95 pub const fn permanent(
96 source: Source,
97 strategy: Strategy,
98 limit: RestartLimit,
99 release: RestartRelease,
100 ) -> Self {
101 Self {
102 decision: RecoveryDecision::Permanent {
103 source,
104 strategy,
105 limit,
106 release,
107 },
108 }
109 }
110
111 #[must_use]
113 pub const fn transient(
114 source: Source,
115 strategy: Strategy,
116 limit: RestartLimit,
117 release: RestartRelease,
118 ) -> Self {
119 Self {
120 decision: RecoveryDecision::Transient {
121 source,
122 strategy,
123 limit,
124 release,
125 },
126 }
127 }
128}
129
130impl Recovery<behavior::Never> {
131 #[must_use]
133 pub const fn temporary() -> Self {
134 Self {
135 decision: RecoveryDecision::Temporary,
136 }
137 }
138}
139
140#[derive(Clone, Copy, Debug, Eq, PartialEq)]
142pub enum FailureReaction {
143 RetireMember,
145 StopSupervisor,
147}
148
149pub struct FixedBuilder<Factory, Role, Source, DiagnosticRoute, LifecycleRoute> {
151 factory: Factory,
152 roles: OrderedRoles<Role>,
153 activation: ActivationPolicy,
154 recovery: Recovery<Source>,
155 failure_reaction: FailureReaction,
156 actor_drain: ActorDrainPolicy,
157 diagnostics: DiagnosticDisposition<DiagnosticRoute>,
158 lifecycle: Option<LifecycleRoute>,
159}
160
161pub struct FixedSupervisor<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute>
163where
164 Role: Send + Sync,
165 Source: WorkerSource<Role, Worker, Plan>,
166 Worker: Behavior + Send,
167 Plan: ActivationPlan,
168 BehaviorAddr<Worker>: EndpointAddress,
169 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
170{
171 roster: FixedRoster<Role, Worker, Plan>,
172 creations: CreationSequence,
173 activation: ActivationPolicy,
174 recovery: SupervisorRecovery<Source, WorkerPreparationRetirement<Source, Role, Worker, Plan>>,
175 failure_reaction: FailureReaction,
176 actor_drain: ActorDrainPolicy,
177 diagnostics: DiagnosticDisposition<DiagnosticRoute>,
178 lifecycle: Option<LifecycleRoute>,
179}
180
181pub enum FixedSupervisorError<Role, Worker, Plan, PreparationStart, PreparationReturn>
187where
188 Worker: Behavior,
189 Plan: ActivationPlan,
190 BehaviorAddr<Worker>: EndpointAddress,
191 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
192{
193 InitializationUnavailable,
195 ProxyCreationsExhausted {
197 members: Vec<PreparedWorker<Role, Worker, Plan>>,
199 },
200 ProxyCreationRejected {
202 proxies: CreationsSettled<BehaviorAddr<Worker>, StableProxy<Worker, Plan>>,
204 },
205 InputRejected {
207 input: FixedSupervisorEvent<Role, Worker, Plan, PreparationStart, PreparationReturn>,
209 },
210 RosterStateRejected {
215 members: Vec<PreparedWorker<Role, Worker, Plan>>,
217 },
218 AuthorizationStateRejected {
220 occupied: usize,
222 maximum: usize,
224 },
225 OperatingStateRejected,
227 InitialOutcomeContradiction {
229 report: behavior::ChildReport<super::ProxyOutcome<Worker, Plan>>,
231 },
232}
233
234pub struct FixedConstructionRejected<
236 Factory,
237 Role,
238 Worker,
239 Plan,
240 Rejection,
241 Source,
242 DiagnosticRoute,
243 LifecycleRoute,
244> {
245 pub workers: InitialWorkerRejection<Factory, Role, Worker, Plan, Rejection>,
247 pub activation: ActivationPolicy,
249 pub recovery: Recovery<Source>,
251 pub failure_reaction: FailureReaction,
253 pub actor_drain: ActorDrainPolicy,
255 pub diagnostics: DiagnosticDisposition<DiagnosticRoute>,
257 pub lifecycle: Option<LifecycleRoute>,
259}
260
261#[must_use]
263pub fn fixed<Factory, Role, Source, DiagnosticRoute>(
264 factory: Factory,
265 roles: OrderedRoles<Role>,
266 activation: ActivationPolicy,
267 recovery: Recovery<Source>,
268 failure_reaction: FailureReaction,
269 actor_drain: ActorDrainPolicy,
270 diagnostics: DiagnosticDisposition<DiagnosticRoute>,
271) -> FixedBuilder<Factory, Role, Source, DiagnosticRoute, Infallible> {
272 FixedBuilder {
273 factory,
274 roles,
275 activation,
276 recovery,
277 failure_reaction,
278 actor_drain,
279 diagnostics,
280 lifecycle: None,
281 }
282}
283
284impl<Factory, Role, Source, DiagnosticRoute>
285 FixedBuilder<Factory, Role, Source, DiagnosticRoute, Infallible>
286{
287 #[must_use]
289 pub fn publish_lifecycle<Route>(
290 self,
291 route: Route,
292 ) -> FixedBuilder<Factory, Role, Source, DiagnosticRoute, Route> {
293 FixedBuilder {
294 factory: self.factory,
295 roles: self.roles,
296 activation: self.activation,
297 recovery: self.recovery,
298 failure_reaction: self.failure_reaction,
299 actor_drain: self.actor_drain,
300 diagnostics: self.diagnostics,
301 lifecycle: Some(route),
302 }
303 }
304}
305
306impl<Factory, Role, Source, DiagnosticRoute, LifecycleRoute>
307 FixedBuilder<Factory, Role, Source, DiagnosticRoute, LifecycleRoute>
308{
309 pub fn build<Worker, Plan, Rejection>(
311 self,
312 ) -> Result<
313 FixedSupervisor<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute>,
314 FixedConstructionRejected<
315 Factory,
316 Role,
317 Worker,
318 Plan,
319 Rejection,
320 Source,
321 DiagnosticRoute,
322 LifecycleRoute,
323 >,
324 >
325 where
326 Factory: FnMut(&Role) -> Result<WorkerSubmission<Worker, Plan>, Rejection>,
327 Role: Send + Sync,
328 Source: WorkerSource<Role, Worker, Plan>,
329 Worker: behavior::Behavior + Send,
330 Plan: ActivationPlan,
331 BehaviorAddr<Worker>: EndpointAddress,
332 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
333 {
334 let Self {
335 factory,
336 roles,
337 activation,
338 recovery,
339 failure_reaction,
340 actor_drain,
341 diagnostics,
342 lifecycle,
343 } = self;
344 let prepared = match prepare_initial_workers(factory, roles) {
345 Ok(prepared) => prepared,
346 Err(workers) => {
347 return Err(FixedConstructionRejected {
348 workers,
349 activation,
350 recovery,
351 failure_reaction,
352 actor_drain,
353 diagnostics,
354 lifecycle,
355 });
356 }
357 };
358 Ok(FixedSupervisor {
359 roster: FixedRoster::new(prepared),
360 creations: CreationSequence::new(),
361 activation,
362 recovery: recovery.into(),
363 failure_reaction,
364 actor_drain,
365 diagnostics,
366 lifecycle,
367 })
368 }
369}
370
371impl<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute> BehaviorBase
372 for FixedSupervisor<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute>
373where
374 Role: Send + Sync,
375 Worker: Behavior + Send,
376 Plan: ActivationPlan,
377 Source: WorkerSource<Role, Worker, Plan>,
378 BehaviorAddr<Worker>: EndpointAddress,
379 <BehaviorAddr<Worker> as Address>::Nonce: Send,
380 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
381{
382 type Base = Self;
383
384 fn base(&self) -> &Self::Base {
385 self
386 }
387}
388
389impl<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute>
390 FixedSupervisor<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute>
391where
392 Role: Eq + Send + Sync,
393 Worker: Behavior + Send,
394 Plan: ActivationPlan,
395 Source: WorkerSource<Role, Worker, Plan>,
396 BehaviorAddr<Worker>: EndpointAddress,
397 <BehaviorAddr<Worker> as Address>::Nonce: Send,
398 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
399 EstablishedActor<StableProxy<Worker, Plan>>: Clone + Send,
400 DiagnosticRoute: crate::DiagnosticRoute<FixedDiagnostic<Role, Worker, Plan, Source>> + Clone,
401 LifecycleRoute: FixedLifecycleRoute<Role, Worker, Plan> + Clone,
402 LifecycleRoute::Sends: behavior::SendsFor<
403 FixedSupervisorEvent<
404 Role,
405 Worker,
406 Plan,
407 ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
408 WorkerPreparation<Source, Role, Worker, Plan>,
409 >,
410 >,
411{
412 fn authorize_roster(
413 &self,
414 roster: FixedRoster<Role, Worker, Plan>,
415 ) -> Result<
416 (
417 FixedRoster<Role, Worker, Plan>,
418 SourceActions<ProxyOperation<behavior::Here, Worker, Plan>>,
419 ),
420 (
421 FixedRoster<Role, Worker, Plan>,
422 FixedSupervisorError<
423 Role,
424 Worker,
425 Plan,
426 ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
427 WorkerPreparation<Source, Role, Worker, Plan>,
428 >,
429 ),
430 > {
431 let (roster, operations) = match roster.authorize(self.activation.maximum()) {
432 Ok(authorized) => authorized,
433 Err(rejected) => return Err(rejected),
434 };
435 let mut proxy_operations = SourceActions::empty();
436 for operation in operations {
437 proxy_operations.send(operation);
438 }
439 Ok((roster, proxy_operations))
440 }
441
442 fn authorize_waiting(
443 &mut self,
444 roster: FixedRoster<Role, Worker, Plan>,
445 ) -> BehaviorActed<Self> {
446 let (next, proxy_operations) = match self.authorize_roster(roster) {
447 Ok(authorized) => authorized,
448 Err((next, rejection)) => {
449 self.roster = next;
450 return Err(rejection);
451 }
452 };
453 self.roster = next;
454 let mut requests = FixedSupervisorRequests::empty();
455 requests.proxy_operations = proxy_operations;
456 Ok(Actions::send(requests))
457 }
458
459 fn member_retired(
460 &mut self,
461 roster: FixedRoster<Role, Worker, Plan>,
462 role: RoleName<Role>,
463 ) -> BehaviorActed<Self> {
464 let mut actions = match self.authorize_waiting(roster) {
465 Ok(actions) => actions,
466 Err(rejection) => return Err(rejection),
467 };
468 match self.lifecycle.clone() {
469 Some(route) => {
470 actions.sends.lifecycle = route.deliver(FixedLifecycle::member_retired(role));
471 }
472 None => {}
473 }
474 Ok(actions)
475 }
476
477 fn admit_recovery(
478 &mut self,
479 roster: FixedRoster<Role, Worker, Plan>,
480 recovery: SupervisorRecovery<
481 Source,
482 WorkerPreparationRetirement<Source, Role, Worker, Plan>,
483 >,
484 schedule: Option<ScheduleAfter>,
485 ) -> BehaviorActed<Self> {
486 let (roster, operations) = match roster.authorize(self.activation.maximum()) {
487 Ok(authorized) => authorized,
488 Err((roster, rejection)) => {
489 self.roster = roster;
490 self.recovery = recovery;
491 return Err(rejection);
492 }
493 };
494 self.roster = roster;
495 self.recovery = recovery;
496 let mut proxy_operations = SourceActions::empty();
497 for operation in operations {
498 proxy_operations.send(operation);
499 }
500 let mut restart_schedules = SourceActions::empty();
501 match schedule {
502 Some(schedule) => restart_schedules.send(schedule),
503 None => {}
504 }
505 let mut requests = FixedSupervisorRequests::empty();
506 requests.proxy_operations = proxy_operations;
507 requests.restart_schedules = restart_schedules;
508 Ok(Actions::send(requests))
509 }
510
511 fn accept_restart_schedule(
512 &mut self,
513 roster: FixedRoster<Role, Worker, Plan>,
514 input: ActionItemResult<ScheduleAfter>,
515 ) -> BehaviorActed<Self> {
516 let roster = match roster {
517 FixedRoster::ShuttingDown(shutdown) => {
518 return match shutdown.accept_schedule(input) {
519 Ok(shutdown) => self.retain_shutdown(shutdown, Vec::new()),
520 Err((shutdown, input)) => {
521 self.roster = FixedRoster::ShuttingDown(shutdown);
522 Err(FixedSupervisorError::InputRejected {
523 input: FixedSupervisorEvent::RestartScheduleSettled(input),
524 })
525 }
526 };
527 }
528 roster => roster,
529 };
530 match roster.accept_restart_schedule(input) {
531 RestartScheduleAdmission::Waiting(roster) => self.authorize_waiting(roster),
532 RestartScheduleAdmission::Rejected(rejected) => {
533 let (failed, diagnostic) = rejected.into_parts();
534 self.apply_recovery_failure(
535 failed,
536 FixedDiagnostic::RestartScheduleFailed(diagnostic),
537 )
538 }
539 RestartScheduleAdmission::Unrelated { roster, input } => {
540 self.roster = roster;
541 Err(FixedSupervisorError::InputRejected {
542 input: FixedSupervisorEvent::RestartScheduleSettled(input),
543 })
544 }
545 }
546 }
547
548 fn apply_recovery_failure(
549 &mut self,
550 failed: RecoveryFailure<Role, Worker, Plan>,
551 diagnostic: FixedDiagnostic<Role, Worker, Plan, Source>,
552 ) -> BehaviorActed<Self> {
553 match (&self.diagnostics, self.failure_reaction) {
554 (DiagnosticDisposition::Terminate, _) => {
555 self.roster = failed.terminate();
556 let mut requests = FixedSupervisorRequests::empty();
557 requests.diagnostics =
558 InterpreterRequests::one(DiagnosticAction::terminal(diagnostic));
559 Ok(Actions::new(
560 requests,
561 Creations::empty(),
562 Step::Stop(behavior::Stopped),
563 ))
564 }
565 (DiagnosticDisposition::DeliverTo(_), FailureReaction::RetireMember) => {
566 let (roster, operation) = failed.retire_member();
567 self.roster = roster;
568 let mut proxy_operations = SourceActions::empty();
569 proxy_operations.send(operation);
570 let mut requests = FixedSupervisorRequests::empty();
571 requests.proxy_operations = proxy_operations;
572 requests.diagnostics =
573 InterpreterRequests::one(self.diagnostics.action(diagnostic));
574 Ok(Actions::send(requests))
575 }
576 (DiagnosticDisposition::DeliverTo(_), FailureReaction::StopSupervisor) => {
577 let (shutdown, operations, schedule) = failed.stop_supervisor(self.actor_drain);
578 self.roster = FixedRoster::ShuttingDown(shutdown);
579 let mut proxy_operations = SourceActions::empty();
580 for operation in operations {
581 proxy_operations.send(operation);
582 }
583 let mut restart_schedules = SourceActions::empty();
584 match schedule {
585 Some(schedule) => restart_schedules.send(schedule),
586 None => {}
587 }
588 let mut requests = FixedSupervisorRequests::empty();
589 requests.proxy_operations = proxy_operations;
590 requests.restart_schedules = restart_schedules;
591 requests.diagnostics =
592 InterpreterRequests::one(self.diagnostics.action(diagnostic));
593 Ok(Actions::send(requests))
594 }
595 }
596 }
597
598 fn accept_restart_elapsed(
599 &mut self,
600 roster: FixedRoster<Role, Worker, Plan>,
601 elapsed: crate::TimerElapsed,
602 ) -> BehaviorActed<Self> {
603 let roster = match roster {
604 FixedRoster::ShuttingDown(shutdown) => {
605 return match shutdown.accept_elapsed(elapsed) {
606 Ok(shutdown) => self.retain_shutdown(shutdown, Vec::new()),
607 Err((shutdown, elapsed)) => {
608 self.roster = FixedRoster::ShuttingDown(shutdown);
609 Err(FixedSupervisorError::InputRejected {
610 input: FixedSupervisorEvent::RestartElapsed(elapsed),
611 })
612 }
613 };
614 }
615 roster => roster,
616 };
617 match roster.accept_restart_elapsed(elapsed) {
618 Ok(roster) => self.authorize_waiting(roster),
619 Err((roster, elapsed)) => {
620 self.roster = roster;
621 Err(FixedSupervisorError::InputRejected {
622 input: FixedSupervisorEvent::RestartElapsed(elapsed),
623 })
624 }
625 }
626 }
627
628 fn retain_shutdown(
629 &mut self,
630 shutdown: FixedShutdown<Role, Worker, Plan>,
631 operations: Vec<ProxyOperation<behavior::Here, Worker, Plan>>,
632 ) -> BehaviorActed<Self> {
633 let (shutdown, become_) = shutdown.into_step(self.recovery.retirement());
634 self.roster = FixedRoster::ShuttingDown(shutdown);
635 let mut proxy_operations = SourceActions::empty();
636 for operation in operations {
637 proxy_operations.send(operation);
638 }
639 let mut requests = FixedSupervisorRequests::empty();
640 requests.proxy_operations = proxy_operations;
641 Ok(Actions::new(requests, Creations::empty(), become_))
642 }
643
644 fn accept_proxy_birth(
645 &mut self,
646 roster: FixedRoster<Role, Worker, Plan>,
647 proxies: CreationsSettled<BehaviorAddr<Worker>, StableProxy<Worker, Plan>>,
648 ) -> BehaviorActed<Self> {
649 match roster {
650 FixedRoster::ShuttingDown(shutdown) => match shutdown.accept_births(proxies) {
651 Ok((shutdown, operations)) => self.retain_shutdown(shutdown, operations),
652 Err((shutdown, proxies)) => {
653 self.roster = FixedRoster::ShuttingDown(shutdown);
654 Err(FixedSupervisorError::InputRejected {
655 input: FixedSupervisorEvent::ProxyCreationsSettled(proxies),
656 })
657 }
658 },
659 roster => {
660 let next = match roster.accept_births(proxies) {
661 Ok(next) => next,
662 Err((next, rejection)) => {
663 self.roster = next;
664 return Err(rejection);
665 }
666 };
667 self.authorize_waiting(next)
668 }
669 }
670 }
671
672 fn accept_proxy_operation(
673 &mut self,
674 roster: FixedRoster<Role, Worker, Plan>,
675 settlement: crate::ProxyInputResult<behavior::Here, Worker, Plan>,
676 ) -> BehaviorActed<Self> {
677 match roster {
678 FixedRoster::ShuttingDown(shutdown) => match shutdown.accept_operation(settlement) {
679 Ok(shutdown) => self.retain_shutdown(shutdown, Vec::new()),
680 Err((shutdown, settlement)) => {
681 self.roster = FixedRoster::ShuttingDown(shutdown);
682 Err(FixedSupervisorError::InputRejected {
683 input: FixedSupervisorEvent::ProxyInputSettled(settlement),
684 })
685 }
686 },
687 roster => {
688 let (roster, settlement) = match roster.accept_replacement_operation(settlement) {
689 Ok(ControlFlow::Continue(roster)) => return self.authorize_waiting(roster),
690 Ok(ControlFlow::Break((owners, completed))) => {
691 return self.accept_replacement_completion(owners, completed);
692 }
693 Err(returned) => returned,
694 };
695 let (next, retired) = match roster.accept_operation(settlement) {
696 Ok(next) => next,
697 Err((next, rejection)) => {
698 self.roster = next;
699 return Err(rejection);
700 }
701 };
702 match retired {
703 Some(role) => self.member_retired(next, role),
704 None => self.authorize_waiting(next),
705 }
706 }
707 }
708 }
709
710 fn accept_proxy_report(
711 &mut self,
712 roster: FixedRoster<Role, Worker, Plan>,
713 report: behavior::ChildReport<super::ProxyOutcome<Worker, Plan>>,
714 ) -> BehaviorActed<Self> {
715 let behavior::ChildReport {
716 child,
717 report: outcome,
718 } = report;
719 let report = match outcome {
720 super::ProxyOutcome::Unavailable {
721 sender,
722 phase,
723 command,
724 } => {
725 return self.accept_unavailable(roster, child, sender, phase, command);
726 }
727 outcome => behavior::ChildReport::new(child, outcome),
728 };
729 match roster {
730 FixedRoster::ShuttingDown(shutdown) => match shutdown.accept_outcome(report) {
731 Ok(shutdown) => self.retain_shutdown(shutdown, Vec::new()),
732 Err((shutdown, report)) => {
733 self.roster = FixedRoster::ShuttingDown(shutdown);
734 Err(FixedSupervisorError::InputRejected {
735 input: FixedSupervisorEvent::ProxyReported(report),
736 })
737 }
738 },
739 roster => {
740 let (roster, report) = match roster.accept_recovery_report(report) {
741 Ok(ControlFlow::Continue(roster)) => return self.authorize_waiting(roster),
742 Ok(ControlFlow::Break((owners, completed))) => {
743 return self.accept_replacement_completion(owners, completed);
744 }
745 Err(returned) => returned,
746 };
747 match report.report {
748 outcome @ super::ProxyOutcome::WorkerStopped { .. } => {
749 let report = behavior::ChildReport::new(report.child, outcome);
750 match roster.accept_worker_stop(report) {
751 WorkerStopDecision::Accepted(stopped) => {
752 self.accept_worker_stop(stopped)
753 }
754 WorkerStopDecision::Rejected { roster, report } => {
755 self.roster = roster;
756 Err(FixedSupervisorError::InputRejected {
757 input: FixedSupervisorEvent::ProxyReported(report),
758 })
759 }
760 }
761 }
762 outcome => {
763 let report = behavior::ChildReport::new(report.child, outcome);
764 match roster.accept_initial_outcome(report) {
765 InitialProxyDecision::Ready { mut owners, member } => {
766 let lifecycle = match self.lifecycle.clone() {
767 Some(route_to_lifecycle) => {
768 route_to_lifecycle.deliver(FixedLifecycle::started(
769 member.role.name(),
770 member.proxy.clone(),
771 ))
772 }
773 None => LifecycleRoute::Sends::empty(),
774 };
775 owners.push(RosterOwner::Online(member));
776 let (roster, proxy_operations) =
777 match self.authorize_roster(FixedRoster::Operating(owners)) {
778 Ok(authorized) => authorized,
779 Err((roster, rejection)) => {
780 self.roster = roster;
781 return Err(rejection);
782 }
783 };
784 self.roster = roster;
785 let mut requests = FixedSupervisorRequests::empty();
786 requests.proxy_operations = proxy_operations;
787 requests.lifecycle = lifecycle;
788 Ok(Actions::send(requests))
789 }
790 InitialProxyDecision::Failed {
791 owners,
792 child,
793 role,
794 outcome,
795 } => self.accept_initial_failure(owners, child, role, outcome),
796 InitialProxyDecision::Rejected { roster, rejection } => {
797 self.roster = roster;
798 Err(rejection)
799 }
800 }
801 }
802 }
803 }
804 }
805 }
806
807 fn accept_unavailable(
808 &mut self,
809 roster: FixedRoster<Role, Worker, Plan>,
810 child: CreationId,
811 sender: BehaviorAddr<Worker>,
812 phase: super::ProxyPhase,
813 command: <Worker::Protocol as behavior::Protocol>::Msg,
814 ) -> BehaviorActed<Self> {
815 let role = match roster.live_proxy_role(child) {
816 Some(role) => role,
817 None => {
818 self.roster = roster;
819 return Err(FixedSupervisorError::InputRejected {
820 input: FixedSupervisorEvent::ProxyReported(behavior::ChildReport::new(
821 child,
822 super::ProxyOutcome::Unavailable {
823 sender,
824 phase,
825 command,
826 },
827 )),
828 });
829 }
830 };
831 let unavailable = WorkerUnavailable::new(role, sender, phase, command);
832 self.roster = roster;
833 match self.lifecycle.clone() {
834 Some(route_to_lifecycle) => {
835 let mut requests = FixedSupervisorRequests::empty();
836 requests.lifecycle =
837 route_to_lifecycle.deliver(FixedLifecycle::unavailable(unavailable));
838 Ok(Actions::send(requests))
839 }
840 None => {
841 let become_ = match &self.diagnostics {
842 DiagnosticDisposition::DeliverTo(_) => Step::Continue,
843 DiagnosticDisposition::Terminate => Step::Stop(behavior::Stopped),
844 };
845 let diagnostic = FixedDiagnostic::WorkerUnavailable(unavailable);
846 let mut requests = FixedSupervisorRequests::empty();
847 requests.diagnostics =
848 InterpreterRequests::one(self.diagnostics.action(diagnostic));
849 Ok(Actions::new(requests, Creations::empty(), become_))
850 }
851 }
852 }
853
854 fn accept_replacement_completion(
855 &mut self,
856 mut owners: Vec<RosterOwner<Role, Worker, Plan>>,
857 completed: ReplacementCompletion<Role, Worker, Plan>,
858 ) -> BehaviorActed<Self> {
859 match completed {
860 ReplacementCompletion::Restarted {
861 role,
862 creation,
863 proxy,
864 previous_readiness,
865 stopped,
866 worker,
867 readiness,
868 recovery,
869 } => {
870 let name = role.name();
871 let lifecycle =
872 match self.lifecycle.clone() {
873 Some(route_to_lifecycle) => {
874 let mut lifecycle = LifecycleRoute::Sends::empty();
875 lifecycle.append(route_to_lifecycle.clone().deliver(
876 FixedLifecycle::worker_stopped_after_admission(
877 name.clone(),
878 stopped,
879 recovery.0,
880 ),
881 ));
882 lifecycle.append(route_to_lifecycle.deliver(
883 FixedLifecycle::restarted(name, proxy.clone(), recovery.0),
884 ));
885 lifecycle
886 }
887 None => {
888 let _retired_predecessor = stopped;
889 LifecycleRoute::Sends::empty()
890 }
891 };
892 drop(previous_readiness);
893 owners.push(RosterOwner::Online(OnlineMember {
894 role,
895 creation,
896 proxy,
897 worker,
898 readiness,
899 }));
900 let (roster, proxy_operations) =
901 match self.authorize_roster(FixedRoster::Operating(owners)) {
902 Ok(authorized) => authorized,
903 Err((roster, rejection)) => {
904 self.roster = roster;
905 return Err(rejection);
906 }
907 };
908 self.roster = roster;
909 let mut requests = FixedSupervisorRequests::empty();
910 requests.proxy_operations = proxy_operations;
911 requests.lifecycle = lifecycle;
912 Ok(Actions::send(requests))
913 }
914 ReplacementCompletion::InputRejected {
915 role,
916 creation,
917 proxy,
918 previous,
919 previous_readiness,
920 stopped,
921 operation,
922 reason,
923 recovery,
924 } => {
925 let diagnostic = FixedDiagnostic::ProxyInputRejected(ProxyInputFailure::new(
926 role.name(),
927 operation,
928 reason,
929 ));
930 self.apply_replacement_failure(
931 owners,
932 OnlineMember {
933 role,
934 creation,
935 proxy,
936 worker: previous,
937 readiness: previous_readiness,
938 },
939 stopped,
940 recovery,
941 diagnostic,
942 )
943 }
944 ReplacementCompletion::ProxyFailed {
945 role,
946 creation,
947 proxy,
948 previous,
949 previous_readiness,
950 stopped,
951 outcome,
952 recovery,
953 } => {
954 let diagnostic = FixedDiagnostic::ProxyOutcomeFailed(ProxyOutcomeFailure::new(
955 role.name(),
956 super::ProxyOutcome::Replacement { outcome },
957 ));
958 self.apply_replacement_failure(
959 owners,
960 OnlineMember {
961 role,
962 creation,
963 proxy,
964 worker: previous,
965 readiness: previous_readiness,
966 },
967 Some(stopped),
968 recovery,
969 diagnostic,
970 )
971 }
972 }
973 }
974
975 fn apply_replacement_failure(
976 &mut self,
977 mut owners: Vec<RosterOwner<Role, Worker, Plan>>,
978 previous: OnlineMember<Role, Worker, Plan>,
979 stopped: Option<crate::ChildStopped<BehaviorAddr<Worker>>>,
980 recovery: RecoveryTicket,
981 diagnostic: FixedDiagnostic<Role, Worker, Plan, Source>,
982 ) -> BehaviorActed<Self> {
983 let OnlineMember {
984 role,
985 creation,
986 proxy,
987 worker,
988 readiness,
989 } = previous;
990 let position = role.position();
991 let (member, operation, lifecycle) = match (stopped, self.lifecycle.clone()) {
992 (Some(stopped), Some(route_to_lifecycle)) => {
993 let lifecycle =
994 route_to_lifecycle.deliver(FixedLifecycle::worker_stopped_after_admission(
995 role.name(),
996 stopped,
997 recovery.0,
998 ));
999 drop(worker);
1000 drop(readiness);
1001 let (member, operation) =
1002 UnrecoveredMember::begin_after_worker_stop_transfer(role, creation, proxy);
1003 (member, operation, lifecycle)
1004 }
1005 (stopped, None) | (stopped, Some(_)) => {
1006 let (member, operation) = UnrecoveredMember::begin_replacement_failure(
1007 role, creation, proxy, worker, readiness, stopped,
1008 );
1009 (member, operation, LifecycleRoute::Sends::empty())
1010 }
1011 };
1012 owners.push(RosterOwner::Unrecovered(member));
1013 match (&self.diagnostics, self.failure_reaction) {
1014 (DiagnosticDisposition::Terminate, _) => {
1015 self.roster = FixedRoster::Terminating {
1016 owners,
1017 prepared: Vec::new(),
1018 };
1019 let mut proxy_operations = SourceActions::empty();
1020 proxy_operations.send(operation);
1021 let mut requests = FixedSupervisorRequests::empty();
1022 requests.proxy_operations = proxy_operations;
1023 requests.lifecycle = lifecycle;
1024 requests.diagnostics =
1025 InterpreterRequests::one(DiagnosticAction::terminal(diagnostic));
1026 Ok(Actions::new(
1027 requests,
1028 Creations::empty(),
1029 Step::Stop(behavior::Stopped),
1030 ))
1031 }
1032 (DiagnosticDisposition::DeliverTo(_), FailureReaction::RetireMember) => {
1033 self.roster = FixedRoster::Operating(owners);
1034 let mut proxy_operations = SourceActions::empty();
1035 proxy_operations.send(operation);
1036 let mut requests = FixedSupervisorRequests::empty();
1037 requests.proxy_operations = proxy_operations;
1038 requests.lifecycle = lifecycle;
1039 requests.diagnostics =
1040 InterpreterRequests::one(self.diagnostics.action(diagnostic));
1041 Ok(Actions::send(requests))
1042 }
1043 (DiagnosticDisposition::DeliverTo(_), FailureReaction::StopSupervisor) => {
1044 let (shutdown, mut operations, schedule) =
1045 FixedShutdown::begin(owners, self.actor_drain);
1046 operations.push((position, operation));
1047 operations.sort_by_key(|(position, _)| *position);
1048 self.roster = FixedRoster::ShuttingDown(shutdown);
1049 let mut proxy_operations = SourceActions::empty();
1050 for (_, operation) in operations {
1051 proxy_operations.send(operation);
1052 }
1053 let mut restart_schedules = SourceActions::empty();
1054 match schedule {
1055 Some(schedule) => restart_schedules.send(schedule),
1056 None => {}
1057 }
1058 let mut requests = FixedSupervisorRequests::empty();
1059 requests.proxy_operations = proxy_operations;
1060 requests.restart_schedules = restart_schedules;
1061 requests.lifecycle = lifecycle;
1062 requests.diagnostics =
1063 InterpreterRequests::one(self.diagnostics.action(diagnostic));
1064 Ok(Actions::send(requests))
1065 }
1066 }
1067 }
1068
1069 fn accept_worker_stop(
1070 &mut self,
1071 stopped: AcceptedWorkerStop<Role, Worker, Plan>,
1072 ) -> BehaviorActed<Self> {
1073 let recovery = mem::replace(&mut self.recovery, SupervisorRecovery::temporary());
1074 match recovery.choose(stopped.kind()) {
1075 RecoveryChoice::LeaveEmpty(recovery) => {
1076 self.recovery = recovery;
1077 match self.lifecycle.clone() {
1078 Some(route_to_lifecycle) => {
1079 let (roster, role, worker_stop) = stopped.release_stop();
1080 self.roster = roster;
1081 let mut requests = FixedSupervisorRequests::empty();
1082 requests.lifecycle = route_to_lifecycle
1083 .deliver(FixedLifecycle::worker_stopped_ineligible(role, worker_stop));
1084 Ok(Actions::send(requests))
1085 }
1086 None => {
1087 self.roster = stopped.discharge_stop();
1088 Ok(Actions::cont())
1089 }
1090 }
1091 }
1092 RecoveryChoice::SourceUnavailable(recovery) => {
1093 self.recovery = recovery;
1094 let (roster, report) = stopped.reject();
1095 self.roster = roster;
1096 Err(FixedSupervisorError::InputRejected {
1097 input: FixedSupervisorEvent::ProxyReported(report),
1098 })
1099 }
1100 RecoveryChoice::Prepare(preparation) => match stopped.begin_recovery(preparation) {
1101 Ok((recovery, roster, request)) => {
1102 self.recovery = recovery;
1103 self.roster = roster;
1104 let mut worker_preparations = SourceActions::empty();
1105 worker_preparations.send(request);
1106 let mut requests = FixedSupervisorRequests::empty();
1107 requests.worker_preparations = worker_preparations;
1108 Ok(Actions::send(requests))
1109 }
1110 Err((stopped, preparation)) => {
1111 self.recovery = preparation.restore();
1112 let (roster, report) = stopped.reject();
1113 self.roster = roster;
1114 Err(FixedSupervisorError::InputRejected {
1115 input: FixedSupervisorEvent::ProxyReported(report),
1116 })
1117 }
1118 },
1119 }
1120 }
1121
1122 fn apply_failed_preparation(
1123 &mut self,
1124 awaiting: recovery::AwaitingWorkerSource,
1125 failed: recovery::FailedPreparation<Role, Worker, Plan, Source>,
1126 ) -> BehaviorActed<Self> {
1127 match (&self.diagnostics, self.failure_reaction) {
1128 (DiagnosticDisposition::Terminate, _) => {
1129 let (roster, source, failure) = failed.terminate();
1130 self.recovery = awaiting.restore(source);
1131 self.roster = roster;
1132 let mut requests = FixedSupervisorRequests::empty();
1133 requests.diagnostics = InterpreterRequests::one(DiagnosticAction::terminal(
1134 FixedDiagnostic::WorkerPreparationFailed(failure),
1135 ));
1136 Ok(Actions::new(
1137 requests,
1138 Creations::empty(),
1139 Step::Stop(behavior::Stopped),
1140 ))
1141 }
1142 (DiagnosticDisposition::DeliverTo(_), FailureReaction::RetireMember) => {
1143 let (roster, source, failure, operation) = failed.retire_member();
1144 self.recovery = awaiting.restore(source);
1145 self.roster = roster;
1146 let mut proxy_operations = SourceActions::empty();
1147 proxy_operations.send(operation);
1148 let mut requests = FixedSupervisorRequests::empty();
1149 requests.proxy_operations = proxy_operations;
1150 requests.diagnostics = InterpreterRequests::one(
1151 self.diagnostics
1152 .action(FixedDiagnostic::WorkerPreparationFailed(failure)),
1153 );
1154 Ok(Actions::send(requests))
1155 }
1156 (DiagnosticDisposition::DeliverTo(_), FailureReaction::StopSupervisor) => {
1157 let (shutdown, source, failure, operations, schedule) =
1158 failed.stop_supervisor(self.actor_drain);
1159 self.recovery = awaiting.restore(source);
1160 self.roster = FixedRoster::ShuttingDown(shutdown);
1161 let mut proxy_operations = SourceActions::empty();
1162 for operation in operations {
1163 proxy_operations.send(operation);
1164 }
1165 let mut restart_schedules = SourceActions::empty();
1166 match schedule {
1167 Some(schedule) => restart_schedules.send(schedule),
1168 None => {}
1169 }
1170 let mut requests = FixedSupervisorRequests::empty();
1171 requests.proxy_operations = proxy_operations;
1172 requests.restart_schedules = restart_schedules;
1173 requests.diagnostics = InterpreterRequests::one(
1174 self.diagnostics
1175 .action(FixedDiagnostic::WorkerPreparationFailed(failure)),
1176 );
1177 Ok(Actions::send(requests))
1178 }
1179 }
1180 }
1181
1182 fn accept_worker_preparation_start(
1183 &mut self,
1184 roster: FixedRoster<Role, Worker, Plan>,
1185 input: ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
1186 ) -> BehaviorActed<Self> {
1187 let recovery = mem::replace(&mut self.recovery, SupervisorRecovery::temporary());
1188 let awaiting = match recovery.expect_source() {
1189 Ok(awaiting) => awaiting,
1190 Err(recovery) => {
1191 self.recovery = recovery;
1192 self.roster = roster;
1193 return Err(FixedSupervisorError::InputRejected {
1194 input: FixedSupervisorEvent::WorkerPreparationStarted(input),
1195 });
1196 }
1197 };
1198 let roster = match roster {
1199 FixedRoster::ShuttingDown(shutdown) => {
1200 return match shutdown.accept_preparation_start(input) {
1201 Ok(ControlFlow::Continue(shutdown)) => {
1202 self.recovery = awaiting.waiting();
1203 self.retain_shutdown(shutdown, Vec::new())
1204 }
1205 Ok(ControlFlow::Break((shutdown, returned))) => {
1206 self.recovery = awaiting
1207 .retain_for_retirement(WorkerPreparationRetirement::Start(returned));
1208 self.retain_shutdown(shutdown, Vec::new())
1209 }
1210 Err((shutdown, input)) => {
1211 self.recovery = awaiting.waiting();
1212 self.roster = FixedRoster::ShuttingDown(shutdown);
1213 Err(FixedSupervisorError::InputRejected {
1214 input: FixedSupervisorEvent::WorkerPreparationStarted(input),
1215 })
1216 }
1217 };
1218 }
1219 roster => roster,
1220 };
1221 match roster.accept_preparation_start(input) {
1222 Ok(ControlFlow::Continue(roster)) => {
1223 self.recovery = awaiting.waiting();
1224 self.roster = roster;
1225 Ok(Actions::cont())
1226 }
1227 Ok(ControlFlow::Break(failed)) => self.apply_failed_preparation(awaiting, failed),
1228 Err((roster, input)) => {
1229 self.recovery = awaiting.waiting();
1230 self.roster = roster;
1231 Err(FixedSupervisorError::InputRejected {
1232 input: FixedSupervisorEvent::WorkerPreparationStarted(input),
1233 })
1234 }
1235 }
1236 }
1237
1238 fn accept_worker_preparation_return(
1239 &mut self,
1240 roster: FixedRoster<Role, Worker, Plan>,
1241 returned: WorkerPreparation<Source, Role, Worker, Plan>,
1242 ) -> BehaviorActed<Self> {
1243 let recovery = mem::replace(&mut self.recovery, SupervisorRecovery::temporary());
1244 let awaiting = match recovery.expect_source() {
1245 Ok(awaiting) => awaiting,
1246 Err(recovery) => {
1247 self.recovery = recovery;
1248 self.roster = roster;
1249 return Err(FixedSupervisorError::InputRejected {
1250 input: FixedSupervisorEvent::WorkerPreparationReturned(returned),
1251 });
1252 }
1253 };
1254 let roster = match roster {
1255 FixedRoster::ShuttingDown(shutdown) => {
1256 return match shutdown.accept_preparation_return(returned) {
1257 Ok((shutdown, returned)) => {
1258 self.recovery = awaiting
1259 .retain_for_retirement(WorkerPreparationRetirement::Returned(returned));
1260 self.retain_shutdown(shutdown, Vec::new())
1261 }
1262 Err((shutdown, returned)) => {
1263 self.recovery = awaiting.waiting();
1264 self.roster = FixedRoster::ShuttingDown(shutdown);
1265 Err(FixedSupervisorError::InputRejected {
1266 input: FixedSupervisorEvent::WorkerPreparationReturned(returned),
1267 })
1268 }
1269 };
1270 }
1271 roster => roster,
1272 };
1273 match roster.accept_preparation_return(returned) {
1274 Ok(PreparationAcceptance::Prepared {
1275 owners,
1276 position,
1277 prepared,
1278 source,
1279 }) => match awaiting.decide_prepared(owners, position, prepared, source) {
1280 PreparedRecoveryDecision::Admitted {
1281 roster,
1282 recovery,
1283 schedule,
1284 } => self.admit_recovery(roster, recovery, schedule),
1285 PreparedRecoveryDecision::Denied {
1286 failed,
1287 recovery,
1288 diagnostic,
1289 } => {
1290 self.recovery = recovery;
1291 self.apply_recovery_failure(failed, FixedDiagnostic::RecoveryDenied(diagnostic))
1292 }
1293 },
1294 Ok(PreparationAcceptance::Failed(failed)) => {
1295 self.apply_failed_preparation(awaiting, failed)
1296 }
1297 Err((roster, returned)) => {
1298 self.recovery = awaiting.waiting();
1299 self.roster = roster;
1300 Err(FixedSupervisorError::InputRejected {
1301 input: FixedSupervisorEvent::WorkerPreparationReturned(returned),
1302 })
1303 }
1304 }
1305 }
1306
1307 fn accept_initial_failure(
1308 &mut self,
1309 owners: Vec<RosterOwner<Role, Worker, Plan>>,
1310 child: CreationId,
1311 role: RoleName<Role>,
1312 outcome: super::ProxyOutcome<Worker, Plan>,
1313 ) -> BehaviorActed<Self> {
1314 match (&self.diagnostics, self.failure_reaction) {
1315 (DiagnosticDisposition::Terminate, _) => {
1316 self.roster = FixedRoster::Terminating {
1317 owners,
1318 prepared: Vec::new(),
1319 };
1320 let diagnostic =
1321 FixedDiagnostic::ProxyOutcomeFailed(ProxyOutcomeFailure::new(role, outcome));
1322 let mut requests = FixedSupervisorRequests::empty();
1323 requests.diagnostics =
1324 InterpreterRequests::one(DiagnosticAction::terminal(diagnostic));
1325 Ok(Actions::new(
1326 requests,
1327 Creations::empty(),
1328 Step::Stop(behavior::Stopped),
1329 ))
1330 }
1331 (DiagnosticDisposition::DeliverTo(_), FailureReaction::RetireMember) => {
1332 let roster = FixedRoster::Operating(owners);
1333 let (roster, shutdown) = match roster.begin_stop_after_initial_failure(child) {
1334 Ok(started) => started,
1335 Err(roster) => {
1336 self.roster = roster;
1337 return Err(FixedSupervisorError::InitialOutcomeContradiction {
1338 report: behavior::ChildReport::new(child, outcome),
1339 });
1340 }
1341 };
1342 self.roster = roster;
1343 let diagnostic =
1344 FixedDiagnostic::ProxyOutcomeFailed(ProxyOutcomeFailure::new(role, outcome));
1345 let mut proxy_operations = SourceActions::empty();
1346 proxy_operations.send(shutdown);
1347 let mut requests = FixedSupervisorRequests::empty();
1348 requests.proxy_operations = proxy_operations;
1349 requests.diagnostics =
1350 InterpreterRequests::one(self.diagnostics.action(diagnostic));
1351 Ok(Actions::send(requests))
1352 }
1353 (DiagnosticDisposition::DeliverTo(_), FailureReaction::StopSupervisor) => {
1354 let roster = FixedRoster::Operating(owners);
1355 let (shutdown, operations, schedule) = match roster.begin_shutdown(self.actor_drain)
1356 {
1357 Ok(started) => started,
1358 Err(roster) => {
1359 self.roster = roster;
1360 return Err(FixedSupervisorError::InitialOutcomeContradiction {
1361 report: behavior::ChildReport::new(child, outcome),
1362 });
1363 }
1364 };
1365 let diagnostic =
1366 FixedDiagnostic::ProxyOutcomeFailed(ProxyOutcomeFailure::new(role, outcome));
1367 let mut actions = match self.retain_shutdown(shutdown, operations) {
1368 Ok(actions) => actions,
1369 Err(rejection) => return Err(rejection),
1370 };
1371 match schedule {
1372 Some(schedule) => actions.sends.restart_schedules.send(schedule),
1373 None => {}
1374 }
1375 actions.sends.diagnostics =
1376 InterpreterRequests::one(self.diagnostics.action(diagnostic));
1377 Ok(actions)
1378 }
1379 }
1380 }
1381
1382 fn accept_proxy_exit(
1383 &mut self,
1384 roster: FixedRoster<Role, Worker, Plan>,
1385 stopped: crate::ChildStopped<BehaviorAddr<Worker>>,
1386 ) -> BehaviorActed<Self> {
1387 match roster {
1388 FixedRoster::ShuttingDown(shutdown) => match shutdown.accept_stop(stopped) {
1389 Ok(shutdown) => self.retain_shutdown(shutdown, Vec::new()),
1390 Err((shutdown, stopped)) => {
1391 self.roster = FixedRoster::ShuttingDown(shutdown);
1392 Err(FixedSupervisorError::InputRejected {
1393 input: FixedSupervisorEvent::ProxyStopped(stopped),
1394 })
1395 }
1396 },
1397 roster => {
1398 let (next, retired) = match roster.accept_proxy_stop(stopped) {
1399 Ok(next) => next,
1400 Err((next, rejection)) => {
1401 self.roster = next;
1402 return Err(rejection);
1403 }
1404 };
1405 match retired {
1406 Some(role) => self.member_retired(next, role),
1407 None => self.authorize_waiting(next),
1408 }
1409 }
1410 }
1411 }
1412
1413 fn accept_command(
1414 &mut self,
1415 roster: FixedRoster<Role, Worker, Plan>,
1416 from: BehaviorAddr<Worker>,
1417 message: FixedCommand<BehaviorAddr<Worker>, Role, Worker::Protocol>,
1418 ) -> BehaviorActed<Self> {
1419 match message {
1420 FixedCommand::Shutdown => match roster.begin_shutdown(self.actor_drain) {
1421 Ok((shutdown, operations, schedule)) => {
1422 let mut actions = match self.retain_shutdown(shutdown, operations) {
1423 Ok(actions) => actions,
1424 Err(rejection) => return Err(rejection),
1425 };
1426 match schedule {
1427 Some(schedule) => actions.sends.restart_schedules.send(schedule),
1428 None => {}
1429 }
1430 Ok(actions)
1431 }
1432 Err(roster) => {
1433 self.roster = roster;
1434 Err(FixedSupervisorError::InputRejected {
1435 input: FixedSupervisorEvent::Command(User::new(
1436 from,
1437 FixedCommand::Shutdown,
1438 )),
1439 })
1440 }
1441 },
1442 FixedCommand::Status { reply_to } => match roster.snapshot() {
1443 Some(snapshot) => {
1444 self.roster = roster;
1445 let mut requests = FixedSupervisorRequests::empty();
1446 requests.status_replies = reply_to.deliver(snapshot);
1447 Ok(Actions::send(requests))
1448 }
1449 None => {
1450 self.roster = roster;
1451 Err(FixedSupervisorError::InputRejected {
1452 input: FixedSupervisorEvent::Command(User::new(
1453 from,
1454 FixedCommand::Status { reply_to },
1455 )),
1456 })
1457 }
1458 },
1459 FixedCommand::Capability { role, reply_to } => match roster.capability(role) {
1460 Ok(result) => {
1461 self.roster = roster;
1462 let mut requests = FixedSupervisorRequests::empty();
1463 requests.capability_replies = reply_to.deliver(result);
1464 Ok(Actions::send(requests))
1465 }
1466 Err(role) => {
1467 self.roster = roster;
1468 Err(FixedSupervisorError::InputRejected {
1469 input: FixedSupervisorEvent::Command(User::new(
1470 from,
1471 FixedCommand::Capability { role, reply_to },
1472 )),
1473 })
1474 }
1475 },
1476 }
1477 }
1478
1479 fn unexpected_input(
1480 &mut self,
1481 input: FixedSupervisorEvent<
1482 Role,
1483 Worker,
1484 Plan,
1485 ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
1486 WorkerPreparation<Source, Role, Worker, Plan>,
1487 >,
1488 ) -> BehaviorActed<Self> {
1489 let become_ = match &self.diagnostics {
1490 DiagnosticDisposition::DeliverTo(_) => Step::Continue,
1491 DiagnosticDisposition::Terminate => Step::Stop(behavior::Stopped),
1492 };
1493 let mut requests = FixedSupervisorRequests::empty();
1494 requests.diagnostics = InterpreterRequests::one(
1495 self.diagnostics
1496 .action(FixedDiagnostic::UnexpectedInput { input }),
1497 );
1498 Ok(Actions::new(requests, Creations::empty(), become_))
1499 }
1500}
1501
1502impl<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute> Behavior
1503 for FixedSupervisor<Role, Worker, Plan, Source, DiagnosticRoute, LifecycleRoute>
1504where
1505 Role: Eq + Send + Sync,
1506 Worker: Behavior + Send,
1507 Plan: ActivationPlan,
1508 Source: WorkerSource<Role, Worker, Plan>,
1509 BehaviorAddr<Worker>: EndpointAddress,
1510 <BehaviorAddr<Worker> as Address>::Nonce: Send,
1511 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
1512 EstablishedActor<StableProxy<Worker, Plan>>: Clone + Send,
1513 DiagnosticRoute: crate::DiagnosticRoute<FixedDiagnostic<Role, Worker, Plan, Source>> + Clone,
1514 LifecycleRoute: FixedLifecycleRoute<Role, Worker, Plan> + Clone,
1515 LifecycleRoute::Sends: behavior::SendsFor<
1516 FixedSupervisorEvent<
1517 Role,
1518 Worker,
1519 Plan,
1520 ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
1521 WorkerPreparation<Source, Role, Worker, Plan>,
1522 >,
1523 >,
1524{
1525 type Protocol = MessageProtocol<
1526 BehaviorAddr<Worker>,
1527 FixedCommand<BehaviorAddr<Worker>, Role, Worker::Protocol>,
1528 >;
1529 type Event = FixedSupervisorEvent<
1530 Role,
1531 Worker,
1532 Plan,
1533 ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
1534 WorkerPreparation<Source, Role, Worker, Plan>,
1535 >;
1536 type Sends = FixedSupervisorRequests<
1537 InterpreterRequests<ObserveChild<Worker::Protocol, ChildHead>>,
1538 SourceActions<PrepareWorkers<Source, Role, Worker, Plan>>,
1539 SourceActions<ProxyOperation<behavior::Here, Worker, Plan>>,
1540 SourceActions<ScheduleAfter>,
1541 LifecycleRoute::Sends,
1542 <ReplyRoute<
1543 MessageProtocol<BehaviorAddr<Worker>, FixedSnapshot<Worker::Protocol>>,
1544 > as DeliveryRoute>::Sends,
1545 <ReplyRoute<
1546 MessageProtocol<
1547 BehaviorAddr<Worker>,
1548 CapabilityResult<Role, Worker::Protocol>,
1549 >,
1550 > as DeliveryRoute>::Sends,
1551 InterpreterRequests<
1552 DiagnosticAction<DiagnosticRoute, FixedDiagnostic<Role, Worker, Plan, Source>>,
1553 >,
1554 >;
1555 type Ph = Never;
1556 type Error = FixedSupervisorError<
1557 Role,
1558 Worker,
1559 Plan,
1560 ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
1561 WorkerPreparation<Source, Role, Worker, Plan>,
1562 >;
1563 type Birth = Births<StableProxy<Worker, Plan>>;
1564
1565 fn init(&mut self, _: InitializationTurn) -> BehaviorActed<Self> {
1566 let roster = mem::replace(&mut self.roster, FixedRoster::Stopped);
1567 match roster.begin(&mut self.creations) {
1568 Ok((next, observations, proxies)) => {
1569 self.roster = next;
1570 let mut sends = FixedSupervisorRequests::empty();
1571 sends.proxy_observations = InterpreterRequests::new(observations);
1572 Ok(Actions::new(sends, proxies, Step::Continue))
1573 }
1574 Err((roster, rejection)) => {
1575 self.roster = roster;
1576 Err(rejection)
1577 }
1578 }
1579 }
1580
1581 fn transition(&mut self, _: ActiveTurn, input: Self::Event) -> BehaviorActed<Self> {
1582 let roster = mem::replace(&mut self.roster, FixedRoster::Stopped);
1583 let acted = match input {
1584 FixedSupervisorEvent::ProxyCreationsSettled(proxies) => {
1585 self.accept_proxy_birth(roster, proxies)
1586 }
1587 FixedSupervisorEvent::ProxyInputSettled(settlement) => {
1588 self.accept_proxy_operation(roster, settlement)
1589 }
1590 FixedSupervisorEvent::ProxyReported(report) => self.accept_proxy_report(roster, report),
1591 FixedSupervisorEvent::ProxyDiagnosed(report) => {
1592 self.roster = roster;
1593 Err(FixedSupervisorError::InputRejected {
1594 input: FixedSupervisorEvent::ProxyDiagnosed(report),
1595 })
1596 }
1597 FixedSupervisorEvent::ProxyStopped(stopped) => self.accept_proxy_exit(roster, stopped),
1598 FixedSupervisorEvent::WorkerPreparationStarted(started) => {
1599 self.accept_worker_preparation_start(roster, started)
1600 }
1601 FixedSupervisorEvent::WorkerPreparationReturned(returned) => {
1602 self.accept_worker_preparation_return(roster, returned)
1603 }
1604 FixedSupervisorEvent::RestartScheduleSettled(input) => {
1605 self.accept_restart_schedule(roster, input)
1606 }
1607 FixedSupervisorEvent::RestartElapsed(elapsed) => {
1608 self.accept_restart_elapsed(roster, elapsed)
1609 }
1610 FixedSupervisorEvent::Command(User { from, message }) => {
1611 self.accept_command(roster, from, message)
1612 }
1613 };
1614 match acted {
1615 Err(FixedSupervisorError::InputRejected { input }) => self.unexpected_input(input),
1616 acted => acted,
1617 }
1618 }
1619}
1620
1621#[cfg(test)]
1622mod tests {
1623 use crate::DiagnosticDisposition;
1624
1625 use super::{
1626 ActivationPolicy, ActorDrainPolicy, FailureReaction, FixedBuilder, OrderedRoles, Recovery,
1627 fixed,
1628 };
1629
1630 #[test]
1631 fn builder_owns_each_policy_directly() {
1632 let roles = OrderedRoles::new(1_u8, [2_u8]).expect("roles are distinct");
1633 let builder = fixed(
1634 |_: &u8| (),
1635 roles,
1636 ActivationPolicy::new(1).expect("activation capacity is positive"),
1637 Recovery::temporary(),
1638 FailureReaction::StopSupervisor,
1639 ActorDrainPolicy::WaitForActorGraph,
1640 DiagnosticDisposition::terminate(),
1641 );
1642
1643 let FixedBuilder {
1644 factory: _,
1645 roles: _,
1646 activation: _,
1647 recovery: _,
1648 failure_reaction: _,
1649 actor_drain: _,
1650 diagnostics: _,
1651 lifecycle: _,
1652 } = builder;
1653 }
1654}