Skip to main content

behavior_actors/routing/
sequencer.rs

1//! Gap-buffered delivery in an explicit sequence domain.
2
3use 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/// A position in a [`Sequencer`] stream.
15#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
16pub struct Sequence(pub u64);
17
18/// Complete lifecycle state of a [`Sequencer`].
19#[derive(Debug, Clone, PartialEq, Eq)]
20pub enum SequencerState {
21    /// The next position can still be accepted.
22    Active {
23        /// Position whose delivery releases the next contiguous run.
24        expected: Sequence,
25        /// Number of accepted values waiting beyond a gap.
26        buffered: usize,
27    },
28    /// `Sequence(u64::MAX)` was delivered; no successor is representable.
29    Exhausted,
30}
31
32/// Factual result of offering one sequenced value.
33#[derive(Debug, PartialEq, Eq)]
34pub enum SequencerOutcome<T> {
35    /// The value was accepted and this many contiguous values were delivered.
36    Accepted {
37        /// Number released by this offer, including the offered value when it
38        /// filled the current position.
39        released: usize,
40        /// Number still retained beyond a gap.
41        buffered: usize,
42    },
43    /// The position precedes the current expected position.
44    Stale {
45        /// Rejected sequence position.
46        sequence: Sequence,
47        /// Rejected value, returned without cloning or loss.
48        value: T,
49        /// Current expected position.
50        expected: Sequence,
51    },
52    /// A value already occupies this future position.
53    Duplicate {
54        /// Rejected value, returned without replacing the accepted value.
55        value: T,
56        /// Occupied position.
57        sequence: Sequence,
58    },
59    /// The sequence domain has no representable successor.
60    Exhausted {
61        /// Rejected sequence position.
62        sequence: Sequence,
63        /// Rejected value.
64        value: T,
65    },
66}
67
68/// Commands accepted by [`Sequencer`].
69pub enum SequencerMessage<T, TargetRoute, ReplyRoute> {
70    /// Offer one value for delivery after all preceding positions.
71    Offer {
72        /// Explicit sequence position.
73        sequence: Sequence,
74        /// Owned value.
75        value: T,
76        /// Destination used when this position becomes contiguous.
77        to: TargetRoute,
78        /// Recipient for the complete offer outcome.
79        reply_to: ReplyRoute,
80    },
81}
82
83struct Pending<T, TargetRoute> {
84    value: T,
85    to: TargetRoute,
86}
87
88/// Deterministic gap-buffered ordered-delivery policy.
89///
90/// `Sequencer` accepts each position at most once, retains future positions in
91/// a `BTreeMap`, and releases only the contiguous run beginning at `expected`.
92/// Stale, duplicate, and exhausted offers return ownership through
93/// [`SequencerOutcome`]; an accepted value is never overwritten. Delivering
94/// `Sequence(u64::MAX)` moves to the explicit terminal policy state
95/// [`SequencerState::Exhausted`] rather than wrapping. Initialization is empty,
96/// the template creates no actors, and it does not terminate the hosting actor.
97/// Ordered release and bounded integer exhaustion are deliberate Bombay policy;
98/// mailbox FIFO and physical delivery remain `bombay-communication`
99/// responsibilities. The fold requires no capability beyond ordinary typed
100/// sends and has no panic path.
101pub 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    /// Construct an empty sequencer beginning at `first`.
119    #[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    /// Borrow the complete observable lifecycle state.
129    #[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}