Skip to main content

behavior_actors/lifecycle/
shutdown_coordinator.rs

1//! Phased child-topology shutdown.
2
3use 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/// Validated ordered shutdown phases.
14#[derive(Debug, Clone, PartialEq, Eq)]
15pub struct ShutdownPlan<N> {
16    phases: Vec<Vec<N>>,
17}
18
19impl<N: Copy + Eq> ShutdownPlan<N> {
20    /// Validate non-empty phases and globally unique child IDs.
21    ///
22    /// # Errors
23    /// Returns the offending phase or duplicate child ID.
24    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/// Invalid static shutdown topology.
48#[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
56/// Dependency topology compiled into ordered shutdown phases.
57///
58/// Each `(dependent, dependency)` edge means the dependent must stop before
59/// the dependency. Independent nodes share a phase in declaration order.
60pub struct ShutdownTree<N> {
61    plan: ShutdownPlan<N>,
62}
63
64impl<N: Copy + Eq> ShutdownTree<N> {
65    /// Validate a closed acyclic topology and derive shutdown layers.
66    ///
67    /// # Errors
68    /// Returns duplicate/unknown child evidence or `Cycle`.
69    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/// Invalid dependency-ordered shutdown topology.
120#[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
130/// One member of a closed heterogeneous child-protocol sum.
131///
132/// A root gives the recursive sum a topology-specific alias. `Child` selects
133/// the protocol at this position; `Other` selects one of the remaining
134/// protocols. The value carries only the creator-local creation ID, never an erased
135/// actor, address, request, or runtime protocol key.
136///
137/// A child whose event algebra has no direct shutdown owner cannot enter a
138/// validated heterogeneous plan:
139///
140/// ```compile_fail
141/// struct Plain;
142/// impl behavior::Protocol for Plain { type Addr = behavior::MailAddr; type Msg = (); }
143/// impl behavior::Behavior for Plain {
144///     type Protocol = Self;
145///     type Event = behavior::User<behavior::MailAddr, ()>;
146///     type Sends = Vec<behavior::Never>;
147///     type Ph = behavior::Never;
148///     type Error = behavior::Never;
149///     type Birth = behavior::NoBirths;
150///     fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> behavior::BehaviorActed<Self> {
151///         Ok(behavior::Actions::cont())
152///     }
153/// }
154/// type Targets = behavior_actors::ShutdownChoice<Plain, behavior_actors::NoShutdownTargets<behavior::MailAddr>>;
155/// let child = behavior::CreationSequence::new()
156///     .issue()
157///     .expect("the first child creation ID exists");
158/// let _ = behavior_actors::HeterogeneousShutdownPlan::new([vec![Targets::child(child)]]);
159/// ```
160pub enum ShutdownChoice<C: Behavior, Tail> {
161    Child {
162        creation: CreationId,
163        child: core::marker::PhantomData<fn() -> C>,
164    },
165    Other(Tail),
166}
167
168/// Uninhabited end of a heterogeneous shutdown choice.
169pub 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
196/// Structural construction of an existing [`ShutdownChoice`] at `Position`.
197///
198/// Implementations preserve the creation ID exactly. `ChildHead` selects the
199/// current branch; `ChildTail<P>` delegates to the existing tail. This trait
200/// introduces no alternate target product, runtime lookup, or shutdown effect.
201pub trait ShutdownTargetAt<Child: Behavior, Position>: named_shutdown::Target + Sized {
202    /// Lower one creator-local child ID into its statically selected branch.
203    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/// Lower one Behavior-owned named child ID into an existing heterogeneous
231/// shutdown target sum.
232///
233/// `Parent` fixes the [`ChildRole`] implementation, allowing the compiler to
234/// select the exact structural position even when several roles share one
235/// child behavior type. An unrelated role cannot select a target. The creation
236/// ID remains opaque creator-local correlation; it is not statically branded
237/// with the selected role.
238///
239/// ```compile_fail
240/// struct Worker;
241/// #[behavior::behavior(
242///     addr = behavior::MailAddr,
243///     message = behavior::Never,
244/// )]
245/// impl Worker {
246///     fn receive(
247///         &mut self,
248///         _: behavior::MailAddr,
249///         message: behavior::Never,
250///     ) -> behavior::BehaviorActed<Self> {
251///         match message {}
252///     }
253/// }
254/// struct Parent;
255/// #[behavior::behavior(
256///     addr = behavior::MailAddr,
257///     message = behavior::Never,
258///     births = { primary: Worker, fallback: Worker },
259///     creation_settlements = retain_for_retirement,
260/// )]
261/// impl Parent {
262///     fn receive(
263///         &mut self,
264///         _: behavior::MailAddr,
265///         message: behavior::Never,
266///     ) -> behavior::BehaviorActed<Self> {
267///         match message {}
268///     }
269/// }
270/// struct UnrelatedRole;
271/// type Targets = behavior_actors::ShutdownChoice<
272///     Worker,
273///     behavior_actors::ShutdownChoice<
274///         Worker,
275///         behavior_actors::NoShutdownTargets<behavior::MailAddr>,
276///     >,
277/// >;
278/// let mut creations = behavior::CreationSequence::new();
279/// let child = creations.issue().expect("fixture creation ID");
280/// let _: Targets =
281///     behavior_actors::shutdown_target::<Parent, _, Targets>(UnrelatedRole, child);
282/// ```
283#[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/// Settlement shape of one closed heterogeneous shutdown choice.
313#[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/// Validated shutdown phases over an arbitrary closed child-protocol sum.
591#[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    /// Validate non-empty phases and global uniqueness in the creator's one
610    /// child namespace, including collisions across protocol lanes.
611    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
637/// Ordered heterogeneous shutdown requests. Static dispatch occurs per item,
638/// preserving phase declaration order across protocol alternatives.
639pub 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    /// Borrow the phase-ordered shutdown selections emitted by this turn.
657    ///
658    /// The order is the declaration order of the active phase. Retained
659    /// selections from later phases are not exposed until that phase starts.
660    #[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/// Complete phase of coordinated shutdown, including plan installation.
850///
851/// The installed plan remains owned by the exact phase in which it is valid.
852/// A shutdown request received before installation has its own state and is
853/// discharged immediately when a plan later arrives. No correlated readiness
854/// flag or optional plan can describe a contradictory combination.
855#[derive(Debug, Clone, PartialEq, Eq)]
856pub enum ShutdownState<P, N> {
857    /// Child establishment has not yet produced the plan.
858    AwaitingPlan,
859    /// Shutdown was requested while child establishment was incomplete.
860    AwaitingPlanAfterShutdown,
861    /// The validated plan is installed and shutdown has not been requested.
862    Ready { plan: P },
863    /// One plan phase is active and owns its outstanding child IDs.
864    Stopping {
865        plan: P,
866        phase: usize,
867        awaiting: Vec<N>,
868    },
869    /// Every configured phase completed, or the installed plan was empty.
870    Completed,
871}
872
873enum ShutdownMove<P> {
874    None,
875    StartPhase { plan: P, phase: usize },
876    Stop,
877}
878
879/// Install one validated plan into a coordinator that was started without it.
880///
881/// The plan type is part of the coordinator's event sum. Homogeneous and
882/// heterogeneous plans therefore cannot be confused at installation:
883///
884/// ```compile_fail
885/// struct Probe;
886/// impl behavior::Protocol for Probe { type Addr = behavior::MailAddr; type Msg = (); }
887/// impl behavior::Behavior for Probe {
888///     type Protocol = Self;
889///     type Event = behavior::User<behavior::MailAddr, ()>;
890///     type Sends = Vec<behavior::Never>;
891///     type Ph = behavior::Never;
892///     type Error = behavior::Never;
893///     type Birth = behavior::NoBirths;
894///     fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> behavior::BehaviorActed<Self> {
895///         Ok(behavior::Actions::cont())
896///     }
897/// }
898/// type Targets = behavior_actors::ShutdownChoice<behavior_actors::StopOnShutdown<Probe>, behavior_actors::NoShutdownTargets<behavior::MailAddr>>;
899/// let mut coordinator = behavior_actors::Activate::initialize(
900///     behavior_actors::ShutdownCoordinator::<Probe, behavior_actors::StopOnShutdown<Probe>, behavior::ChildHead>::awaiting_plan(Probe)
901/// ).unwrap().behavior;
902/// let heterogeneous = behavior_actors::HeterogeneousShutdownPlan::<Targets>::new([]).unwrap();
903/// coordinator.on_path(behavior_actors::InstallShutdownPlan::new(heterogeneous)).unwrap();
904/// ```
905pub struct InstallShutdownPlan<P> {
906    plan: P,
907}
908
909impl<P> InstallShutdownPlan<P> {
910    /// Construct the explicit plan-installation input.
911    #[must_use]
912    pub const fn new(plan: P) -> Self {
913        Self { plan }
914    }
915
916    /// Consume the input into the validated plan it owns.
917    #[must_use]
918    pub fn into_plan(self) -> P {
919        self.plan
920    }
921}
922
923/// Interpreter request reporting one validated plan to its owning coordinator.
924///
925/// The request owns the plan. Interpretation enqueues one ordinary event for
926/// the same actor incarnation through its source-indexed [`EventIngress`].
927pub 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    /// Borrow the complete validated plan without changing the request.
957    #[must_use]
958    pub const fn plan(&self) -> &P {
959        &self.installation.plan
960    }
961
962    /// Build the exact root event selected for the current actor source.
963    #[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
1014/// Event sum accepted by [`ShutdownCoordinator`].
1015pub 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/// Controlled coordinated-shutdown failure.
1087#[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    /// A child-stop fact did not name a child awaited by the current phase.
1097    #[error("child-stop fact does not belong to the active shutdown phase")]
1098    UnexpectedChildStopped(ChildStopped<A>),
1099    /// A child-shutdown rejection did not name a child awaited by the current phase.
1100    #[error("child-shutdown rejection does not belong to the active shutdown phase")]
1101    UnexpectedChildRejection {
1102        child: CreationId,
1103        reason: ChildShutdownRejection,
1104    },
1105    /// A plan was supplied after one had already been installed or completed.
1106    #[error("shutdown plan was already installed")]
1107    PlanAlreadyInstalled(P),
1108}
1109
1110/// Pure phased shutdown wrapper over an explicitly validated homogeneous child
1111/// topology.
1112///
1113/// `B` is the wrapped coordinator behavior and `C` is the one concrete child
1114/// protocol selected by every creation ID in the plan. Starting a phase emits one
1115/// typed [`ShutdownChild<C, Occurrence>`] request per member in plan order. A phase advances
1116/// only after every matching [`ChildStopped`] fact arrives. An acceptance
1117/// rejection returns [`ShutdownCoordinatorError::ChildRejected`] without
1118/// changing the phase. A stale or foreign child fact is returned intact as a
1119/// typed error. This is Bombay lifecycle policy, not an actor-model allocation
1120/// or ordering guarantee.
1121///
1122/// The fold introduces no panic conditions.
1123///
1124/// A child protocol without the shutdown input cannot form an executable
1125/// coordinator:
1126///
1127/// ```compile_fail
1128///
1129/// struct Plain;
1130/// impl behavior::Protocol for Plain {
1131///     type Addr = behavior::MailAddr;
1132///     type Msg = ();
1133/// }
1134/// impl behavior::Behavior for Plain {
1135///     type Protocol = Self;
1136///     type Event = behavior::User<behavior::MailAddr, ()>;
1137///     type Sends = Vec<behavior::Never>;
1138///     type Ph = behavior::Never;
1139///     type Error = behavior::Never;
1140///     type Birth = behavior::NoBirths;
1141///     fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> behavior::BehaviorActed<Self> {
1142///         Ok(behavior::Actions::cont())
1143///     }
1144/// }
1145///
1146/// fn require_behavior<B: behavior::Behavior>(_: B) {}
1147/// let child = behavior::CreationSequence::new()
1148///     .issue()
1149///     .expect("the first child creation ID exists");
1150/// let plan = behavior_actors::ShutdownPlan::new([vec![child]]).unwrap();
1151/// require_behavior(behavior_actors::ShutdownCoordinator::<Plain, Plain, behavior::ChildHead>::new(Plain, plan));
1152/// ```
1153pub 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    /// Start the wrapper before committed child creation can supply its plan.
1183    ///
1184    /// [`InstallShutdownPlan`] later installs exactly one validated plan. A
1185    /// shutdown request received first is retained by the state machine and
1186    /// begins that plan immediately on installation.
1187    #[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
1401/// Pure phased shutdown over an arbitrary closed child-protocol sum.
1402///
1403/// Each choice is interpreted in plan order through its exact concrete
1404/// `ShutdownChild<C, Occurrence>` request. Phase completion consumes the shared
1405/// creator-local creation ID because the child namespace is globally unique. This
1406/// phased ordering is Bombay policy, not an actor-model guarantee.
1407pub 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    /// Start the wrapper before committed heterogeneous children supply a plan.
1435    #[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/// Homogeneous dependency-ordered shutdown over one concrete child protocol.
1650#[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}