1use crate::{ChildShutdownRejected, ChildShutdownRejection, ChildStopped, ShutdownChild};
4use behavior::{
5 ActionItem, ActionItemResult, Actions, Address, Behavior, BehaviorActed, BirthMode, ChildHead,
6 ChildRole, ChildTail, ClassifySettlement, CreationId, EventIngress, Here, InjectEvent, Inside,
7 Interpretation, InterpretationProgress, InterpreterRequests, ItemSettlement, SendEffects,
8 SendLayer, SettledItem, SourceCustody, SourceProgress, SourceSettlementCustody, finish_item,
9 prepare_item,
10};
11use behavior::{User, UserEvent};
12
13#[derive(Debug, Clone, PartialEq, Eq)]
15pub struct ShutdownPlan<N> {
16 phases: Vec<Vec<N>>,
17}
18
19impl<N: Copy + Eq> ShutdownPlan<N> {
20 pub fn new(phases: impl IntoIterator<Item = Vec<N>>) -> Result<Self, ShutdownPlanError<N>> {
25 let phases: Vec<_> = phases.into_iter().collect();
26 let mut seen = Vec::new();
27 for (phase, children) in phases.iter().enumerate() {
28 if children.is_empty() {
29 return Err(ShutdownPlanError::EmptyPhase { phase });
30 }
31 for &child in children {
32 if seen.contains(&child) {
33 return Err(ShutdownPlanError::DuplicateChild(child));
34 }
35 seen.push(child);
36 }
37 }
38 Ok(Self { phases })
39 }
40
41 #[must_use]
42 pub fn phases(&self) -> &[Vec<N>] {
43 &self.phases
44 }
45}
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
49pub enum ShutdownPlanError<N> {
50 #[error("shutdown phase {phase} has no children")]
51 EmptyPhase { phase: usize },
52 #[error("a child occurs in more than one shutdown position")]
53 DuplicateChild(N),
54}
55
56pub struct ShutdownTree<N> {
61 plan: ShutdownPlan<N>,
62}
63
64impl<N: Copy + Eq> ShutdownTree<N> {
65 pub fn new(
70 nodes: impl IntoIterator<Item = N>,
71 edges: impl IntoIterator<Item = (N, N)>,
72 ) -> Result<Self, ShutdownTreeError<N>> {
73 let nodes: Vec<_> = nodes.into_iter().collect();
74 let mut unique = Vec::new();
75 for &node in &nodes {
76 if unique.contains(&node) {
77 return Err(ShutdownTreeError::DuplicateChild(node));
78 }
79 unique.push(node);
80 }
81 let edges: Vec<_> = edges.into_iter().collect();
82 for &(dependent, dependency) in &edges {
83 if !nodes.contains(&dependent) {
84 return Err(ShutdownTreeError::UnknownChild(dependent));
85 }
86 if !nodes.contains(&dependency) {
87 return Err(ShutdownTreeError::UnknownChild(dependency));
88 }
89 }
90 let mut remaining = nodes;
91 let mut phases = Vec::new();
92 while !remaining.is_empty() {
93 let phase: Vec<_> = remaining
94 .iter()
95 .copied()
96 .filter(|candidate| {
97 !edges.iter().any(|(dependent, dependency)| {
98 dependency == candidate && remaining.contains(dependent)
99 })
100 })
101 .collect();
102 if phase.is_empty() {
103 return Err(ShutdownTreeError::Cycle);
104 }
105 remaining.retain(|node| !phase.contains(node));
106 phases.push(phase);
107 }
108 Ok(Self {
109 plan: ShutdownPlan { phases },
110 })
111 }
112
113 #[must_use]
114 pub fn into_plan(self) -> ShutdownPlan<N> {
115 self.plan
116 }
117}
118
119#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
121pub enum ShutdownTreeError<N> {
122 #[error("a child is declared more than once")]
123 DuplicateChild(N),
124 #[error("a dependency edge names an undeclared child")]
125 UnknownChild(N),
126 #[error("the shutdown dependency topology contains a cycle")]
127 Cycle,
128}
129
130pub enum ShutdownChoice<C: Behavior, Tail> {
161 Child {
162 creation: CreationId,
163 child: core::marker::PhantomData<fn() -> C>,
164 },
165 Other(Tail),
166}
167
168pub struct NoShutdownTargets<A: Address> {
170 never: behavior::Never,
171 address: core::marker::PhantomData<fn() -> A>,
172}
173
174impl<A: Address> Copy for NoShutdownTargets<A> {}
175impl<A: Address> Clone for NoShutdownTargets<A> {
176 fn clone(&self) -> Self {
177 *self
178 }
179}
180
181impl<C: Behavior, Tail> ShutdownChoice<C, Tail> {
182 #[must_use]
183 pub const fn child(creation: CreationId) -> Self {
184 Self::Child {
185 creation,
186 child: core::marker::PhantomData,
187 }
188 }
189
190 #[must_use]
191 pub const fn other(target: Tail) -> Self {
192 Self::Other(target)
193 }
194}
195
196pub trait ShutdownTargetAt<Child: Behavior, Position>: named_shutdown::Target + Sized {
202 fn shutdown_target_at(creation: CreationId) -> Self;
204}
205
206mod named_shutdown {
207 pub trait Target {}
208}
209
210impl<Child: Behavior, Tail> named_shutdown::Target for ShutdownChoice<Child, Tail> {}
211
212impl<Child: Behavior, Tail> ShutdownTargetAt<Child, ChildHead> for ShutdownChoice<Child, Tail> {
213 fn shutdown_target_at(creation: CreationId) -> Self {
214 Self::child(creation)
215 }
216}
217
218impl<Head, Tail, Child, Position> ShutdownTargetAt<Child, ChildTail<Position>>
219 for ShutdownChoice<Head, Tail>
220where
221 Head: Behavior,
222 Child: Behavior,
223 Tail: ShutdownTargetAt<Child, Position>,
224{
225 fn shutdown_target_at(creation: CreationId) -> Self {
226 Self::other(Tail::shutdown_target_at(creation))
227 }
228}
229
230#[must_use]
284pub fn shutdown_target<Parent, Role, Targets>(_: Role, creation: CreationId) -> Targets
285where
286 Parent: Behavior,
287 Role: ChildRole<Parent>,
288 Targets: ShutdownTargetAt<Role::Child, Role::Position>,
289{
290 Targets::shutdown_target_at(creation)
291}
292
293impl<C, Tail> Copy for ShutdownChoice<C, Tail>
294where
295 C: Behavior,
296 <behavior::BehaviorAddr<C> as Address>::Nonce: Copy,
297 Tail: Copy,
298{
299}
300
301impl<C, Tail> Clone for ShutdownChoice<C, Tail>
302where
303 C: Behavior,
304 <behavior::BehaviorAddr<C> as Address>::Nonce: Copy,
305 Tail: Copy,
306{
307 fn clone(&self) -> Self {
308 *self
309 }
310}
311
312#[doc(hidden)]
314pub enum HeterogeneousShutdownItem<Child, Tail> {
315 Child(Child),
316 Other(Tail),
317}
318
319impl<Child, Tail> behavior::ClassifySettlement for HeterogeneousShutdownItem<Child, Tail>
320where
321 Child: behavior::ClassifySettlement,
322 Tail: behavior::ClassifySettlement,
323{
324 fn settlement_status(&self) -> behavior::SettlementStatus {
325 match self {
326 Self::Child(child) => child.settlement_status(),
327 Self::Other(other) => other.settlement_status(),
328 }
329 }
330}
331
332pub(crate) mod heterogeneous {
333 use super::*;
334
335 pub trait Selection: Sized + Send {
336 type Addr: Address;
337 fn creation(&self) -> CreationId;
338 }
339
340 #[doc(hidden)]
341 pub trait ChoiceSettlements<Occurrence>: Sized + Send {
342 type Settlements: Send + ClassifySettlement;
343 type InterpretationCustody;
344 fn prepare_interpretation(
345 progress: &mut Option<
346 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
347 >,
348 );
349 fn finish_interpretation(
350 progress: &mut Option<
351 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
352 >,
353 );
354 fn unattempted(
355 progress: &mut Option<
356 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
357 >,
358 );
359 }
360
361 pub trait InterpretChoice<I, E, Path, Occurrence>: ChoiceSettlements<Occurrence> {
362 fn settle(
363 progress: &mut Option<
364 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
365 >,
366 interpreter: &mut I,
367 ) -> impl core::future::Future<Output = ()> + Send;
368 }
369
370 impl<A: Address> Selection for NoShutdownTargets<A> {
371 type Addr = A;
372 fn creation(&self) -> CreationId {
373 match self.never {}
374 }
375 }
376
377 impl<A: Address, Occurrence> ChoiceSettlements<Occurrence> for NoShutdownTargets<A> {
378 type Settlements = behavior::Never;
379 type InterpretationCustody = Self;
380 fn prepare_interpretation(
381 progress: &mut Option<
382 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
383 >,
384 ) {
385 match progress.take() {
386 Some(
387 InterpretationProgress::Original(original)
388 | InterpretationProgress::Interpreting(original),
389 ) => match original.never {},
390 retained => *progress = retained,
391 }
392 }
393 fn finish_interpretation(
394 progress: &mut Option<
395 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
396 >,
397 ) {
398 <Self as ChoiceSettlements<Occurrence>>::prepare_interpretation(progress);
399 }
400 fn unattempted(
401 progress: &mut Option<
402 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
403 >,
404 ) {
405 <Self as ChoiceSettlements<Occurrence>>::prepare_interpretation(progress);
406 }
407 }
408
409 impl<I: Send, E, Path, A: Address, Occurrence> InterpretChoice<I, E, Path, Occurrence>
410 for NoShutdownTargets<A>
411 {
412 async fn settle(
413 progress: &mut Option<
414 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
415 >,
416 _: &mut I,
417 ) {
418 <Self as ChoiceSettlements<Occurrence>>::prepare_interpretation(progress);
419 }
420 }
421
422 impl<C, Tail> Selection for ShutdownChoice<C, Tail>
423 where
424 C: Behavior,
425 C::Event: InjectEvent<crate::ShutdownRequested, Here>,
426 Tail: Selection<Addr = behavior::BehaviorAddr<C>>,
427 {
428 type Addr = behavior::BehaviorAddr<C>;
429 fn creation(&self) -> CreationId {
430 match self {
431 Self::Child { creation, .. } => *creation,
432 Self::Other(target) => target.creation(),
433 }
434 }
435 }
436
437 impl<C, Tail, Occurrence> ChoiceSettlements<Occurrence> for ShutdownChoice<C, Tail>
438 where
439 C: Behavior,
440 ShutdownChild<C, Occurrence>: ActionItem,
441 Tail: ChoiceSettlements<ChildTail<Occurrence>>,
442 {
443 type Settlements = HeterogeneousShutdownItem<
444 ActionItemResult<ShutdownChild<C, Occurrence>>,
445 Tail::Settlements,
446 >;
447 type InterpretationCustody = HeterogeneousShutdownItem<
448 (
449 Option<ShutdownChild<C, Occurrence>>,
450 Option<
451 ItemSettlement<
452 ShutdownChild<C, Occurrence>,
453 <ShutdownChild<C, Occurrence> as ActionItem>::Accepted,
454 <ShutdownChild<C, Occurrence> as ActionItem>::Rejection,
455 <ShutdownChild<C, Occurrence> as ActionItem>::Prerequisite,
456 >,
457 >,
458 ),
459 Option<InterpretationProgress<Tail, Tail::InterpretationCustody, Tail::Settlements>>,
460 >;
461 fn prepare_interpretation(
462 progress: &mut Option<
463 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
464 >,
465 ) {
466 *progress = match progress.take() {
467 Some(InterpretationProgress::Original(Self::Child { creation, .. })) => Some(
468 InterpretationProgress::Interpreting(HeterogeneousShutdownItem::Child((
469 Some(ShutdownChild::new(creation)),
470 None,
471 ))),
472 ),
473 Some(InterpretationProgress::Original(Self::Other(target))) => Some(
474 InterpretationProgress::Interpreting(HeterogeneousShutdownItem::Other(Some(
475 InterpretationProgress::Original(target),
476 ))),
477 ),
478 retained => retained,
479 };
480 }
481 fn finish_interpretation(
482 progress: &mut Option<
483 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
484 >,
485 ) {
486 let complete = matches!(
487 progress,
488 Some(InterpretationProgress::Interpreting(
489 HeterogeneousShutdownItem::Child((None, Some(_)))
490 )) | Some(InterpretationProgress::Interpreting(
491 HeterogeneousShutdownItem::Other(Some(InterpretationProgress::Completed(_)))
492 ))
493 );
494 if !complete {
495 return;
496 }
497 *progress = match progress.take() {
498 Some(InterpretationProgress::Interpreting(HeterogeneousShutdownItem::Child((
499 None,
500 Some(received),
501 )))) => {
502 let received = match received {
503 received @ ItemSettlement::Corrupt { .. } => Interpretation::Corrupt(
504 HeterogeneousShutdownItem::Child(SettledItem::Attempted(received)),
505 ),
506 received => Interpretation::Complete(HeterogeneousShutdownItem::Child(
507 SettledItem::Attempted(received),
508 )),
509 };
510 Some(InterpretationProgress::Completed(received))
511 }
512 Some(InterpretationProgress::Interpreting(HeterogeneousShutdownItem::Other(
513 Some(InterpretationProgress::Completed(received)),
514 ))) => Some(InterpretationProgress::Completed(
515 received.map(HeterogeneousShutdownItem::Other),
516 )),
517 retained => retained,
518 };
519 }
520 fn unattempted(
521 progress: &mut Option<
522 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
523 >,
524 ) {
525 if let Some(InterpretationProgress::Original(Self::Child { .. })) = progress {
526 *progress = match progress.take() {
527 Some(InterpretationProgress::Original(Self::Child { creation, .. })) => {
528 Some(InterpretationProgress::Completed(Interpretation::Complete(
529 HeterogeneousShutdownItem::Child(SettledItem::Unattempted(
530 ShutdownChild::new(creation),
531 )),
532 )))
533 }
534 retained => retained,
535 };
536 return;
537 }
538 <Self as ChoiceSettlements<Occurrence>>::prepare_interpretation(progress);
539 if let Some(InterpretationProgress::Interpreting(HeterogeneousShutdownItem::Other(
540 target,
541 ))) = progress
542 {
543 Tail::unattempted(target);
544 }
545 <Self as ChoiceSettlements<Occurrence>>::finish_interpretation(progress);
546 }
547 }
548
549 impl<I, E, Path, C, Tail, Occurrence> InterpretChoice<I, E, Path, Occurrence>
550 for ShutdownChoice<C, Tail>
551 where
552 I: behavior::InterpretItem<ShutdownChild<C, Occurrence>, E, Path> + Send,
553 C: Behavior,
554 <behavior::BehaviorAddr<C> as Address>::Nonce: Send,
555 Tail: InterpretChoice<I, E, Path, ChildTail<Occurrence>>,
556 Tail::InterpretationCustody: Send,
557 {
558 async fn settle(
559 progress: &mut Option<
560 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
561 >,
562 interpreter: &mut I,
563 ) {
564 <Self as ChoiceSettlements<Occurrence>>::prepare_interpretation(progress);
565 match progress {
566 Some(InterpretationProgress::Interpreting(HeterogeneousShutdownItem::Child((
567 input,
568 received,
569 )))) => {
570 if matches!((&*input, &*received), (Some(_), None)) {
571 <I as behavior::InterpretItem<ShutdownChild<C,Occurrence>,E,Path>>::interpret_item(interpreter,input,received).await;
572 }
573 }
574 Some(InterpretationProgress::Interpreting(HeterogeneousShutdownItem::Other(
575 target,
576 ))) => {
577 Tail::settle(target, interpreter).await;
578 Tail::finish_interpretation(target);
579 }
580 _ => return,
581 }
582 <Self as ChoiceSettlements<Occurrence>>::finish_interpretation(progress);
583 }
584 }
585}
586
587#[doc(hidden)]
588pub use heterogeneous::ChoiceSettlements as HeterogeneousShutdownChoiceSettlement;
589
590#[derive(Clone, PartialEq, Eq)]
592pub struct HeterogeneousShutdownPlan<T: heterogeneous::Selection> {
593 phases: Vec<Vec<T>>,
594}
595
596impl<T: heterogeneous::Selection> core::fmt::Debug for HeterogeneousShutdownPlan<T> {
597 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
598 formatter
599 .debug_struct("HeterogeneousShutdownPlan")
600 .field("phase_count", &self.phases.len())
601 .finish_non_exhaustive()
602 }
603}
604
605impl<T> HeterogeneousShutdownPlan<T>
606where
607 T: heterogeneous::Selection,
608{
609 pub fn new(
612 phases: impl IntoIterator<Item = Vec<T>>,
613 ) -> Result<Self, ShutdownPlanError<CreationId>> {
614 let phases: Vec<_> = phases.into_iter().collect();
615 let mut seen = Vec::new();
616 for (phase, children) in phases.iter().enumerate() {
617 if children.is_empty() {
618 return Err(ShutdownPlanError::EmptyPhase { phase });
619 }
620 for child in children {
621 let creation = heterogeneous::Selection::creation(child);
622 if seen.contains(&creation) {
623 return Err(ShutdownPlanError::DuplicateChild(creation));
624 }
625 seen.push(creation);
626 }
627 }
628 Ok(Self { phases })
629 }
630
631 #[must_use]
632 pub fn phases(&self) -> &[Vec<T>] {
633 &self.phases
634 }
635}
636
637pub struct HeterogeneousShutdownSends<T> {
640 requests: Vec<T>,
641}
642
643impl<T> SendEffects for HeterogeneousShutdownSends<T> {
644 fn empty() -> Self {
645 Self {
646 requests: Vec::new(),
647 }
648 }
649
650 fn append(&mut self, other: Self) {
651 self.requests.extend(other.requests);
652 }
653}
654
655impl<T> HeterogeneousShutdownSends<T> {
656 #[must_use]
661 pub fn as_slice(&self) -> &[T] {
662 &self.requests
663 }
664}
665
666impl<E, T: Send> behavior::SendsFor<E> for HeterogeneousShutdownSends<T> {}
667
668impl<T> behavior::ClassifySettlement for HeterogeneousShutdownSends<T>
669where
670 T: behavior::ClassifySettlement,
671{
672 fn settlement_status(&self) -> behavior::SettlementStatus {
673 self.requests.settlement_status()
674 }
675}
676
677impl<T> behavior::SendSettlements for HeterogeneousShutdownSends<T>
678where
679 T: heterogeneous::ChoiceSettlements<ChildHead>,
680{
681 type Settlements = HeterogeneousShutdownSends<T::Settlements>;
682 type SourceCustody = Self::Settlements;
683 type InterpretationCustody =
684 Vec<Option<InterpretationProgress<T, T::InterpretationCustody, T::Settlements>>>;
685 fn prepare_interpretation(
686 progress: &mut Option<
687 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
688 >,
689 ) {
690 *progress = match progress.take() {
691 Some(InterpretationProgress::Original(original)) => {
692 Some(InterpretationProgress::Interpreting(
693 original
694 .requests
695 .into_iter()
696 .map(|request| Some(InterpretationProgress::Original(request)))
697 .collect(),
698 ))
699 }
700 retained => retained,
701 };
702 }
703 fn finish_interpretation(
704 progress: &mut Option<
705 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
706 >,
707 ) {
708 let Some(InterpretationProgress::Interpreting(requests)) = progress.as_ref() else {
709 return;
710 };
711 if !requests
712 .iter()
713 .all(|request| matches!(request, Some(InterpretationProgress::Completed(_))))
714 {
715 return;
716 }
717 let Some(InterpretationProgress::Interpreting(requests)) = progress.take() else {
718 return;
719 };
720 let mut remaining = requests.into_iter();
721 let mut settled = Vec::with_capacity(remaining.len());
722 while let Some(request) = remaining.next() {
723 match request {
724 Some(InterpretationProgress::Completed(Interpretation::Complete(received))) => {
725 settled.push(Interpretation::Complete(received))
726 }
727 Some(InterpretationProgress::Completed(Interpretation::Corrupt(received))) => {
728 settled.push(Interpretation::Corrupt(received));
729 }
730 request => {
731 *progress = Some(InterpretationProgress::Interpreting(
732 settled
733 .into_iter()
734 .map(|received| Some(InterpretationProgress::Completed(received)))
735 .chain(core::iter::once(request))
736 .chain(remaining)
737 .collect(),
738 ));
739 return;
740 }
741 }
742 }
743 let disposition = settled.iter().fold(
744 Interpretation::Complete(()),
745 |disposition, received| match (disposition, received) {
746 (Interpretation::Corrupt(()), _) | (_, Interpretation::Corrupt(_)) => {
747 Interpretation::Corrupt(())
748 }
749 (Interpretation::Complete(()), Interpretation::Complete(_)) => {
750 Interpretation::Complete(())
751 }
752 },
753 );
754 let requests = settled
755 .into_iter()
756 .map(Interpretation::into_settlement)
757 .collect();
758 let received = HeterogeneousShutdownSends { requests };
759 *progress = Some(InterpretationProgress::Completed(
760 disposition.map(|()| received),
761 ));
762 }
763 fn unattempted(
764 progress: &mut Option<
765 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
766 >,
767 ) {
768 <Self as behavior::SendSettlements>::prepare_interpretation(progress);
769 if let Some(InterpretationProgress::Interpreting(requests)) = progress {
770 for request in requests {
771 T::unattempted(request);
772 }
773 }
774 <Self as behavior::SendSettlements>::finish_interpretation(progress);
775 }
776}
777
778impl<Host, RootEvent, T> SourceSettlementCustody<Host, RootEvent> for HeterogeneousShutdownSends<T>
779where
780 T: Send,
781{
782 type Custody = Self;
783
784 fn prepare_source(progress: &mut Option<SourceProgress<Self, Self::Custody>>) {
785 *progress = match progress.take() {
786 Some(SourceProgress::Original(original)) => Some(SourceProgress::Completed(
787 SourceCustody::Exhausted(original),
788 )),
789 retained => retained,
790 };
791 }
792
793 fn offer_next_to_source(
794 _: &mut Self::Custody,
795 _: &mut Host,
796 ) -> impl core::future::Future<Output = ()> + Send {
797 core::future::ready(())
798 }
799
800 fn finish_source(progress: &mut Option<SourceProgress<Self, Self::Custody>>) {
801 *progress = match progress.take() {
802 Some(SourceProgress::Original(original) | SourceProgress::Offering(original)) => Some(
803 SourceProgress::Completed(SourceCustody::Exhausted(original)),
804 ),
805 retained => retained,
806 };
807 }
808}
809
810impl<I, E, Path, T> behavior::InterpretSends<I, E, Path> for HeterogeneousShutdownSends<T>
811where
812 I: Send,
813 T: heterogeneous::InterpretChoice<I, E, Path, ChildHead>,
814 T::InterpretationCustody: Send,
815{
816 fn interpret(
817 progress: &mut Option<
818 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
819 >,
820 interpreter: &mut I,
821 ) -> impl core::future::Future<Output = ()> + Send {
822 async move {
823 <Self as behavior::SendSettlements>::prepare_interpretation(progress);
824 let Some(InterpretationProgress::Interpreting(requests)) = progress else {
825 return;
826 };
827 let mut requests = requests.iter_mut();
828 while let Some(request) = requests.next() {
829 if !matches!(request, Some(InterpretationProgress::Completed(_))) {
830 T::settle(request, interpreter).await;
831 T::finish_interpretation(request);
832 }
833 match request {
834 Some(InterpretationProgress::Completed(Interpretation::Complete(_))) => {}
835 Some(InterpretationProgress::Completed(Interpretation::Corrupt(_))) => {
836 for untouched in requests {
837 T::unattempted(untouched);
838 }
839 break;
840 }
841 _ => return,
842 }
843 }
844 <Self as behavior::SendSettlements>::finish_interpretation(progress);
845 }
846 }
847}
848
849#[derive(Debug, Clone, PartialEq, Eq)]
856pub enum ShutdownState<P, N> {
857 AwaitingPlan,
859 AwaitingPlanAfterShutdown,
861 Ready { plan: P },
863 Stopping {
865 plan: P,
866 phase: usize,
867 awaiting: Vec<N>,
868 },
869 Completed,
871}
872
873enum ShutdownMove<P> {
874 None,
875 StartPhase { plan: P, phase: usize },
876 Stop,
877}
878
879pub struct InstallShutdownPlan<P> {
906 plan: P,
907}
908
909impl<P> InstallShutdownPlan<P> {
910 #[must_use]
912 pub const fn new(plan: P) -> Self {
913 Self { plan }
914 }
915
916 #[must_use]
918 pub fn into_plan(self) -> P {
919 self.plan
920 }
921}
922
923pub struct ReportShutdownPlan<P> {
928 installation: InstallShutdownPlan<P>,
929}
930
931impl<P: core::fmt::Debug> core::fmt::Debug for ReportShutdownPlan<P> {
932 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
933 formatter
934 .debug_struct("ReportShutdownPlan")
935 .field("plan", &self.installation.plan)
936 .finish()
937 }
938}
939
940impl<P: PartialEq> PartialEq for ReportShutdownPlan<P> {
941 fn eq(&self, other: &Self) -> bool {
942 self.installation.plan == other.installation.plan
943 }
944}
945
946impl<P: Eq> Eq for ReportShutdownPlan<P> {}
947
948impl<P> ReportShutdownPlan<P> {
949 #[must_use]
950 pub fn new(plan: P) -> Self {
951 Self {
952 installation: InstallShutdownPlan::new(plan),
953 }
954 }
955
956 #[must_use]
958 pub const fn plan(&self) -> &P {
959 &self.installation.plan
960 }
961
962 #[must_use]
964 pub fn into_event<Event>(self) -> Event
965 where
966 Event: EventIngress<Here, InstallShutdownPlan<P>>,
967 {
968 Event::ingress(self.installation)
969 }
970}
971
972impl<P> behavior::InterpreterRequest for ReportShutdownPlan<P> {
973 type ReturnToEmitter = behavior::NoReturnToEmitter;
974 type LogicalProtocols = behavior::NoBirthProtocols;
975}
976
977impl<P: Send> behavior::ActionItem for ReportShutdownPlan<P> {
978 type Custody = (Option<Self>, Option<Self::Reply>);
979 type Input<'a>
980 = &'a mut Option<Self>
981 where
982 Self: 'a;
983 type Reply = ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>;
984
985 fn prepare_interpretation(
986 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
987 ) {
988 prepare_item::<Self>(progress);
989 }
990
991 fn interpretation_input<'a>(
992 custody: &'a mut Self::Custody,
993 ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
994 where
995 Self: 'a,
996 {
997 match custody {
998 (input @ Some(_), received @ None) => Some((input, received)),
999 _ => None,
1000 }
1001 }
1002
1003 fn finish_interpretation(
1004 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
1005 ) {
1006 finish_item::<Self>(progress);
1007 }
1008
1009 type Accepted = ();
1010 type Rejection = behavior::Never;
1011 type Prerequisite = behavior::Never;
1012}
1013
1014pub enum ShutdownCoordinatorEvent<E: UserEvent, P> {
1016 Behavior(E),
1017 Plan(InstallShutdownPlan<P>),
1018 Requested(crate::ShutdownRequested),
1019 ChildStopped(ChildStopped<E::Addr>),
1020 ChildRejected(ChildShutdownRejected),
1021}
1022
1023impl<E: UserEvent, P> UserEvent for ShutdownCoordinatorEvent<E, P> {
1024 type Addr = E::Addr;
1025 type Message = E::Message;
1026 fn user(from: Self::Addr, message: Self::Message) -> Self {
1027 Self::Behavior(E::user(from, message))
1028 }
1029 fn into_user(self) -> Result<User<Self::Addr, Self::Message>, Self> {
1030 match self {
1031 Self::Behavior(e) => e.into_user().map_err(Self::Behavior),
1032 other => Err(other),
1033 }
1034 }
1035}
1036
1037impl<E: UserEvent, P> behavior::ComposedEvent for ShutdownCoordinatorEvent<E, P> {
1038 type Inner = E;
1039
1040 fn from_inner(event: E) -> Self {
1041 Self::Behavior(event)
1042 }
1043}
1044
1045impl<E: UserEvent, P> InjectEvent<InstallShutdownPlan<P>, Here> for ShutdownCoordinatorEvent<E, P> {
1046 fn inject_at(value: InstallShutdownPlan<P>) -> Self {
1047 Self::Plan(value)
1048 }
1049}
1050
1051impl<E: UserEvent, P> EventIngress<Here, InstallShutdownPlan<P>>
1052 for ShutdownCoordinatorEvent<E, P>
1053{
1054 fn ingress(value: InstallShutdownPlan<P>) -> Self {
1055 Self::Plan(value)
1056 }
1057}
1058
1059impl<E: UserEvent, P> InjectEvent<crate::ShutdownRequested, Here>
1060 for ShutdownCoordinatorEvent<E, P>
1061{
1062 fn inject_at(value: crate::ShutdownRequested) -> Self {
1063 Self::Requested(value)
1064 }
1065}
1066impl<E: UserEvent, P> InjectEvent<ChildStopped<E::Addr>, Here> for ShutdownCoordinatorEvent<E, P> {
1067 fn inject_at(value: ChildStopped<E::Addr>) -> Self {
1068 Self::ChildStopped(value)
1069 }
1070}
1071impl<E: UserEvent, P> InjectEvent<ChildShutdownRejected, Here> for ShutdownCoordinatorEvent<E, P> {
1072 fn inject_at(value: ChildShutdownRejected) -> Self {
1073 Self::ChildRejected(value)
1074 }
1075}
1076
1077impl<E, P, Input, Path> InjectEvent<Input, Inside<Path>> for ShutdownCoordinatorEvent<E, P>
1078where
1079 E: UserEvent + InjectEvent<Input, Path>,
1080{
1081 fn inject_at(input: Input) -> Self {
1082 Self::Behavior(E::inject_at(input))
1083 }
1084}
1085
1086#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
1088pub enum ShutdownCoordinatorError<E, A: Address, P> {
1089 #[error("wrapped behavior rejected its transition")]
1090 Behavior(#[source] E),
1091 #[error("child shutdown was rejected")]
1092 ChildRejected {
1093 child: CreationId,
1094 reason: ChildShutdownRejection,
1095 },
1096 #[error("child-stop fact does not belong to the active shutdown phase")]
1098 UnexpectedChildStopped(ChildStopped<A>),
1099 #[error("child-shutdown rejection does not belong to the active shutdown phase")]
1101 UnexpectedChildRejection {
1102 child: CreationId,
1103 reason: ChildShutdownRejection,
1104 },
1105 #[error("shutdown plan was already installed")]
1107 PlanAlreadyInstalled(P),
1108}
1109
1110pub struct ShutdownCoordinator<B: Behavior, C: Behavior, Occurrence>
1154where
1155 C::Protocol: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
1156{
1157 inner: B,
1158 state: ShutdownState<ShutdownPlan<CreationId>, CreationId>,
1159 child: core::marker::PhantomData<fn() -> (C, Occurrence)>,
1160}
1161
1162type ShutdownCoordinatorActions<B, C, Occurrence> = Actions<
1163 behavior::BehaviorAddr<B>,
1164 <B as Behavior>::Ph,
1165 SendLayer<InterpreterRequests<ShutdownChild<C, Occurrence>>, <B as Behavior>::Sends>,
1166 <B as Behavior>::Birth,
1167>;
1168
1169impl<B: Behavior, C: Behavior, Occurrence> ShutdownCoordinator<B, C, Occurrence>
1170where
1171 C::Protocol: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
1172{
1173 #[must_use]
1174 pub const fn new(inner: B, plan: ShutdownPlan<CreationId>) -> Self {
1175 Self {
1176 inner,
1177 state: ShutdownState::Ready { plan },
1178 child: core::marker::PhantomData,
1179 }
1180 }
1181
1182 #[must_use]
1188 pub const fn awaiting_plan(inner: B) -> Self {
1189 Self {
1190 inner,
1191 state: ShutdownState::AwaitingPlan,
1192 child: core::marker::PhantomData,
1193 }
1194 }
1195
1196 #[must_use]
1197 pub fn state(&self) -> &ShutdownState<ShutdownPlan<CreationId>, CreationId> {
1198 &self.state
1199 }
1200
1201 fn wrap(
1202 actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
1203 ) -> ShutdownCoordinatorActions<B, C, Occurrence> {
1204 actions.map_sends(|inner| SendLayer::new(InterpreterRequests::empty(), inner))
1205 }
1206
1207 fn phase_actions(
1208 plan: &ShutdownPlan<CreationId>,
1209 phase: usize,
1210 ) -> ShutdownCoordinatorActions<B, C, Occurrence> {
1211 let shutdowns = InterpreterRequests::new(
1212 plan.phases[phase]
1213 .iter()
1214 .copied()
1215 .map(ShutdownChild::<C, Occurrence>::new)
1216 .collect(),
1217 );
1218 Actions::send(SendLayer::new(shutdowns, B::Sends::empty()))
1219 }
1220
1221 fn start_plan(
1222 &mut self,
1223 plan: ShutdownPlan<CreationId>,
1224 ) -> ShutdownMove<ShutdownPlan<CreationId>> {
1225 if plan.phases.is_empty() {
1226 self.state = ShutdownState::Completed;
1227 return ShutdownMove::Stop;
1228 }
1229 let awaiting = plan.phases[0].clone();
1230 let selected = plan.clone();
1231 self.state = ShutdownState::Stopping {
1232 plan,
1233 phase: 0,
1234 awaiting,
1235 };
1236 ShutdownMove::StartPhase {
1237 plan: selected,
1238 phase: 0,
1239 }
1240 }
1241
1242 fn install_plan(
1243 &mut self,
1244 plan: ShutdownPlan<CreationId>,
1245 ) -> Result<ShutdownMove<ShutdownPlan<CreationId>>, ShutdownPlan<CreationId>> {
1246 match self.state {
1247 ShutdownState::AwaitingPlan => {
1248 self.state = ShutdownState::Ready { plan };
1249 Ok(ShutdownMove::None)
1250 }
1251 ShutdownState::AwaitingPlanAfterShutdown => Ok(self.start_plan(plan)),
1252 ShutdownState::Ready { .. }
1253 | ShutdownState::Stopping { .. }
1254 | ShutdownState::Completed => Err(plan),
1255 }
1256 }
1257
1258 fn request_shutdown(&mut self) -> ShutdownMove<ShutdownPlan<CreationId>> {
1259 let ready = match &self.state {
1260 ShutdownState::AwaitingPlan => {
1261 self.state = ShutdownState::AwaitingPlanAfterShutdown;
1262 None
1263 }
1264 ShutdownState::Ready { plan } => Some(plan.clone()),
1265 ShutdownState::AwaitingPlanAfterShutdown
1266 | ShutdownState::Stopping { .. }
1267 | ShutdownState::Completed => None,
1268 };
1269 ready.map_or(ShutdownMove::None, |plan| self.start_plan(plan))
1270 }
1271
1272 fn child_stopped(&mut self, child: CreationId) -> ShutdownMove<ShutdownPlan<CreationId>> {
1273 let ShutdownState::Stopping {
1274 plan,
1275 phase,
1276 awaiting,
1277 } = &mut self.state
1278 else {
1279 return ShutdownMove::None;
1280 };
1281 let Some(position) = awaiting.iter().position(|candidate| *candidate == child) else {
1282 return ShutdownMove::None;
1283 };
1284 awaiting.remove(position);
1285 if !awaiting.is_empty() {
1286 return ShutdownMove::None;
1287 }
1288 let next = *phase + 1;
1289 if next == plan.phases.len() {
1290 self.state = ShutdownState::Completed;
1291 ShutdownMove::Stop
1292 } else {
1293 *phase = next;
1294 *awaiting = plan.phases[next].clone();
1295 ShutdownMove::StartPhase {
1296 plan: plan.clone(),
1297 phase: next,
1298 }
1299 }
1300 }
1301
1302 fn move_actions(
1303 next: ShutdownMove<ShutdownPlan<CreationId>>,
1304 ) -> ShutdownCoordinatorActions<B, C, Occurrence> {
1305 match next {
1306 ShutdownMove::None => Actions::cont(),
1307 ShutdownMove::Stop => Actions::stop(),
1308 ShutdownMove::StartPhase { plan, phase } => Self::phase_actions(&plan, phase),
1309 }
1310 }
1311}
1312
1313impl<B, C, Occurrence> behavior::BehaviorBase for ShutdownCoordinator<B, C, Occurrence>
1314where
1315 B: Behavior + behavior::BehaviorBase,
1316 C: Behavior,
1317 C::Protocol: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
1318{
1319 type Base = B::Base;
1320 fn base(&self) -> &Self::Base {
1321 self.inner.base()
1322 }
1323}
1324
1325impl<B, C, Occurrence> crate::StashStatus for ShutdownCoordinator<B, C, Occurrence>
1326where
1327 B: Behavior + crate::StashStatus,
1328 C: Behavior,
1329 C::Protocol: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
1330{
1331 fn stashed_messages(&self) -> usize {
1332 self.inner.stashed_messages()
1333 }
1334}
1335
1336impl<B, C, Occurrence, A, Ph, S, Br> Behavior for ShutdownCoordinator<B, C, Occurrence>
1337where
1338 A: Address,
1339 S: SendEffects + behavior::SendsFor<B::Event>,
1340 Br: BirthMode,
1341 B: Behavior<Ph = Ph, Sends = S, Birth = Br>,
1342 B::Protocol: behavior::Protocol<Addr = A>,
1343 C: Behavior,
1344 C::Protocol: behavior::Protocol<Addr = A>,
1345 C::Event: InjectEvent<crate::ShutdownRequested, Here>,
1346{
1347 type Protocol = B::Protocol;
1348 type Event = ShutdownCoordinatorEvent<B::Event, ShutdownPlan<CreationId>>;
1349 type Sends = SendLayer<InterpreterRequests<ShutdownChild<C, Occurrence>>, S>;
1350 type Ph = Ph;
1351 type Error = ShutdownCoordinatorError<B::Error, A, ShutdownPlan<CreationId>>;
1352 type Birth = Br;
1353 fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
1354 behavior::initialize(&mut self.inner)
1355 .map(Self::wrap)
1356 .map_err(ShutdownCoordinatorError::Behavior)
1357 }
1358 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
1359 match event {
1360 ShutdownCoordinatorEvent::Plan(installation) => {
1361 let next = self
1362 .install_plan(installation.into_plan())
1363 .map_err(ShutdownCoordinatorError::PlanAlreadyInstalled)?;
1364 Ok(Self::move_actions(next))
1365 }
1366 ShutdownCoordinatorEvent::Requested(_) => {
1367 let next = self.request_shutdown();
1368 Ok(Self::move_actions(next))
1369 }
1370 ShutdownCoordinatorEvent::ChildStopped(stopped) => {
1371 let matching = matches!(&self.state, ShutdownState::Stopping { awaiting, .. } if awaiting.contains(&stopped.child));
1372 if !matching {
1373 return Err(ShutdownCoordinatorError::UnexpectedChildStopped(stopped));
1374 }
1375 let next = self.child_stopped(stopped.child);
1376 Ok(Self::move_actions(next))
1377 }
1378 ShutdownCoordinatorEvent::ChildRejected(rejected) => {
1379 let matching = matches!(&self.state, ShutdownState::Stopping { awaiting, .. } if awaiting.contains(&rejected.child));
1380 if matching {
1381 Err(ShutdownCoordinatorError::ChildRejected {
1382 child: rejected.child,
1383 reason: rejected.reason,
1384 })
1385 } else {
1386 Err(ShutdownCoordinatorError::UnexpectedChildRejection {
1387 child: rejected.child,
1388 reason: rejected.reason,
1389 })
1390 }
1391 }
1392 ShutdownCoordinatorEvent::Behavior(inner) => {
1393 behavior::delegate_transition(&mut self.inner, inner)
1394 .map(Self::wrap)
1395 .map_err(ShutdownCoordinatorError::Behavior)
1396 }
1397 }
1398 }
1399}
1400
1401pub struct HeterogeneousShutdownCoordinator<B: Behavior, T>
1408where
1409 T: heterogeneous::Selection<Addr = behavior::BehaviorAddr<B>>,
1410{
1411 inner: B,
1412 state: ShutdownState<HeterogeneousShutdownPlan<T>, CreationId>,
1413}
1414
1415type HeterogeneousShutdownActions<B, T> = Actions<
1416 behavior::BehaviorAddr<B>,
1417 <B as Behavior>::Ph,
1418 SendLayer<HeterogeneousShutdownSends<T>, <B as Behavior>::Sends>,
1419 <B as Behavior>::Birth,
1420>;
1421
1422impl<B: Behavior, T> HeterogeneousShutdownCoordinator<B, T>
1423where
1424 T: heterogeneous::Selection<Addr = behavior::BehaviorAddr<B>> + Copy,
1425{
1426 #[must_use]
1427 pub const fn new(inner: B, plan: HeterogeneousShutdownPlan<T>) -> Self {
1428 Self {
1429 inner,
1430 state: ShutdownState::Ready { plan },
1431 }
1432 }
1433
1434 #[must_use]
1436 pub const fn awaiting_plan(inner: B) -> Self {
1437 Self {
1438 inner,
1439 state: ShutdownState::AwaitingPlan,
1440 }
1441 }
1442
1443 #[must_use]
1444 pub fn state(&self) -> &ShutdownState<HeterogeneousShutdownPlan<T>, CreationId> {
1445 &self.state
1446 }
1447
1448 fn wrap(
1449 actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
1450 ) -> HeterogeneousShutdownActions<B, T> {
1451 actions.map_sends(|inner| SendLayer::new(HeterogeneousShutdownSends::empty(), inner))
1452 }
1453
1454 fn phase_actions(
1455 plan: &HeterogeneousShutdownPlan<T>,
1456 phase: usize,
1457 ) -> HeterogeneousShutdownActions<B, T> {
1458 let sends = HeterogeneousShutdownSends {
1459 requests: plan.phases[phase].clone(),
1460 };
1461 Actions::send(SendLayer::new(sends, B::Sends::empty()))
1462 }
1463
1464 fn phase_creations(plan: &HeterogeneousShutdownPlan<T>, phase: usize) -> Vec<CreationId> {
1465 plan.phases[phase]
1466 .iter()
1467 .map(heterogeneous::Selection::creation)
1468 .collect()
1469 }
1470
1471 fn start_plan(
1472 &mut self,
1473 plan: HeterogeneousShutdownPlan<T>,
1474 ) -> ShutdownMove<HeterogeneousShutdownPlan<T>> {
1475 if plan.phases.is_empty() {
1476 self.state = ShutdownState::Completed;
1477 return ShutdownMove::Stop;
1478 }
1479 let awaiting = Self::phase_creations(&plan, 0);
1480 let selected = plan.clone();
1481 self.state = ShutdownState::Stopping {
1482 plan,
1483 phase: 0,
1484 awaiting,
1485 };
1486 ShutdownMove::StartPhase {
1487 plan: selected,
1488 phase: 0,
1489 }
1490 }
1491
1492 fn install_plan(
1493 &mut self,
1494 plan: HeterogeneousShutdownPlan<T>,
1495 ) -> Result<ShutdownMove<HeterogeneousShutdownPlan<T>>, HeterogeneousShutdownPlan<T>> {
1496 match self.state {
1497 ShutdownState::AwaitingPlan => {
1498 self.state = ShutdownState::Ready { plan };
1499 Ok(ShutdownMove::None)
1500 }
1501 ShutdownState::AwaitingPlanAfterShutdown => Ok(self.start_plan(plan)),
1502 ShutdownState::Ready { .. }
1503 | ShutdownState::Stopping { .. }
1504 | ShutdownState::Completed => Err(plan),
1505 }
1506 }
1507
1508 fn request_shutdown(&mut self) -> ShutdownMove<HeterogeneousShutdownPlan<T>> {
1509 let ready = match &self.state {
1510 ShutdownState::AwaitingPlan => {
1511 self.state = ShutdownState::AwaitingPlanAfterShutdown;
1512 None
1513 }
1514 ShutdownState::Ready { plan } => Some(plan.clone()),
1515 ShutdownState::AwaitingPlanAfterShutdown
1516 | ShutdownState::Stopping { .. }
1517 | ShutdownState::Completed => None,
1518 };
1519 ready.map_or(ShutdownMove::None, |plan| self.start_plan(plan))
1520 }
1521
1522 fn child_stopped(&mut self, child: CreationId) -> ShutdownMove<HeterogeneousShutdownPlan<T>> {
1523 let ShutdownState::Stopping {
1524 plan,
1525 phase,
1526 awaiting,
1527 } = &mut self.state
1528 else {
1529 return ShutdownMove::None;
1530 };
1531 let Some(position) = awaiting.iter().position(|candidate| *candidate == child) else {
1532 return ShutdownMove::None;
1533 };
1534 awaiting.remove(position);
1535 if !awaiting.is_empty() {
1536 return ShutdownMove::None;
1537 }
1538 let next = *phase + 1;
1539 if next == plan.phases.len() {
1540 self.state = ShutdownState::Completed;
1541 ShutdownMove::Stop
1542 } else {
1543 *phase = next;
1544 *awaiting = Self::phase_creations(plan, next);
1545 ShutdownMove::StartPhase {
1546 plan: plan.clone(),
1547 phase: next,
1548 }
1549 }
1550 }
1551
1552 fn move_actions(
1553 next: ShutdownMove<HeterogeneousShutdownPlan<T>>,
1554 ) -> HeterogeneousShutdownActions<B, T> {
1555 match next {
1556 ShutdownMove::None => Actions::cont(),
1557 ShutdownMove::Stop => Actions::stop(),
1558 ShutdownMove::StartPhase { plan, phase } => Self::phase_actions(&plan, phase),
1559 }
1560 }
1561}
1562
1563impl<B, T> behavior::BehaviorBase for HeterogeneousShutdownCoordinator<B, T>
1564where
1565 B: Behavior + behavior::BehaviorBase,
1566 T: heterogeneous::Selection<Addr = behavior::BehaviorAddr<B>>,
1567{
1568 type Base = B::Base;
1569 fn base(&self) -> &Self::Base {
1570 self.inner.base()
1571 }
1572}
1573
1574impl<B, T> crate::StashStatus for HeterogeneousShutdownCoordinator<B, T>
1575where
1576 B: Behavior + crate::StashStatus,
1577 T: heterogeneous::Selection<Addr = behavior::BehaviorAddr<B>>,
1578{
1579 fn stashed_messages(&self) -> usize {
1580 self.inner.stashed_messages()
1581 }
1582}
1583
1584impl<B, T, A, Ph, Sends, Br> Behavior for HeterogeneousShutdownCoordinator<B, T>
1585where
1586 A: Address,
1587 Sends: SendEffects + behavior::SendsFor<B::Event>,
1588 Br: BirthMode,
1589 B: Behavior<Ph = Ph, Sends = Sends, Birth = Br>,
1590 B::Protocol: behavior::Protocol<Addr = A>,
1591 T: heterogeneous::Selection<Addr = A> + Copy,
1592{
1593 type Protocol = B::Protocol;
1594 type Event = ShutdownCoordinatorEvent<B::Event, HeterogeneousShutdownPlan<T>>;
1595 type Sends = SendLayer<HeterogeneousShutdownSends<T>, Sends>;
1596 type Ph = Ph;
1597 type Error = ShutdownCoordinatorError<B::Error, A, HeterogeneousShutdownPlan<T>>;
1598 type Birth = Br;
1599
1600 fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
1601 behavior::initialize(&mut self.inner)
1602 .map(Self::wrap)
1603 .map_err(ShutdownCoordinatorError::Behavior)
1604 }
1605
1606 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
1607 match event {
1608 ShutdownCoordinatorEvent::Plan(installation) => {
1609 let next = self
1610 .install_plan(installation.into_plan())
1611 .map_err(ShutdownCoordinatorError::PlanAlreadyInstalled)?;
1612 Ok(Self::move_actions(next))
1613 }
1614 ShutdownCoordinatorEvent::Requested(_) => {
1615 let next = self.request_shutdown();
1616 Ok(Self::move_actions(next))
1617 }
1618 ShutdownCoordinatorEvent::ChildStopped(stopped) => {
1619 let matching = matches!(&self.state, ShutdownState::Stopping { awaiting, .. } if awaiting.contains(&stopped.child));
1620 if !matching {
1621 return Err(ShutdownCoordinatorError::UnexpectedChildStopped(stopped));
1622 }
1623 let next = self.child_stopped(stopped.child);
1624 Ok(Self::move_actions(next))
1625 }
1626 ShutdownCoordinatorEvent::ChildRejected(rejected) => {
1627 let matching = matches!(&self.state, ShutdownState::Stopping { awaiting, .. } if awaiting.contains(&rejected.child));
1628 if matching {
1629 Err(ShutdownCoordinatorError::ChildRejected {
1630 child: rejected.child,
1631 reason: rejected.reason,
1632 })
1633 } else {
1634 Err(ShutdownCoordinatorError::UnexpectedChildRejection {
1635 child: rejected.child,
1636 reason: rejected.reason,
1637 })
1638 }
1639 }
1640 ShutdownCoordinatorEvent::Behavior(inner) => {
1641 behavior::delegate_transition(&mut self.inner, inner)
1642 .map(Self::wrap)
1643 .map_err(ShutdownCoordinatorError::Behavior)
1644 }
1645 }
1646 }
1647}
1648
1649#[cfg(test)]
1651mod tests {
1652 use core::future::Future;
1653 use std::time::Instant;
1654
1655 use super::heterogeneous::Selection;
1656 use super::*;
1657 use crate::Activate as _;
1658 use crate::{Exit, ShutdownRequested};
1659 use behavior::{CreationSequence, MailAddr, Never, NoBirths, NoSends, Step};
1660
1661 struct Probe;
1662
1663 impl behavior::BehaviorBase for Probe {
1664 type Base = Self;
1665 fn base(&self) -> &Self {
1666 self
1667 }
1668 }
1669
1670 impl behavior::Protocol for Probe {
1671 type Addr = MailAddr;
1672 type Msg = u8;
1673 }
1674
1675 impl Behavior for Probe {
1676 type Protocol = Self;
1677 type Event = User<MailAddr, u8>;
1678 type Sends = Vec<u8>;
1679 type Ph = Never;
1680 type Error = Never;
1681 type Birth = NoBirths;
1682 fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
1683 Ok(Actions::send(vec![1]))
1684 }
1685 fn transition(
1686 &mut self,
1687 _: behavior::ActiveTurn,
1688 event: Self::Event,
1689 ) -> BehaviorActed<Self> {
1690 Ok(Actions::send(vec![event.message]))
1691 }
1692 }
1693
1694 struct NamedParent;
1695
1696 #[behavior::behavior(
1697 addr = MailAddr,
1698 message = Never,
1699 births = {
1700 primary: crate::StopOnShutdown<Probe>,
1701 pool: crate::StopOnShutdown<Probe>,
1702 fallback: crate::StopOnShutdown<Probe>,
1703 },
1704 creation_settlements = retain_for_retirement,
1705 )]
1706 impl NamedParent {
1707 fn receive(&mut self, _: MailAddr, message: Never) -> BehaviorActed<Self> {
1708 match message {}
1709 }
1710 }
1711
1712 #[derive(Clone, Copy)]
1713 struct ApplicationCreations {
1714 primary: CreationId,
1715 pool: CreationId,
1716 fallback: CreationId,
1717 unrelated: CreationId,
1718 }
1719
1720 impl ApplicationCreations {
1721 fn issue() -> Self {
1722 let mut sequence = CreationSequence::new();
1723 let primary = sequence.issue().expect("the primary creation ID exists");
1724 let pool = sequence.issue().expect("the pool creation ID exists");
1725 let fallback = sequence.issue().expect("the fallback creation ID exists");
1726 let unrelated = sequence.issue().expect("the unrelated creation ID exists");
1727 Self {
1728 primary,
1729 pool,
1730 fallback,
1731 unrelated,
1732 }
1733 }
1734 }
1735
1736 fn stopped(child: CreationId) -> ChildStopped<MailAddr> {
1737 ChildStopped::new(child, Ok(Exit::Normal), Instant::now())
1738 }
1739
1740 #[test]
1741 fn plan_rejects_empty_phases_and_duplicate_children() {
1742 assert!(matches!(
1743 ShutdownPlan::<u64>::new([vec![]]),
1744 Err(ShutdownPlanError::EmptyPhase { phase: 0 })
1745 ));
1746 assert!(matches!(
1747 ShutdownPlan::new([vec![1, 2], vec![2]]),
1748 Err(ShutdownPlanError::DuplicateChild(2))
1749 ));
1750 }
1751
1752 #[test]
1753 fn tree_derives_stable_dependency_layers_and_rejects_invalid_graphs() {
1754 let tree = ShutdownTree::new([1, 2, 3, 4], [(1, 3), (2, 3), (3, 4)]).unwrap();
1755 let plan = tree.into_plan();
1756 assert_eq!(plan.phases(), &[vec![1, 2], vec![3], vec![4]]);
1757 assert!(matches!(
1758 ShutdownTree::new([1, 1], []),
1759 Err(ShutdownTreeError::DuplicateChild(1))
1760 ));
1761 assert!(matches!(
1762 ShutdownTree::new([1], [(1, 2)]),
1763 Err(ShutdownTreeError::UnknownChild(2))
1764 ));
1765 assert!(matches!(
1766 ShutdownTree::new([1, 2], [(1, 2), (2, 1)]),
1767 Err(ShutdownTreeError::Cycle)
1768 ));
1769 }
1770
1771 #[test]
1772 fn phases_advance_only_after_every_current_child_stops() {
1773 let children = ApplicationCreations::issue();
1774 let plan = ShutdownPlan::new([
1775 vec![children.primary, children.pool],
1776 vec![children.fallback],
1777 ])
1778 .unwrap();
1779 let initialized =
1780 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::new(Probe, plan)
1781 .initialize()
1782 .unwrap();
1783 assert_eq!(initialized.actions.sends.inner, [1]);
1784 assert!(initialized.actions.sends.owned.is_empty());
1785 let mut active = initialized.behavior;
1786
1787 let first = active.on_path(ShutdownRequested).unwrap();
1788 assert_eq!(
1789 first.sends.owned.as_slice(),
1790 [
1791 ShutdownChild::new(children.primary),
1792 ShutdownChild::new(children.pool),
1793 ]
1794 );
1795 assert!(matches!(
1796 active.state(),
1797 ShutdownState::Stopping {
1798 phase: 0,
1799 awaiting,
1800 ..
1801 } if awaiting == &[children.primary, children.pool]
1802 ));
1803
1804 let one = active.on_path(stopped(children.pool)).unwrap();
1805 assert_eq!(one.sends, SendLayer::empty());
1806 assert!(matches!(
1807 active.state(),
1808 ShutdownState::Stopping {
1809 phase: 0,
1810 awaiting,
1811 ..
1812 } if awaiting == &[children.primary]
1813 ));
1814
1815 let second = active.on_path(stopped(children.primary)).unwrap();
1816 assert_eq!(
1817 second.sends.owned.as_slice(),
1818 [ShutdownChild::new(children.fallback)]
1819 );
1820 assert!(matches!(
1821 active.state(),
1822 ShutdownState::Stopping {
1823 phase: 1,
1824 awaiting,
1825 ..
1826 } if awaiting == &[children.fallback]
1827 ));
1828
1829 let complete = active.on_path(stopped(children.fallback)).unwrap();
1830 assert!(matches!(complete.become_, Step::Stop(_)));
1831 assert_eq!(active.state(), &ShutdownState::Completed);
1832 }
1833
1834 #[test]
1835 fn heterogeneous_phases_preserve_cross_protocol_order_and_await_the_union() {
1836 type SupervisorChild = crate::StopOnShutdown<Probe>;
1837 type PoolChild = crate::StopOnShutdown<Probe>;
1838 type RootTargets = ShutdownChoice<
1839 SupervisorChild,
1840 ShutdownChoice<PoolChild, ShutdownChoice<SupervisorChild, NoShutdownTargets<MailAddr>>>,
1841 >;
1842 let children = ApplicationCreations::issue();
1843 let primary = shutdown_target::<NamedParent, _, RootTargets>(
1844 NamedParentChild::Primary,
1845 children.primary,
1846 );
1847 let pool =
1848 shutdown_target::<NamedParent, _, RootTargets>(NamedParentChild::Pool, children.pool);
1849 let fallback = shutdown_target::<NamedParent, _, RootTargets>(
1850 NamedParentChild::Fallback,
1851 children.fallback,
1852 );
1853
1854 assert!(matches!(
1855 primary,
1856 ShutdownChoice::Other(ShutdownChoice::Other(ShutdownChoice::Child {
1857 creation,
1858 ..
1859 })) if creation == children.primary
1860 ));
1861 assert!(matches!(
1862 pool,
1863 ShutdownChoice::Other(ShutdownChoice::Child { creation, .. })
1864 if creation == children.pool
1865 ));
1866 assert!(matches!(
1867 fallback,
1868 ShutdownChoice::Child { creation, .. } if creation == children.fallback
1869 ));
1870
1871 let plan = HeterogeneousShutdownPlan::new([vec![primary, pool], vec![fallback]]).unwrap();
1872 let mut active =
1873 HeterogeneousShutdownCoordinator::<NamedParent, RootTargets>::new(NamedParent, plan)
1874 .initialize()
1875 .unwrap()
1876 .behavior;
1877
1878 let first = active.on_path(ShutdownRequested).unwrap();
1879 assert_eq!(
1880 first
1881 .sends
1882 .owned
1883 .requests
1884 .iter()
1885 .map(|target| target.creation())
1886 .collect::<Vec<_>>(),
1887 [children.primary, children.pool]
1888 );
1889 assert!(matches!(first.sends.inner, NoSends));
1890 assert!(first.creates.is_empty());
1891 assert!(matches!(first.become_, Step::Continue));
1892 let retained = active.on_path(stopped(children.pool)).unwrap();
1893 assert!(retained.sends.owned.requests.is_empty());
1894 assert!(matches!(retained.sends.inner, NoSends));
1895 assert!(retained.creates.is_empty());
1896 assert!(matches!(retained.become_, Step::Continue));
1897 let second = active.on_path(stopped(children.primary)).unwrap();
1898 assert_eq!(second.sends.owned.requests[0].creation(), children.fallback);
1899 assert!(matches!(second.sends.inner, NoSends));
1900 assert!(second.creates.is_empty());
1901 assert!(matches!(second.become_, Step::Continue));
1902 let completed = active.on_path(stopped(children.fallback)).unwrap();
1903 assert!(completed.sends.owned.requests.is_empty());
1904 assert!(matches!(completed.sends.inner, NoSends));
1905 assert!(completed.creates.is_empty());
1906 assert!(matches!(completed.become_, Step::Stop(_)));
1907 }
1908
1909 #[test]
1910 fn heterogeneous_plan_rejects_cross_protocol_creation_collisions() {
1911 type SupervisorChild = crate::StopOnShutdown<Probe>;
1912 type PoolChild = crate::StopOnShutdown<Probe>;
1913 type RootTargets = ShutdownChoice<
1914 SupervisorChild,
1915 ShutdownChoice<PoolChild, ShutdownChoice<SupervisorChild, NoShutdownTargets<MailAddr>>>,
1916 >;
1917 let children = ApplicationCreations::issue();
1918 assert!(matches!(
1919 HeterogeneousShutdownPlan::new([vec![
1920 shutdown_target::<NamedParent, _, RootTargets>(
1921 NamedParentChild::Primary,
1922 children.primary,
1923 ),
1924 shutdown_target::<NamedParent, _, RootTargets>(
1925 NamedParentChild::Pool,
1926 children.primary,
1927 ),
1928 ]]),
1929 Err(ShutdownPlanError::DuplicateChild(creation)) if creation == children.primary
1930 ));
1931 }
1932
1933 #[tokio::test]
1934 async fn heterogeneous_requests_interpret_once_each_in_cross_protocol_plan_order() {
1935 type PrimaryWorker = crate::StopOnShutdown<Probe>;
1936 type PoolWorker = crate::StopOnShutdown<Probe>;
1937 type Targets = ShutdownChoice<
1938 PrimaryWorker,
1939 ShutdownChoice<PoolWorker, ShutdownChoice<PrimaryWorker, NoShutdownTargets<MailAddr>>>,
1940 >;
1941 type Event =
1942 ShutdownCoordinatorEvent<User<MailAddr, u8>, HeterogeneousShutdownPlan<Targets>>;
1943
1944 struct Recording(Vec<CreationId>);
1945 impl behavior::InterpretItem<ShutdownChild<PrimaryWorker, ChildHead>, Event, Here> for Recording {
1946 fn interpret_item<'a>(
1947 &'a mut self,
1948 input: &'a mut Option<ShutdownChild<PrimaryWorker, ChildHead>>,
1949 received: &'a mut Option<
1950 <ShutdownChild<PrimaryWorker, ChildHead> as ActionItem>::Reply,
1951 >,
1952 ) -> impl Future<Output = ()> + Send + 'a
1953 where
1954 ShutdownChild<PrimaryWorker, ChildHead>: 'a,
1955 {
1956 async move {
1957 if received.is_some() {
1958 return;
1959 }
1960 let Some(request) = input.take() else {
1961 return;
1962 };
1963 self.0.push(request.child);
1964 *received = Some(ItemSettlement::Accepted(()));
1965 }
1966 }
1967 }
1968 impl behavior::InterpretItem<ShutdownChild<PoolWorker, ChildTail<ChildHead>>, Event, Here>
1969 for Recording
1970 {
1971 fn interpret_item<'a>(
1972 &'a mut self,
1973 input: &'a mut Option<ShutdownChild<PoolWorker, ChildTail<ChildHead>>>,
1974 received: &'a mut Option<
1975 <ShutdownChild<PoolWorker, ChildTail<ChildHead>> as ActionItem>::Reply,
1976 >,
1977 ) -> impl Future<Output = ()> + Send + 'a
1978 where
1979 ShutdownChild<PoolWorker, ChildTail<ChildHead>>: 'a,
1980 {
1981 async move {
1982 if received.is_some() {
1983 return;
1984 }
1985 let Some(request) = input.take() else {
1986 return;
1987 };
1988 self.0.push(request.child);
1989 *received = Some(ItemSettlement::Accepted(()));
1990 }
1991 }
1992 }
1993 impl
1994 behavior::InterpretItem<
1995 ShutdownChild<PrimaryWorker, ChildTail<ChildTail<ChildHead>>>,
1996 Event,
1997 Here,
1998 > for Recording
1999 {
2000 fn interpret_item<'a>(
2001 &'a mut self,
2002 input: &'a mut Option<
2003 ShutdownChild<PrimaryWorker, ChildTail<ChildTail<ChildHead>>>,
2004 >,
2005 received: &'a mut Option<<ShutdownChild<PrimaryWorker, ChildTail<ChildTail<ChildHead>>> as ActionItem>::Reply>,
2006 ) -> impl Future<Output = ()> + Send + 'a
2007 where
2008 ShutdownChild<PrimaryWorker, ChildTail<ChildTail<ChildHead>>>: 'a,
2009 {
2010 async move {
2011 if received.is_some() {
2012 return;
2013 }
2014 let Some(request) = input.take() else {
2015 return;
2016 };
2017 self.0.push(request.child);
2018 *received = Some(ItemSettlement::Accepted(()));
2019 }
2020 }
2021 }
2022
2023 let children = ApplicationCreations::issue();
2024 let sends = HeterogeneousShutdownSends::<Targets> {
2025 requests: vec![
2026 shutdown_target::<NamedParent, _, Targets>(NamedParentChild::Pool, children.pool),
2027 shutdown_target::<NamedParent, _, Targets>(
2028 NamedParentChild::Fallback,
2029 children.fallback,
2030 ),
2031 shutdown_target::<NamedParent, _, Targets>(
2032 NamedParentChild::Primary,
2033 children.primary,
2034 ),
2035 ],
2036 };
2037 let mut interpreter = Recording(Vec::new());
2038 let mut progress = Some(InterpretationProgress::Original(sends));
2039 <_ as behavior::InterpretSends<_, Event, Here>>::interpret(&mut progress, &mut interpreter)
2040 .await;
2041 let Some(InterpretationProgress::Completed(settlement)) = progress else {
2042 panic!("the actual heterogeneous shutdown plan must finish each whole request");
2043 };
2044 assert!(matches!(settlement, behavior::Interpretation::Complete(_)));
2045 assert_eq!(
2046 interpreter.0,
2047 [children.pool, children.fallback, children.primary]
2048 );
2049 }
2050
2051 #[test]
2052 fn coordinator_routes_shutdown_to_children_before_root_stop() {
2053 let children = ApplicationCreations::issue();
2054 let plan = ShutdownPlan::new([vec![children.primary]]).unwrap();
2055 let initialized =
2056 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::new(Probe, plan)
2057 .initialize()
2058 .unwrap();
2059 let mut active = initialized.behavior;
2060
2061 let actions = active.on(ShutdownRequested).unwrap();
2062
2063 assert_eq!(
2064 actions.sends.owned.as_slice(),
2065 [ShutdownChild::new(children.primary)]
2066 );
2067 assert!(actions.sends.inner.is_empty());
2068 assert!(actions.creates.is_empty());
2069 assert!(matches!(actions.become_, Step::Continue));
2070 }
2071
2072 #[test]
2073 fn stale_child_stops_are_returned_while_repeated_shutdown_is_idempotent() {
2074 let children = ApplicationCreations::issue();
2075 let plan = ShutdownPlan::new([vec![children.primary, children.pool]]).unwrap();
2076 let mut active =
2077 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::new(Probe, plan)
2078 .initialize()
2079 .unwrap()
2080 .behavior;
2081 let before_start = stopped(children.unrelated);
2082 assert_eq!(
2083 active.on_path(before_start),
2084 Err(ShutdownCoordinatorError::UnexpectedChildStopped(
2085 before_start
2086 ))
2087 );
2088 let started = active.on_path(ShutdownRequested).unwrap();
2089 assert_eq!(started.sends.owned.as_slice().len(), 2);
2090 assert!(started.sends.inner.is_empty());
2091 assert!(started.creates.is_empty());
2092 assert!(matches!(started.become_, Step::Continue));
2093 let repeated = active.on_path(ShutdownRequested).unwrap();
2094 assert_eq!(repeated.sends, SendLayer::empty());
2095 assert!(repeated.creates.is_empty());
2096 assert!(matches!(repeated.become_, Step::Continue));
2097 let retained = active.on_path(stopped(children.primary)).unwrap();
2098 assert_eq!(retained.sends, SendLayer::empty());
2099 assert!(retained.creates.is_empty());
2100 assert!(matches!(retained.become_, Step::Continue));
2101 let duplicate = stopped(children.primary);
2102 assert_eq!(
2103 active.on_path(duplicate),
2104 Err(ShutdownCoordinatorError::UnexpectedChildStopped(duplicate))
2105 );
2106 assert!(matches!(
2107 active.state(),
2108 ShutdownState::Stopping {
2109 phase: 0,
2110 awaiting,
2111 ..
2112 } if awaiting == &[children.pool]
2113 ));
2114 }
2115
2116 #[test]
2117 fn matching_rejection_is_typed_and_does_not_mutate_phase() {
2118 let children = ApplicationCreations::issue();
2119 let plan = ShutdownPlan::new([vec![children.primary]]).unwrap();
2120 let mut active =
2121 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::new(Probe, plan)
2122 .initialize()
2123 .unwrap()
2124 .behavior;
2125 let started = active.on_path(ShutdownRequested).unwrap();
2126 assert_eq!(
2127 started.sends.owned.as_slice(),
2128 [ShutdownChild::new(children.primary)]
2129 );
2130 assert!(started.sends.inner.is_empty());
2131 assert!(started.creates.is_empty());
2132 assert!(matches!(started.become_, Step::Continue));
2133 let before = active.state().clone();
2134 assert_eq!(
2135 active.on_path(ChildShutdownRejected::new(
2136 children.primary,
2137 ChildShutdownRejection::NotEstablished
2138 )),
2139 Err(ShutdownCoordinatorError::ChildRejected {
2140 child: children.primary,
2141 reason: ChildShutdownRejection::NotEstablished
2142 })
2143 );
2144 assert_eq!(active.state(), &before);
2145 }
2146
2147 #[test]
2148 fn empty_plan_stops_immediately_and_user_actions_preserve_named_lanes() {
2149 let plan = ShutdownPlan::<CreationId>::new([]).unwrap();
2150 let mut active =
2151 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::new(Probe, plan)
2152 .initialize()
2153 .unwrap()
2154 .behavior;
2155 let user = active.receive(MailAddr(0), 7).unwrap();
2156 assert_eq!(user.sends.inner, [7]);
2157 assert!(user.sends.owned.is_empty());
2158 assert!(user.creates.is_empty());
2159 assert!(matches!(user.become_, Step::Continue));
2160 let stopped = active.on_path(ShutdownRequested).unwrap();
2161 assert!(stopped.sends.owned.is_empty());
2162 assert!(stopped.sends.inner.is_empty());
2163 assert!(stopped.creates.is_empty());
2164 assert!(matches!(stopped.become_, Step::Stop(_)));
2165 }
2166
2167 #[test]
2168 fn homogeneous_plan_acceptance_is_an_explicit_one_way_lifecycle() {
2169 let mut active =
2170 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::awaiting_plan(
2171 Probe,
2172 )
2173 .initialize()
2174 .unwrap()
2175 .behavior;
2176 assert!(matches!(active.state(), ShutdownState::AwaitingPlan));
2177
2178 let children = ApplicationCreations::issue();
2179 let plan = ShutdownPlan::new([
2180 vec![children.primary, children.pool],
2181 vec![children.fallback],
2182 ])
2183 .unwrap();
2184 let installed = active
2185 .on_path(InstallShutdownPlan::new(plan.clone()))
2186 .unwrap();
2187 assert_eq!(installed.sends, SendLayer::empty());
2188 assert!(matches!(
2189 active.state(),
2190 ShutdownState::Ready { plan: installed } if installed == &plan
2191 ));
2192 assert_eq!(
2193 active
2194 .on_path(InstallShutdownPlan::new(plan.clone()))
2195 .unwrap_err(),
2196 ShutdownCoordinatorError::PlanAlreadyInstalled(plan)
2197 );
2198 }
2199
2200 #[test]
2201 fn topology_owner_reports_a_plan_as_an_ordinary_typed_event() {
2202 type Event = ShutdownCoordinatorEvent<User<MailAddr, u8>, ShutdownPlan<CreationId>>;
2203
2204 let mut active =
2205 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::awaiting_plan(
2206 Probe,
2207 )
2208 .initialize()
2209 .unwrap()
2210 .behavior;
2211 let children = ApplicationCreations::issue();
2212 let plan = ShutdownPlan::new([vec![children.primary, children.pool]]).unwrap();
2213 let report: ReportShutdownPlan<_> = ReportShutdownPlan::new(plan.clone());
2214 let event: Event = report.into_event();
2215
2216 let installed = active.transition(event).unwrap();
2217
2218 assert_eq!(installed.sends, SendLayer::empty());
2219 assert!(matches!(
2220 active.state(),
2221 ShutdownState::Ready { plan: installed } if installed == &plan
2222 ));
2223 }
2224
2225 #[test]
2226 fn shutdown_before_homogeneous_plan_is_retained_and_empty_plan_stops() {
2227 let children = ApplicationCreations::issue();
2228 let mut active =
2229 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::awaiting_plan(
2230 Probe,
2231 )
2232 .initialize()
2233 .unwrap()
2234 .behavior;
2235 let waiting = active.on_path(ShutdownRequested).unwrap();
2236 assert_eq!(waiting.sends, SendLayer::empty());
2237 assert!(waiting.creates.is_empty());
2238 assert!(matches!(waiting.become_, Step::Continue));
2239 assert!(matches!(
2240 active.state(),
2241 ShutdownState::AwaitingPlanAfterShutdown
2242 ));
2243 let started = active
2244 .on_path(InstallShutdownPlan::new(
2245 ShutdownPlan::new([vec![children.primary, children.pool]]).unwrap(),
2246 ))
2247 .unwrap();
2248 assert_eq!(
2249 started.sends.owned.as_slice(),
2250 [
2251 ShutdownChild::new(children.primary),
2252 ShutdownChild::new(children.pool),
2253 ]
2254 );
2255 assert!(started.sends.inner.is_empty());
2256 assert!(started.creates.is_empty());
2257 assert!(matches!(started.become_, Step::Continue));
2258
2259 let mut empty =
2260 ShutdownCoordinator::<Probe, crate::StopOnShutdown<Probe>, ChildHead>::awaiting_plan(
2261 Probe,
2262 )
2263 .initialize()
2264 .unwrap()
2265 .behavior;
2266 let waiting = empty.on_path(ShutdownRequested).unwrap();
2267 assert_eq!(waiting.sends, SendLayer::empty());
2268 assert!(waiting.creates.is_empty());
2269 assert!(matches!(waiting.become_, Step::Continue));
2270 let stopped = empty
2271 .on_path(InstallShutdownPlan::new(ShutdownPlan::new([]).unwrap()))
2272 .unwrap();
2273 assert!(matches!(stopped.become_, Step::Stop(_)));
2274 assert_eq!(empty.state(), &ShutdownState::Completed);
2275 }
2276
2277 #[test]
2278 fn heterogeneous_plan_can_be_installed_after_exact_children_are_selected() {
2279 type PrimaryWorker = crate::StopOnShutdown<Probe>;
2280 type PoolWorker = crate::StopOnShutdown<Probe>;
2281 type Targets = ShutdownChoice<
2282 PrimaryWorker,
2283 ShutdownChoice<PoolWorker, ShutdownChoice<PrimaryWorker, NoShutdownTargets<MailAddr>>>,
2284 >;
2285
2286 let mut active =
2287 HeterogeneousShutdownCoordinator::<NamedParent, Targets>::awaiting_plan(NamedParent)
2288 .initialize()
2289 .unwrap()
2290 .behavior;
2291 let waiting = active.on_path(ShutdownRequested).unwrap();
2292 assert!(waiting.sends.owned.requests.is_empty());
2293 assert!(matches!(waiting.sends.inner, NoSends));
2294 assert!(waiting.creates.is_empty());
2295 assert!(matches!(waiting.become_, Step::Continue));
2296 let children = ApplicationCreations::issue();
2297 let plan = HeterogeneousShutdownPlan::new([vec![
2298 shutdown_target::<NamedParent, _, Targets>(NamedParentChild::Pool, children.pool),
2299 shutdown_target::<NamedParent, _, Targets>(NamedParentChild::Primary, children.primary),
2300 ]])
2301 .unwrap();
2302 let started = active.on_path(InstallShutdownPlan::new(plan)).unwrap();
2303 assert_eq!(
2304 started
2305 .sends
2306 .owned
2307 .requests
2308 .iter()
2309 .map(|target| target.creation())
2310 .collect::<Vec<_>>(),
2311 [children.pool, children.primary]
2312 );
2313 assert!(matches!(started.sends.inner, NoSends));
2314 assert!(started.creates.is_empty());
2315 assert!(matches!(started.become_, Step::Continue));
2316 assert!(matches!(
2317 active.state(),
2318 ShutdownState::Stopping {
2319 phase: 0,
2320 awaiting,
2321 ..
2322 } if awaiting == &[children.pool, children.primary]
2323 ));
2324 }
2325}