Skip to main content

behavior/supervision/domain/
fleet.rs

1//! Supervised stable-child topology and lifecycle.
2
3use super::super::Strategy;
4use crate::CreationRejection;
5
6#[derive(Debug, Clone, Copy, PartialEq, Eq)]
7enum SlotState {
8    Available,
9    Retired,
10}
11
12struct Slot<N> {
13    nonce: N,
14    sequence: u64,
15    state: SlotState,
16}
17
18#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
19pub enum FleetError<N> {
20    #[error("unknown child nonce")]
21    UnknownChild(N),
22    #[error("duplicate child nonce")]
23    DuplicateChild(N),
24    #[error("child birth sequence exhausted")]
25    SequenceExhausted,
26}
27
28pub(crate) struct ReplacementCandidate<N> {
29    pub index: usize,
30    pub nonce: N,
31}
32
33/// Ordered topology of stable supervised children.
34pub(crate) struct Fleet<N> {
35    slots: Vec<Slot<N>>,
36    configured: usize,
37    next_sequence: u64,
38}
39
40impl<N: Copy + PartialEq> Fleet<N> {
41    pub fn configured(nonces: impl IntoIterator<Item = N>) -> Result<Self, FleetError<N>> {
42        let mut fleet = Self {
43            slots: Vec::new(),
44            configured: 0,
45            next_sequence: 0,
46        };
47        for nonce in nonces {
48            fleet.register(nonce)?;
49            fleet.configured += 1;
50        }
51        Ok(fleet)
52    }
53
54    pub fn configured_nonces(&self) -> impl Iterator<Item = N> + '_ {
55        self.slots[..self.configured].iter().map(|slot| slot.nonce)
56    }
57
58    #[must_use]
59    pub fn len(&self) -> usize {
60        self.slots.len()
61    }
62
63    pub fn is_available(&self, nonce: N) -> Result<bool, FleetError<N>> {
64        Ok(self.slot(nonce)?.state == SlotState::Available)
65    }
66
67    pub fn register(&mut self, nonce: N) -> Result<(), FleetError<N>> {
68        if self.position(nonce).is_some() {
69            return Err(FleetError::DuplicateChild(nonce));
70        }
71        let sequence = self.next_sequence;
72        self.next_sequence = self
73            .next_sequence
74            .checked_add(1)
75            .ok_or(FleetError::SequenceExhausted)?;
76        self.slots.push(Slot {
77            nonce,
78            sequence,
79            state: SlotState::Available,
80        });
81        Ok(())
82    }
83
84    pub(crate) fn resolve_creation(&mut self, nonce: N, result: Result<(), CreationRejection>) {
85        if let Some(position) = self.position(nonce) {
86            self.slots[position].state = match result {
87                Ok(()) => SlotState::Available,
88                Err(_) => SlotState::Retired,
89            };
90        }
91    }
92
93    pub fn retire(&mut self, nonce: N) -> Result<(), FleetError<N>> {
94        self.slot_mut(nonce)?.state = SlotState::Retired;
95        Ok(())
96    }
97
98    /// A replacement command keeps the stable proxy slot available while its
99    /// worker replacement is pending.
100    pub fn replacement_requested(&mut self, nonce: N) -> Result<(), FleetError<N>> {
101        self.slot_mut(nonce)?.state = SlotState::Available;
102        Ok(())
103    }
104
105    pub fn replacements(
106        &self,
107        failed: N,
108        strategy: Strategy,
109    ) -> Result<Vec<ReplacementCandidate<N>>, FleetError<N>> {
110        let failed = self
111            .position(failed)
112            .ok_or(FleetError::UnknownChild(failed))?;
113        let sequence = self.slots[failed].sequence;
114        Ok(self
115            .slots
116            .iter()
117            .enumerate()
118            .filter(|(index, slot)| match strategy {
119                Strategy::OneForOne => *index == failed,
120                Strategy::OneForAll => slot.state == SlotState::Available,
121                Strategy::RestForOne => {
122                    slot.state == SlotState::Available && slot.sequence >= sequence
123                }
124            })
125            .map(|(index, slot)| ReplacementCandidate {
126                index,
127                nonce: slot.nonce,
128            })
129            .collect())
130    }
131
132    fn position(&self, nonce: N) -> Option<usize> {
133        self.slots.iter().position(|slot| slot.nonce == nonce)
134    }
135
136    fn slot(&self, nonce: N) -> Result<&Slot<N>, FleetError<N>> {
137        self.position(nonce)
138            .map(|position| &self.slots[position])
139            .ok_or(FleetError::UnknownChild(nonce))
140    }
141
142    fn slot_mut(&mut self, nonce: N) -> Result<&mut Slot<N>, FleetError<N>> {
143        self.position(nonce)
144            .map(|position| &mut self.slots[position])
145            .ok_or(FleetError::UnknownChild(nonce))
146    }
147}
148
149impl<N: Copy + PartialEq> TryFrom<Vec<N>> for Fleet<N> {
150    type Error = FleetError<N>;
151
152    fn try_from(nonces: Vec<N>) -> Result<Self, Self::Error> {
153        Self::configured(nonces)
154    }
155}
156
157#[cfg(test)]
158mod tests {
159    use super::*;
160
161    #[test]
162    fn rest_for_one_uses_birth_sequence_and_skips_retired_slots() {
163        let mut fleet = Fleet::configured([7, 2, 9]).unwrap();
164        fleet.retire(2).unwrap();
165        let replacements = fleet.replacements(7, Strategy::RestForOne).unwrap();
166        assert_eq!(
167            replacements
168                .into_iter()
169                .map(|candidate| candidate.nonce)
170                .collect::<Vec<_>>(),
171            [7, 9]
172        );
173    }
174
175    #[test]
176    fn one_for_one_addresses_the_exact_slot_for_each_stop_observation() {
177        let mut fleet = Fleet::configured([7]).unwrap();
178        fleet.retire(7).unwrap();
179        let replacements = fleet.replacements(7, Strategy::OneForOne).unwrap();
180        assert_eq!(replacements.len(), 1);
181        assert_eq!(replacements[0].nonce, 7);
182    }
183
184    #[test]
185    fn creation_resolution_commits_success_and_retires_rejection() {
186        let mut fleet = Fleet::configured([7, 9]).unwrap();
187        fleet.retire(7).unwrap();
188        fleet.resolve_creation(7, Ok(()));
189        assert_eq!(fleet.is_available(7), Ok(true));
190
191        fleet.resolve_creation(9, Err(CreationRejection::EnvironmentFailed));
192        assert_eq!(fleet.is_available(9), Ok(false));
193
194        fleet.resolve_creation(99, Err(CreationRejection::EnvironmentFailed));
195        assert_eq!(fleet.len(), 2);
196    }
197}