Skip to main content

behavior_actors/atomic/
drain.rs

1//! Actor-graph retirement policy shared by atomic actor owners.
2
3use core::time::Duration;
4
5use behavior::{ActionItemResult, ItemSettlement, SettledItem};
6
7use crate::{ScheduleAfter, TimerElapsed, TimerGeneration, TimerId};
8
9use super::schedule::ScheduleKey;
10
11/// Actor-graph retirement policy after aggregate shutdown.
12#[derive(Clone, Copy, Debug, Eq, PartialEq)]
13pub enum ActorDrainPolicy {
14    /// Wait until every owned actor has retired.
15    WaitForActorGraph,
16    /// Force-transfer unresolved actor ownership after this duration.
17    RetireActorGraphAfter {
18        /// Duration from drain start to forced actor retirement.
19        deadline: Duration,
20    },
21}
22
23pub(super) enum ShutdownDeadline {
24    Unlimited,
25    Scheduling(ScheduleKey),
26    Waiting(ScheduleKey),
27    NotScheduled(ActionItemResult<ScheduleAfter>),
28    Elapsed(TimerElapsed),
29}
30
31impl ShutdownDeadline {
32    const KEY: ScheduleKey = ScheduleKey::new(TimerId(0), TimerGeneration(0));
33
34    pub(super) fn begin(policy: ActorDrainPolicy) -> (Self, Option<ScheduleAfter>) {
35        match policy {
36            ActorDrainPolicy::WaitForActorGraph => (Self::Unlimited, None),
37            ActorDrainPolicy::RetireActorGraphAfter { deadline } => {
38                (Self::Scheduling(Self::KEY), Some(Self::KEY.after(deadline)))
39            }
40        }
41    }
42
43    pub(super) fn accept_schedule(
44        &mut self,
45        input: ActionItemResult<ScheduleAfter>,
46    ) -> Result<(), ActionItemResult<ScheduleAfter>> {
47        let Self::Scheduling(expected) = self else {
48            return Err(input);
49        };
50        let expected = *expected;
51        let input = match expected.admit_result(input) {
52            Ok(input) => input,
53            Err(input) => return Err(input),
54        };
55        match input {
56            SettledItem::Attempted(ItemSettlement::Accepted(_)) => {
57                *self = Self::Waiting(expected);
58                Ok(())
59            }
60            SettledItem::Attempted(ItemSettlement::Blocked { prerequisite, .. }) => {
61                match prerequisite {}
62            }
63            input @ (SettledItem::Attempted(
64                ItemSettlement::Rejected { .. } | ItemSettlement::Corrupt { .. },
65            )
66            | SettledItem::Unattempted(_)) => {
67                *self = Self::NotScheduled(input);
68                Ok(())
69            }
70        }
71    }
72
73    pub(super) fn accept_elapsed(&mut self, elapsed: TimerElapsed) -> Result<(), TimerElapsed> {
74        let Self::Waiting(expected) = self else {
75            return Err(elapsed);
76        };
77        let expected = *expected;
78        match expected.admit_elapsed(elapsed) {
79            Ok(elapsed) => {
80                *self = Self::Elapsed(elapsed);
81                Ok(())
82            }
83            Err(elapsed) => Err(elapsed),
84        }
85    }
86}
87
88#[expect(
89    dead_code,
90    reason = "Bombay's retirement custodian receives the complete cause"
91)]
92pub(super) enum ForcedRetirementCause {
93    WorkerShutdownIdsExhausted,
94    DeadlineNotScheduled(ActionItemResult<ScheduleAfter>),
95    DeadlineElapsed(TimerElapsed),
96}