Skip to main content

behavior_actors/routing/
buffer.rs

1//! Bounded policy buffering over ordinary typed delivery.
2
3use std::collections::VecDeque;
4
5#[cfg(test)]
6use behavior::MessageProtocol;
7use behavior::{
8    Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
9    SendEffects, User,
10};
11use thiserror::Error;
12
13use super::DeliveryOutcomes;
14use crate::DeliveryRoute;
15
16/// Exhaustive behavior-owned policy when a [`Buffer`] is full.
17#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
18pub enum OverflowPolicy {
19    /// Refuse the offered value and retain the complete queue.
20    Reject,
21    /// Return the oldest accepted value, then accept the new value.
22    DropOldest,
23    /// Return the newly offered value without changing the queue.
24    DropNewest,
25}
26
27/// Why an offered value was not accepted.
28#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
29pub enum BufferRejection {
30    /// The configured policy rejects offers at capacity.
31    Full,
32    /// The configured policy explicitly selects the newest value for return.
33    DroppedNewest,
34}
35
36/// Factual result delivered for one buffer command.
37#[derive(Debug, PartialEq, Eq)]
38pub enum BufferOutcome<T> {
39    /// The offered value is now owned by the queue.
40    Accepted {
41        /// Queue depth after acceptance.
42        depth: usize,
43    },
44    /// The offer was not accepted; ownership is returned here.
45    Rejected {
46        /// Unaccepted owned value.
47        value: T,
48        /// Policy reason for rejection.
49        reason: BufferRejection,
50    },
51    /// A previously accepted oldest value was evicted and returned.
52    Evicted {
53        /// Evicted owned value.
54        value: T,
55    },
56    /// One value was released to its target.
57    Released {
58        /// Queue depth after release.
59        remaining: usize,
60    },
61    /// A release found no accepted value.
62    Empty,
63}
64
65/// One accepted value and the recipient that must receive an eviction fact.
66pub struct Buffered<T, Route> {
67    /// Value owned by the buffer.
68    pub value: T,
69    /// Typed outcome recipient supplied with the original offer.
70    pub reply_to: Route,
71}
72
73/// Complete valid state product of a [`Buffer`].
74pub struct BufferState<T, Route> {
75    /// Positive maximum accepted queue length.
76    pub capacity: usize,
77    /// Exhaustive policy used at capacity.
78    pub overflow: OverflowPolicy,
79    queued: VecDeque<Buffered<T, Route>>,
80}
81
82impl<T, Route> BufferState<T, Route> {
83    /// Number of accepted values currently owned.
84    #[must_use]
85    pub fn len(&self) -> usize {
86        self.queued.len()
87    }
88
89    /// Whether no accepted value is currently owned.
90    #[must_use]
91    pub fn is_empty(&self) -> bool {
92        self.queued.is_empty()
93    }
94
95    /// Iterate accepted values in release order.
96    #[must_use]
97    pub fn queued(&self) -> impl ExactSizeIterator<Item = &Buffered<T, Route>> {
98        self.queued.iter()
99    }
100}
101
102/// Commands accepted by [`Buffer`].
103pub enum BufferMessage<T, TargetRoute, ReplyRoute> {
104    /// Offer ownership of one value under the configured overflow policy.
105    Offer {
106        /// Offered value.
107        value: T,
108        /// Recipient for acceptance, rejection, or later eviction.
109        reply_to: ReplyRoute,
110    },
111    /// Release the oldest accepted value to a concrete typed destination.
112    Release {
113        /// Destination for the released value.
114        to: TargetRoute,
115        /// Recipient for the release or empty result.
116        reply_to: ReplyRoute,
117    },
118}
119
120/// Invalid buffer definition.
121#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
122pub enum BufferConfigError {
123    /// A zero-capacity buffer could never accept ownership.
124    #[error("buffer capacity must be positive")]
125    ZeroCapacity,
126}
127
128/// Validated, protocol-independent buffer policy.
129#[derive(Debug, Clone, Copy, PartialEq, Eq)]
130pub struct BufferConfiguration {
131    capacity: usize,
132    overflow: OverflowPolicy,
133}
134
135impl BufferConfiguration {
136    /// Validate policy before it is bound to any actor protocol.
137    ///
138    /// # Errors
139    ///
140    /// Returns [`BufferConfigError::ZeroCapacity`] when `capacity` is zero.
141    pub fn new(capacity: usize, overflow: OverflowPolicy) -> Result<Self, BufferConfigError> {
142        if capacity == 0 {
143            return Err(BufferConfigError::ZeroCapacity);
144        }
145        Ok(Self { capacity, overflow })
146    }
147}
148
149/// Bounded FIFO policy behavior over typed actor deliveries.
150///
151/// State is the named [`BufferState`] product. Offers below capacity commit
152/// ownership and report `Accepted`. At capacity, [`OverflowPolicy`] either
153/// returns the new value, or returns the oldest value before accepting the new
154/// one. Release moves exactly the oldest value to the destination lane and
155/// reports the remaining depth; empty release reports `Empty`. No payload is
156/// silently discarded. Initialization is empty, transitions do not create
157/// actors, and the buffer never terminates itself. FIFO and overflow semantics
158/// are Bombay policy. Physical mailbox buffering, admission, fairness, and
159/// backpressure remain `bombay-communication` responsibilities. Construction
160/// returns [`BufferConfigError`] rather than creating an ownership sink. The
161/// internal full-queue eviction uses the proven invariant that positive
162/// capacity plus the full-offer branch implies a non-empty queue; callers
163/// cannot violate that invariant through the public API.
164pub struct Buffer<A, T, TargetRoute, ReplyRoute>
165where
166    A: Address,
167    TargetRoute: DeliveryRoute,
168    TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
169    ReplyRoute: DeliveryRoute,
170    ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
171{
172    state: BufferState<T, ReplyRoute>,
173    marker: core::marker::PhantomData<fn() -> (A, TargetRoute)>,
174}
175
176impl<A, T, TargetRoute, ReplyRoute> Buffer<A, T, TargetRoute, ReplyRoute>
177where
178    A: Address,
179    TargetRoute: DeliveryRoute,
180    TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
181    ReplyRoute: DeliveryRoute,
182    ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
183{
184    /// Bind validated policy to an empty buffer actor.
185    #[must_use]
186    pub fn new(configuration: BufferConfiguration) -> Self {
187        Self {
188            state: BufferState {
189                capacity: configuration.capacity,
190                overflow: configuration.overflow,
191                queued: VecDeque::with_capacity(configuration.capacity),
192            },
193            marker: core::marker::PhantomData,
194        }
195    }
196
197    /// Borrow the complete current buffer state.
198    #[must_use]
199    pub const fn state(&self) -> &BufferState<T, ReplyRoute> {
200        &self.state
201    }
202
203    fn actions(
204        deliveries: TargetRoute::Sends,
205        outcomes: ReplyRoute::Sends,
206    ) -> Actions<A, Never, DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>, NoBirths> {
207        Actions::send(DeliveryOutcomes {
208            deliveries,
209            outcomes,
210        })
211    }
212}
213
214impl<A, T, TargetRoute, ReplyRoute> BehaviorBase for Buffer<A, T, TargetRoute, ReplyRoute>
215where
216    A: Address,
217    TargetRoute: DeliveryRoute,
218    TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
219    ReplyRoute: DeliveryRoute,
220    ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
221{
222    type Base = Self;
223
224    fn base(&self) -> &Self {
225        self
226    }
227}
228
229impl<A, T, TargetRoute, ReplyRoute> behavior::Protocol for Buffer<A, T, TargetRoute, ReplyRoute>
230where
231    A: Address,
232    TargetRoute: DeliveryRoute,
233    TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
234    ReplyRoute: DeliveryRoute,
235    ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
236{
237    type Addr = A;
238    type Msg = BufferMessage<T, TargetRoute, ReplyRoute>;
239}
240
241impl<A, T, TargetRoute, ReplyRoute> Behavior for Buffer<A, T, TargetRoute, ReplyRoute>
242where
243    A: Address,
244    TargetRoute: DeliveryRoute,
245    TargetRoute::Protocol: Protocol<Addr = A, Msg = T>,
246    ReplyRoute: DeliveryRoute + Clone,
247    ReplyRoute::Protocol: Protocol<Addr = A, Msg = BufferOutcome<T>>,
248    TargetRoute::Sends: behavior::SendsFor<User<A, BufferMessage<T, TargetRoute, ReplyRoute>>>,
249    ReplyRoute::Sends: behavior::SendsFor<User<A, BufferMessage<T, TargetRoute, ReplyRoute>>>,
250{
251    type Protocol = Self;
252    type Event = User<A, behavior::BehaviorMessage<Self>>;
253    type Sends = DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>;
254    type Ph = Never;
255    type Error = Never;
256    type Birth = NoBirths;
257
258    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
259        match event.message {
260            BufferMessage::Offer { value, reply_to }
261                if self.state.queued.len() < self.state.capacity =>
262            {
263                self.state.queued.push_back(Buffered {
264                    value,
265                    reply_to: reply_to.clone(),
266                });
267                Ok(Self::actions(
268                    TargetRoute::Sends::empty(),
269                    reply_to.deliver(BufferOutcome::Accepted {
270                        depth: self.state.queued.len(),
271                    }),
272                ))
273            }
274            BufferMessage::Offer { value, reply_to } => match self.state.overflow {
275                OverflowPolicy::Reject => Ok(Self::actions(
276                    TargetRoute::Sends::empty(),
277                    reply_to.deliver(BufferOutcome::Rejected {
278                        value,
279                        reason: BufferRejection::Full,
280                    }),
281                )),
282                OverflowPolicy::DropNewest => Ok(Self::actions(
283                    TargetRoute::Sends::empty(),
284                    reply_to.deliver(BufferOutcome::Rejected {
285                        value,
286                        reason: BufferRejection::DroppedNewest,
287                    }),
288                )),
289                OverflowPolicy::DropOldest => {
290                    // `new` rejects zero capacity and this arm is reached only
291                    // after `len >= capacity`, so the queue is non-empty.
292                    let evicted = self
293                        .state
294                        .queued
295                        .pop_front()
296                        .expect("positive full buffer contains an oldest value");
297                    let mut outcomes = evicted.reply_to.deliver(BufferOutcome::Evicted {
298                        value: evicted.value,
299                    });
300                    self.state.queued.push_back(Buffered {
301                        value,
302                        reply_to: reply_to.clone(),
303                    });
304                    outcomes.append(reply_to.deliver(BufferOutcome::Accepted {
305                        depth: self.state.queued.len(),
306                    }));
307                    Ok(Self::actions(TargetRoute::Sends::empty(), outcomes))
308                }
309            },
310            BufferMessage::Release { to, reply_to } => {
311                let Some(buffered) = self.state.queued.pop_front() else {
312                    return Ok(Self::actions(
313                        TargetRoute::Sends::empty(),
314                        reply_to.deliver(BufferOutcome::Empty),
315                    ));
316                };
317                Ok(Self::actions(
318                    to.deliver(buffered.value),
319                    reply_to.deliver(BufferOutcome::Released {
320                        remaining: self.state.queued.len(),
321                    }),
322                ))
323            }
324        }
325    }
326}
327
328#[cfg(test)]
329mod tests {
330    use super::*;
331    use crate::Activate as _;
332    use behavior::{Delivery, MailAddr, Recipient};
333
334    fn active(
335        policy: OverflowPolicy,
336    ) -> crate::Active<
337        Buffer<
338            MailAddr,
339            u8,
340            Recipient<MessageProtocol<MailAddr, u8>>,
341            Recipient<MessageProtocol<MailAddr, BufferOutcome<u8>>>,
342        >,
343    > {
344        Buffer::new(BufferConfiguration::new(2, policy).unwrap())
345            .initialize()
346            .unwrap()
347            .behavior
348    }
349
350    #[test]
351    fn buffer_uses_the_shared_delivery_outcome_product() {
352        type Target = Recipient<MessageProtocol<MailAddr, u8>>;
353        type Reply = Recipient<MessageProtocol<MailAddr, BufferOutcome<u8>>>;
354        type Expected = crate::DeliveryOutcomes<
355            <Target as DeliveryRoute>::Sends,
356            <Reply as DeliveryRoute>::Sends,
357        >;
358        fn exact<B: Behavior<Sends = Expected>>() {}
359        exact::<Buffer<MailAddr, u8, Target, Reply>>();
360    }
361
362    #[test]
363    fn zero_capacity_is_rejected_before_ownership_is_possible() {
364        assert!(matches!(
365            BufferConfiguration::new(0, OverflowPolicy::Reject),
366            Err(BufferConfigError::ZeroCapacity)
367        ));
368    }
369
370    #[test]
371    fn fifo_release_and_empty_result_use_disjoint_named_lanes() {
372        let reply = Recipient::from(MailAddr(8));
373        let target = Recipient::from(MailAddr(7));
374        let mut buffer = active(OverflowPolicy::Reject);
375        for value in [1, 2] {
376            let accepted = buffer
377                .receive(
378                    MailAddr(9),
379                    BufferMessage::Offer {
380                        value,
381                        reply_to: reply,
382                    },
383                )
384                .unwrap();
385            assert!(accepted.sends.deliveries.is_empty());
386            assert!(matches!(
387                accepted.sends.outcomes.as_slice(),
388                [delivery]
389                    if delivery.message == (BufferOutcome::Accepted {
390                        depth: usize::from(value),
391                    })
392            ));
393        }
394        for (value, remaining) in [(1, 1), (2, 0)] {
395            let released = buffer
396                .receive(
397                    MailAddr(9),
398                    BufferMessage::Release {
399                        to: target,
400                        reply_to: reply,
401                    },
402                )
403                .unwrap();
404            assert!(released.sends.deliveries == vec![Delivery::new(target, value)]);
405            assert!(matches!(
406                released.sends.outcomes.as_slice(),
407                [delivery]
408                    if delivery.message == (BufferOutcome::Released { remaining })
409            ));
410        }
411        let empty = buffer
412            .receive(
413                MailAddr(9),
414                BufferMessage::Release {
415                    to: target,
416                    reply_to: reply,
417                },
418            )
419            .unwrap();
420        assert!(empty.sends.deliveries.is_empty());
421        assert!(matches!(
422            empty.sends.outcomes.as_slice(),
423            [delivery] if delivery.message == BufferOutcome::Empty
424        ));
425    }
426
427    #[test]
428    fn every_overflow_policy_preserves_or_returns_all_owned_values() {
429        let first_reply = Recipient::from(MailAddr(1));
430        let newest_reply = Recipient::from(MailAddr(2));
431        for policy in [
432            OverflowPolicy::Reject,
433            OverflowPolicy::DropNewest,
434            OverflowPolicy::DropOldest,
435        ] {
436            let mut buffer = active(policy);
437            for value in [10, 11] {
438                let offered = buffer
439                    .receive(
440                        MailAddr(9),
441                        BufferMessage::Offer {
442                            value,
443                            reply_to: first_reply,
444                        },
445                    )
446                    .unwrap();
447                assert!(offered.sends.deliveries.is_empty());
448                assert!(matches!(
449                    offered.sends.outcomes.as_slice(),
450                    [delivery]
451                        if delivery.message
452                            == (BufferOutcome::Accepted {
453                                depth: usize::from(value - 9),
454                            })
455                ));
456                assert!(offered.creates.is_empty());
457                assert_eq!(offered.become_, behavior::Step::Continue);
458            }
459            let overflow = buffer
460                .receive(
461                    MailAddr(9),
462                    BufferMessage::Offer {
463                        value: 12,
464                        reply_to: newest_reply,
465                    },
466                )
467                .unwrap();
468            match policy {
469                OverflowPolicy::Reject => assert!(matches!(
470                    overflow.sends.outcomes.as_slice(),
471                    [delivery]
472                        if delivery.message == (BufferOutcome::Rejected {
473                            value: 12,
474                            reason: BufferRejection::Full,
475                        })
476                )),
477                OverflowPolicy::DropNewest => assert!(matches!(
478                    overflow.sends.outcomes.as_slice(),
479                    [delivery]
480                        if delivery.message == (BufferOutcome::Rejected {
481                            value: 12,
482                            reason: BufferRejection::DroppedNewest,
483                        })
484                )),
485                OverflowPolicy::DropOldest => {
486                    assert!(matches!(
487                        overflow.sends.outcomes.as_slice(),
488                        [evicted, accepted]
489                            if evicted.message == (BufferOutcome::Evicted { value: 10 })
490                                && accepted.message == (BufferOutcome::Accepted { depth: 2 })
491                    ));
492                    assert!(buffer.state().queued().map(|item| item.value).eq([11, 12]));
493                }
494            }
495            assert_eq!(buffer.state().len(), 2);
496        }
497    }
498}