1use std::collections::BTreeMap;
4
5use super::DeliveryOutcomes;
6
7use behavior::{
8 Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
9 SendEffects, User,
10};
11
12use crate::DeliveryRoute;
13
14#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct OrderGateState<K> {
17 pub watermark: Option<K>,
19 held: usize,
20}
21
22impl<K> OrderGateState<K> {
23 #[must_use]
25 pub fn held(&self) -> usize {
26 self.held
27 }
28}
29
30#[derive(Debug, PartialEq, Eq)]
32pub enum OrderGateOutcome<K, T> {
33 Held {
35 key: K,
37 held: usize,
39 },
40 Delivered {
42 key: K,
44 },
45 Duplicate {
47 key: K,
49 value: T,
51 },
52 Opened {
54 through: K,
56 released: usize,
58 held: usize,
60 },
61 StaleOpening {
63 requested: K,
65 current: K,
67 },
68}
69
70pub enum OrderGateMessage<K, T, TargetRoute, ReplyRoute> {
72 Hold {
74 key: K,
76 value: T,
78 to: TargetRoute,
80 reply_to: ReplyRoute,
82 },
83 OpenThrough {
85 through: K,
87 reply_to: ReplyRoute,
89 },
90}
91
92struct Held<T, TargetRoute> {
93 value: T,
94 to: TargetRoute,
95}
96
97pub struct OrderGate<
108 A: Address,
109 K,
110 T,
111 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
112 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = OrderGateOutcome<K, T>>>,
113> {
114 watermark: Option<K>,
115 held: BTreeMap<K, Held<T, TargetRoute>>,
116 marker: core::marker::PhantomData<fn() -> (A, ReplyRoute)>,
117}
118
119impl<A, K, T, TargetRoute, ReplyRoute> OrderGate<A, K, T, TargetRoute, ReplyRoute>
120where
121 A: Address,
122 K: Clone + Ord,
123 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
124 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = OrderGateOutcome<K, T>>>,
125{
126 #[must_use]
128 pub const fn new() -> Self {
129 Self {
130 watermark: None,
131 held: BTreeMap::new(),
132 marker: core::marker::PhantomData,
133 }
134 }
135
136 #[must_use]
138 pub fn state(&self) -> OrderGateState<K> {
139 OrderGateState {
140 watermark: self.watermark.clone(),
141 held: self.held.len(),
142 }
143 }
144
145 fn sends(
146 deliveries: TargetRoute::Sends,
147 reply_to: ReplyRoute,
148 outcome: OrderGateOutcome<K, T>,
149 ) -> Actions<A, Never, DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>, NoBirths> {
150 Actions::send(DeliveryOutcomes {
151 deliveries,
152 outcomes: reply_to.deliver(outcome),
153 })
154 }
155
156 fn hold(
157 &mut self,
158 key: K,
159 value: T,
160 to: TargetRoute,
161 reply_to: ReplyRoute,
162 ) -> Actions<A, Never, DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>, NoBirths> {
163 if self
164 .watermark
165 .as_ref()
166 .is_some_and(|watermark| key <= *watermark)
167 {
168 return Self::sends(
169 to.deliver(value),
170 reply_to,
171 OrderGateOutcome::Delivered { key },
172 );
173 }
174 if self.held.contains_key(&key) {
175 return Self::sends(
176 TargetRoute::Sends::empty(),
177 reply_to,
178 OrderGateOutcome::Duplicate { key, value },
179 );
180 }
181 self.held.insert(key.clone(), Held { value, to });
182 Self::sends(
183 TargetRoute::Sends::empty(),
184 reply_to,
185 OrderGateOutcome::Held {
186 key,
187 held: self.held.len(),
188 },
189 )
190 }
191
192 fn open(
193 &mut self,
194 through: K,
195 reply_to: ReplyRoute,
196 ) -> Actions<A, Never, DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>, NoBirths> {
197 if let Some(current) = &self.watermark
198 && through <= *current
199 {
200 return Self::sends(
201 TargetRoute::Sends::empty(),
202 reply_to,
203 OrderGateOutcome::StaleOpening {
204 requested: through,
205 current: current.clone(),
206 },
207 );
208 }
209 let keys = self
210 .held
211 .range(..=through.clone())
212 .map(|(key, _)| key.clone())
213 .collect::<Vec<_>>();
214 let released = keys.len();
215 let mut deliveries = TargetRoute::Sends::empty();
216 for key in keys {
217 if let Some(held) = self.held.remove(&key) {
218 deliveries.append(held.to.deliver(held.value));
219 }
220 }
221 self.watermark = Some(through.clone());
222 Self::sends(
223 deliveries,
224 reply_to,
225 OrderGateOutcome::Opened {
226 through,
227 released,
228 held: self.held.len(),
229 },
230 )
231 }
232}
233
234impl<A, K, T, TargetRoute, ReplyRoute> Default for OrderGate<A, K, T, TargetRoute, ReplyRoute>
235where
236 A: Address,
237 K: Clone + Ord,
238 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
239 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = OrderGateOutcome<K, T>>>,
240{
241 fn default() -> Self {
242 Self::new()
243 }
244}
245
246impl<A, K, T, TargetRoute, ReplyRoute> BehaviorBase for OrderGate<A, K, T, TargetRoute, ReplyRoute>
247where
248 A: Address,
249 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
250 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = OrderGateOutcome<K, T>>>,
251{
252 type Base = Self;
253 fn base(&self) -> &Self {
254 self
255 }
256}
257
258impl<A, K, T, TargetRoute, ReplyRoute> behavior::Protocol
259 for OrderGate<A, K, T, TargetRoute, ReplyRoute>
260where
261 A: Address,
262 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
263 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = OrderGateOutcome<K, T>>>,
264{
265 type Addr = A;
266 type Msg = OrderGateMessage<K, T, TargetRoute, ReplyRoute>;
267}
268
269impl<A, K, T, TargetRoute, ReplyRoute> Behavior for OrderGate<A, K, T, TargetRoute, ReplyRoute>
270where
271 A: Address,
272 K: Clone + Ord,
273 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
274 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = OrderGateOutcome<K, T>>>,
275 TargetRoute::Sends:
276 behavior::SendsFor<User<A, OrderGateMessage<K, T, TargetRoute, ReplyRoute>>>,
277 ReplyRoute::Sends: behavior::SendsFor<User<A, OrderGateMessage<K, T, TargetRoute, ReplyRoute>>>,
278{
279 type Protocol = Self;
280 type Event = User<A, behavior::BehaviorMessage<Self>>;
281 type Sends = DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>;
282 type Ph = Never;
283 type Error = Never;
284 type Birth = NoBirths;
285 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
286 Ok(match event.message {
287 OrderGateMessage::Hold {
288 key,
289 value,
290 to,
291 reply_to,
292 } => self.hold(key, value, to, reply_to),
293 OrderGateMessage::OpenThrough { through, reply_to } => self.open(through, reply_to),
294 })
295 }
296}
297
298#[cfg(test)]
299mod tests {
300 use super::*;
301 use crate::Activate as _;
302 use behavior::{Delivery, MailAddr, Recipient};
303 struct Target;
304 struct Reply;
305 macro_rules! leaf {
306 ($n:ident,$m:ty) => {
307 impl behavior::Protocol for $n {
308 type Addr = MailAddr;
309 type Msg = $m;
310 }
311
312 impl Behavior for $n {
313 type Protocol = Self;
314 type Event = User<MailAddr, $m>;
315 type Sends = Vec<Never>;
316 type Ph = Never;
317 type Error = Never;
318 type Birth = NoBirths;
319 fn transition(
320 &mut self,
321 _: behavior::ActiveTurn,
322 _: Self::Event,
323 ) -> BehaviorActed<Self> {
324 Ok(Actions::cont())
325 }
326 }
327 };
328 }
329 leaf!(Target, u8);
330 leaf!(Reply,OrderGateOutcome<u8,u8>);
331 type Subject = OrderGate<MailAddr, u8, u8, Recipient<Target>, Recipient<Reply>>;
332 fn hold(
333 s: &mut crate::Active<Subject>,
334 key: u8,
335 value: u8,
336 ) -> Actions<
337 MailAddr,
338 Never,
339 DeliveryOutcomes<Vec<Delivery<Target>>, Vec<Delivery<Reply>>>,
340 NoBirths,
341 > {
342 s.receive(
343 MailAddr(9),
344 OrderGateMessage::Hold {
345 key,
346 value,
347 to: Recipient::global(MailAddr(1)),
348 reply_to: Recipient::global(MailAddr(2)),
349 },
350 )
351 .unwrap()
352 }
353 fn open(
354 s: &mut crate::Active<Subject>,
355 through: u8,
356 ) -> Actions<
357 MailAddr,
358 Never,
359 DeliveryOutcomes<Vec<Delivery<Target>>, Vec<Delivery<Reply>>>,
360 NoBirths,
361 > {
362 s.receive(
363 MailAddr(9),
364 OrderGateMessage::OpenThrough {
365 through,
366 reply_to: Recipient::global(MailAddr(2)),
367 },
368 )
369 .unwrap()
370 }
371 #[test]
372 fn opening_releases_in_key_order_and_future_open_keys_deliver_immediately() {
373 let mut s = (Subject::new()).initialize().unwrap().behavior;
374 let held_three = hold(&mut s, 3, 30);
375 assert!(matches!(
376 held_three.sends.outcomes[0].message,
377 OrderGateOutcome::Held { key: 3, held: 1 }
378 ));
379 let held_one = hold(&mut s, 1, 10);
380 assert!(matches!(
381 held_one.sends.outcomes[0].message,
382 OrderGateOutcome::Held { key: 1, held: 2 }
383 ));
384 let a = open(&mut s, 2);
385 assert_eq!(
386 a.sends
387 .deliveries
388 .iter()
389 .map(|d| d.message)
390 .collect::<Vec<_>>(),
391 vec![10]
392 );
393 let delivered = hold(&mut s, 2, 20);
394 assert_eq!(delivered.sends.deliveries.len(), 1);
395 assert_eq!(s.state().watermark, Some(2));
396 assert_eq!(s.state().held(), 1);
397 }
398 #[test]
399 fn duplicate_and_stale_opening_are_atomic() {
400 let mut s = (Subject::new()).initialize().unwrap().behavior;
401 let held = hold(&mut s, 2, 20);
402 assert!(matches!(
403 held.sends.outcomes[0].message,
404 OrderGateOutcome::Held { key: 2, held: 1 }
405 ));
406 let duplicate = hold(&mut s, 2, 21);
407 assert!(matches!(
408 duplicate.sends.outcomes[0].message,
409 OrderGateOutcome::Duplicate { value: 21, .. }
410 ));
411 let opened = open(&mut s, 1);
412 assert!(matches!(
413 opened.sends.outcomes[0].message,
414 OrderGateOutcome::Opened {
415 through: 1,
416 released: 0,
417 held: 1
418 }
419 ));
420 let stale = open(&mut s, 1);
421 assert!(matches!(
422 stale.sends.outcomes[0].message,
423 OrderGateOutcome::StaleOpening { .. }
424 ));
425 assert_eq!(s.state().held, 1);
426 }
427}