behavior_actors/workflow/
latch.rs1use 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#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
14pub struct LatchReleased;
15
16#[derive(Clone, Copy, PartialEq, Eq)]
21pub struct LatchMessage<Route> {
22 pub reply_to: Route,
24}
25
26impl<Route> LatchMessage<Route> {
27 #[must_use]
29 pub const fn arrive(reply_to: Route) -> Self {
30 Self { reply_to }
31 }
32}
33
34pub enum LatchState<Route> {
36 Counting {
38 remaining: usize,
40 waiting: Vec<Route>,
42 },
43 Released,
45}
46
47pub 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 #[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 #[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}