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, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
16pub struct Sequence(pub u64);
17
18#[derive(Debug, Clone, PartialEq, Eq)]
20pub enum SequencerState {
21 Active {
23 expected: Sequence,
25 buffered: usize,
27 },
28 Exhausted,
30}
31
32#[derive(Debug, PartialEq, Eq)]
34pub enum SequencerOutcome<T> {
35 Accepted {
37 released: usize,
40 buffered: usize,
42 },
43 Stale {
45 sequence: Sequence,
47 value: T,
49 expected: Sequence,
51 },
52 Duplicate {
54 value: T,
56 sequence: Sequence,
58 },
59 Exhausted {
61 sequence: Sequence,
63 value: T,
65 },
66}
67
68pub enum SequencerMessage<T, TargetRoute, ReplyRoute> {
70 Offer {
72 sequence: Sequence,
74 value: T,
76 to: TargetRoute,
78 reply_to: ReplyRoute,
80 },
81}
82
83struct Pending<T, TargetRoute> {
84 value: T,
85 to: TargetRoute,
86}
87
88pub struct Sequencer<
102 A: Address,
103 T,
104 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
105 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = SequencerOutcome<T>>>,
106> {
107 expected: Option<Sequence>,
108 pending: BTreeMap<Sequence, Pending<T, TargetRoute>>,
109 marker: core::marker::PhantomData<fn() -> (A, ReplyRoute)>,
110}
111
112impl<A, T, TargetRoute, ReplyRoute> Sequencer<A, T, TargetRoute, ReplyRoute>
113where
114 A: Address,
115 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
116 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = SequencerOutcome<T>>>,
117{
118 #[must_use]
120 pub fn new(first: Sequence) -> Self {
121 Self {
122 expected: Some(first),
123 pending: BTreeMap::new(),
124 marker: core::marker::PhantomData,
125 }
126 }
127
128 #[must_use]
130 pub fn state(&self) -> SequencerState {
131 match self.expected {
132 Some(expected) => SequencerState::Active {
133 expected,
134 buffered: self.pending.len(),
135 },
136 None => SequencerState::Exhausted,
137 }
138 }
139
140 fn actions(
141 deliveries: TargetRoute::Sends,
142 outcomes: ReplyRoute::Sends,
143 ) -> Actions<A, Never, DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>, NoBirths> {
144 Actions::send(DeliveryOutcomes {
145 deliveries,
146 outcomes,
147 })
148 }
149}
150
151impl<A, T, TargetRoute, ReplyRoute> BehaviorBase for Sequencer<A, T, TargetRoute, ReplyRoute>
152where
153 A: Address,
154 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
155 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = SequencerOutcome<T>>>,
156{
157 type Base = Self;
158
159 fn base(&self) -> &Self {
160 self
161 }
162}
163
164impl<A, T, TargetRoute, ReplyRoute> behavior::Protocol for Sequencer<A, T, TargetRoute, ReplyRoute>
165where
166 A: Address,
167 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
168 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = SequencerOutcome<T>>>,
169{
170 type Addr = A;
171 type Msg = SequencerMessage<T, TargetRoute, ReplyRoute>;
172}
173
174impl<A, T, TargetRoute, ReplyRoute> Behavior for Sequencer<A, T, TargetRoute, ReplyRoute>
175where
176 A: Address,
177 TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
178 ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = SequencerOutcome<T>>>,
179 TargetRoute::Sends: behavior::SendsFor<User<A, SequencerMessage<T, TargetRoute, ReplyRoute>>>,
180 ReplyRoute::Sends: behavior::SendsFor<User<A, SequencerMessage<T, TargetRoute, ReplyRoute>>>,
181{
182 type Protocol = Self;
183 type Event = User<A, behavior::BehaviorMessage<Self>>;
184 type Sends = DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>;
185 type Ph = Never;
186 type Error = Never;
187 type Birth = NoBirths;
188
189 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
190 let SequencerMessage::Offer {
191 sequence,
192 value,
193 to,
194 reply_to,
195 } = event.message;
196 let Some(expected) = self.expected else {
197 return Ok(Self::actions(
198 TargetRoute::Sends::empty(),
199 reply_to.deliver(SequencerOutcome::Exhausted { sequence, value }),
200 ));
201 };
202 if sequence < expected {
203 return Ok(Self::actions(
204 TargetRoute::Sends::empty(),
205 reply_to.deliver(SequencerOutcome::Stale {
206 sequence,
207 value,
208 expected,
209 }),
210 ));
211 }
212 if self.pending.contains_key(&sequence) {
213 return Ok(Self::actions(
214 TargetRoute::Sends::empty(),
215 reply_to.deliver(SequencerOutcome::Duplicate { value, sequence }),
216 ));
217 }
218
219 self.pending.insert(sequence, Pending { value, to });
220 let mut cursor = expected;
221 let mut deliveries = TargetRoute::Sends::empty();
222 let mut released = 0;
223 while let Some(pending) = self.pending.remove(&cursor) {
224 deliveries.append(pending.to.deliver(pending.value));
225 released += 1;
226 if cursor.0 == u64::MAX {
227 self.expected = None;
228 break;
229 }
230 cursor = Sequence(cursor.0 + 1);
231 self.expected = Some(cursor);
232 }
233 let outcome = SequencerOutcome::Accepted {
234 released,
235 buffered: self.pending.len(),
236 };
237 Ok(Self::actions(deliveries, reply_to.deliver(outcome)))
238 }
239}
240
241#[cfg(test)]
242mod tests {
243 use crate::Activate as _;
244 use behavior::{Delivery, MailAddr, Recipient, Step};
245
246 use super::*;
247
248 struct Target;
249 struct Reply;
250
251 macro_rules! leaf {
252 ($name:ident, $msg:ty) => {
253 impl behavior::Protocol for $name {
254 type Addr = MailAddr;
255 type Msg = $msg;
256 }
257
258 impl Behavior for $name {
259 type Protocol = Self;
260 type Event = behavior::User<MailAddr, behavior::BehaviorMessage<Self>>;
261 type Sends = Vec<behavior::Never>;
262 type Ph = behavior::Never;
263 type Error = behavior::Never;
264 type Birth = behavior::NoBirths;
265 fn transition(
266 &mut self,
267 _: behavior::ActiveTurn,
268 _: Self::Event,
269 ) -> behavior::BehaviorActed<Self> {
270 Ok(behavior::Actions::cont())
271 }
272 }
273 };
274 }
275
276 leaf!(Target, u8);
277 leaf!(Reply, SequencerOutcome<u8>);
278
279 type Subject = Sequencer<MailAddr, u8, Recipient<Target>, Recipient<Reply>>;
280
281 fn offer(
282 subject: &mut crate::Active<Subject>,
283 sequence: u64,
284 value: u8,
285 ) -> behavior::Actions<
286 MailAddr,
287 behavior::Never,
288 DeliveryOutcomes<Vec<Delivery<Target>>, Vec<Delivery<Reply>>>,
289 behavior::NoBirths,
290 > {
291 subject
292 .receive(
293 MailAddr(9),
294 SequencerMessage::Offer {
295 sequence: Sequence(sequence),
296 value,
297 to: Recipient::global(MailAddr(1)),
298 reply_to: Recipient::global(MailAddr(2)),
299 },
300 )
301 .unwrap()
302 }
303
304 #[test]
305 fn gaps_release_only_after_the_missing_position_arrives() {
306 let mut subject = (Subject::new(Sequence(3))).initialize().unwrap().behavior;
307 let future = offer(&mut subject, 4, 40);
308 assert!(future.sends.deliveries.is_empty());
309 assert_eq!(
310 subject.state(),
311 SequencerState::Active {
312 expected: Sequence(3),
313 buffered: 1
314 }
315 );
316
317 let released = offer(&mut subject, 3, 30);
318 assert_eq!(
319 released
320 .sends
321 .deliveries
322 .iter()
323 .map(|delivery| delivery.message)
324 .collect::<Vec<_>>(),
325 vec![30, 40]
326 );
327 assert_eq!(
328 subject.state(),
329 SequencerState::Active {
330 expected: Sequence(5),
331 buffered: 0
332 }
333 );
334 }
335
336 #[test]
337 fn stale_and_duplicate_offers_return_the_rejected_value() {
338 let mut subject = (Subject::new(Sequence(1))).initialize().unwrap().behavior;
339 let buffered = offer(&mut subject, 2, 20);
340 assert!(matches!(
341 buffered.sends.outcomes[0].message,
342 SequencerOutcome::Accepted {
343 released: 0,
344 buffered: 1
345 }
346 ));
347 let duplicate = offer(&mut subject, 2, 21);
348 assert!(matches!(
349 duplicate.sends.outcomes[0].message,
350 SequencerOutcome::Duplicate {
351 value: 21,
352 sequence: Sequence(2)
353 }
354 ));
355 let released = offer(&mut subject, 1, 10);
356 assert!(matches!(
357 released.sends.outcomes[0].message,
358 SequencerOutcome::Accepted {
359 released: 2,
360 buffered: 0
361 }
362 ));
363 let stale = offer(&mut subject, 1, 11);
364 assert!(matches!(
365 stale.sends.outcomes[0].message,
366 SequencerOutcome::Stale {
367 sequence: Sequence(1),
368 value: 11,
369 expected: Sequence(3)
370 }
371 ));
372 }
373
374 #[test]
375 fn maximum_position_exhausts_without_wrapping() {
376 let mut subject = (Subject::new(Sequence(u64::MAX)))
377 .initialize()
378 .unwrap()
379 .behavior;
380 let delivered = offer(&mut subject, u64::MAX, 1);
381 assert_eq!(delivered.sends.deliveries.len(), 1);
382 assert!(matches!(delivered.become_, Step::Continue));
383 assert_eq!(subject.state(), SequencerState::Exhausted);
384 let rejected = offer(&mut subject, u64::MAX, 2);
385 assert!(matches!(
386 rejected.sends.outcomes[0].message,
387 SequencerOutcome::Exhausted {
388 sequence: Sequence(u64::MAX),
389 value: 2,
390 }
391 ));
392 }
393}