Skip to main content

behavior_actors/workflow/
latch.rs

1//! Countdown latch coordination.
2
3use behavior::{
4    Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
5    SendEffects, User,
6};
7#[cfg(test)]
8use behavior::{Delivery, Recipient};
9
10use crate::DeliveryRoute;
11
12/// Release fact delivered to every accepted arrival.
13#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
14pub struct LatchReleased;
15
16/// One arrival at a [`Latch`].
17///
18/// The reply recipient is retained until release. An arrival after release is
19/// answered immediately and does not reopen the latch.
20#[derive(Clone, Copy, PartialEq, Eq)]
21pub struct LatchMessage<Route> {
22    /// Typed participant recipient awaiting the release fact.
23    pub reply_to: Route,
24}
25
26impl<Route> LatchMessage<Route> {
27    /// Construct one owned arrival.
28    #[must_use]
29    pub const fn arrive(reply_to: Route) -> Self {
30        Self { reply_to }
31    }
32}
33
34/// Complete semantic state of a [`Latch`].
35pub enum LatchState<Route> {
36    /// More arrivals are required; every listed participant is waiting.
37    Counting {
38        /// Number of additional arrivals needed, including the next one.
39        remaining: usize,
40        /// Accepted participants in arrival order.
41        waiting: Vec<Route>,
42    },
43    /// The threshold was reached and cannot be reset by this incarnation.
44    Released,
45}
46
47/// A one-generation countdown latch.
48///
49/// The state sum is [`LatchState::Counting`] or [`LatchState::Released`]. Each
50/// arrival decrements exactly once. The transition reaching zero atomically
51/// changes state to `Released` and emits one [`LatchReleased`] delivery to
52/// every accepted participant in arrival order. Later arrivals receive an
53/// immediate release. Initialization is empty; zero count starts released.
54/// There are no rejection, retry, cancellation, or timer paths. The latch
55/// never terminates itself. Countdown and ordering are Bombay workflow policy,
56/// while typed delivery is interpreted by Address and Communication.
57pub struct Latch<A, Route>
58where
59    A: Address,
60    Route: DeliveryRoute,
61    Route::Protocol: Protocol<Addr = A, Msg = LatchReleased>,
62{
63    state: LatchState<Route>,
64    marker: core::marker::PhantomData<fn() -> A>,
65}
66
67impl<A, Route> Latch<A, Route>
68where
69    A: Address,
70    Route: DeliveryRoute,
71    Route::Protocol: Protocol<Addr = A, Msg = LatchReleased>,
72{
73    /// Construct a latch requiring `count` arrivals.
74    #[must_use]
75    pub fn new(count: usize) -> Self {
76        Self {
77            state: if count == 0 {
78                LatchState::Released
79            } else {
80                LatchState::Counting {
81                    remaining: count,
82                    waiting: Vec::with_capacity(count),
83                }
84            },
85            marker: core::marker::PhantomData,
86        }
87    }
88
89    /// Borrow the complete current semantic state.
90    #[must_use]
91    pub const fn state(&self) -> &LatchState<Route> {
92        &self.state
93    }
94}
95
96impl<A, Route> BehaviorBase for Latch<A, Route>
97where
98    A: Address,
99    Route: DeliveryRoute,
100    Route::Protocol: Protocol<Addr = A, Msg = LatchReleased>,
101{
102    type Base = Self;
103
104    fn base(&self) -> &Self {
105        self
106    }
107}
108
109impl<A, Route> behavior::Protocol for Latch<A, Route>
110where
111    A: Address,
112    Route: DeliveryRoute,
113    Route::Protocol: Protocol<Addr = A, Msg = LatchReleased>,
114{
115    type Addr = A;
116    type Msg = LatchMessage<Route>;
117}
118
119impl<A, Route> Behavior for Latch<A, Route>
120where
121    A: Address,
122    Route: DeliveryRoute,
123    Route::Protocol: Protocol<Addr = A, Msg = LatchReleased>,
124    Route::Sends: behavior::SendsFor<User<A, LatchMessage<Route>>>,
125{
126    type Protocol = Self;
127    type Event = User<A, behavior::BehaviorMessage<Self>>;
128    type Sends = Route::Sends;
129    type Ph = Never;
130    type Error = Never;
131    type Birth = NoBirths;
132
133    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
134        let current = core::mem::replace(&mut self.state, LatchState::Released);
135        let (next, sends) = match current {
136            LatchState::Released => (
137                LatchState::Released,
138                event.message.reply_to.deliver(LatchReleased),
139            ),
140            LatchState::Counting {
141                remaining,
142                mut waiting,
143            } if remaining > 1 => {
144                waiting.push(event.message.reply_to);
145                (
146                    LatchState::Counting {
147                        remaining: remaining - 1,
148                        waiting,
149                    },
150                    Route::Sends::empty(),
151                )
152            }
153            LatchState::Counting { mut waiting, .. } => {
154                waiting.push(event.message.reply_to);
155                let mut sends = Route::Sends::empty();
156                for participant in waiting {
157                    sends.append(participant.deliver(LatchReleased));
158                }
159                (LatchState::Released, sends)
160            }
161        };
162        self.state = next;
163        Ok(Actions::send(sends))
164    }
165}
166
167#[cfg(test)]
168mod tests {
169    use super::*;
170    use crate::Activate as _;
171    use behavior::MailAddr;
172
173    #[test]
174    fn threshold_releases_waiters_in_arrival_order_once() {
175        let one = Recipient::from(MailAddr(1));
176        let two = Recipient::from(MailAddr(2));
177        let late = Recipient::from(MailAddr(3));
178        let mut latch = Latch::new(2).initialize().unwrap().behavior;
179
180        let first = latch
181            .receive(MailAddr(9), LatchMessage::arrive(one))
182            .unwrap();
183        assert!(first.sends.is_empty());
184        assert!(matches!(
185            latch.state(),
186            LatchState::Counting {
187                remaining: 1,
188                waiting,
189            } if waiting == &[one]
190        ));
191
192        let released = latch
193            .receive(MailAddr(9), LatchMessage::arrive(two))
194            .unwrap();
195        assert!(
196            released.sends
197                == vec![
198                    Delivery::new(one, LatchReleased),
199                    Delivery::new(two, LatchReleased),
200                ]
201        );
202        assert!(released.creates.is_empty());
203        assert!(matches!(latch.state(), LatchState::Released));
204
205        let immediate = latch
206            .receive(MailAddr(9), LatchMessage::arrive(late))
207            .unwrap();
208        assert!(immediate.sends == vec![Delivery::new(late, LatchReleased)]);
209    }
210
211    #[test]
212    fn zero_count_starts_released() {
213        let participant = Recipient::from(MailAddr(1));
214        let mut latch = Latch::new(0).initialize().unwrap().behavior;
215        assert!(matches!(latch.state(), LatchState::Released));
216        let actions = latch
217            .receive(MailAddr(9), LatchMessage::arrive(participant))
218            .unwrap();
219        assert!(actions.sends == vec![Delivery::new(participant, LatchReleased)]);
220    }
221}