Skip to main content

behavior_actors/routing/
priority_queue.rs

1//! Bounded stable priority release policy.
2
3use std::cmp::Ordering;
4use std::collections::BinaryHeap;
5
6use super::DeliveryOutcomes;
7
8use behavior::{
9    Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, Protocol,
10    SendEffects, User,
11};
12use thiserror::Error;
13
14use crate::DeliveryRoute;
15
16/// Complete admission phase of a [`PriorityQueue`].
17#[derive(Debug, Clone, Copy, PartialEq, Eq)]
18pub enum PriorityQueueState {
19    /// Further offers can receive a fresh tie-order token.
20    Active {
21        /// Next insertion token.
22        next: u64,
23        /// Current queue depth.
24        queued: usize,
25    },
26    /// The insertion-token domain is exhausted; retained values remain releasable.
27    Exhausted {
28        /// Current queue depth.
29        queued: usize,
30    },
31}
32
33/// Why an offer was not admitted.
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
35pub enum PriorityQueueRejection {
36    /// Positive capacity is full.
37    Full,
38    /// No fresh stable tie-order token is representable.
39    SequenceExhausted,
40}
41
42/// Factual result of one priority-queue operation.
43#[derive(Debug, PartialEq, Eq)]
44pub enum PriorityQueueOutcome<T, P> {
45    /// Value was accepted.
46    Accepted {
47        /// Queue depth after admission.
48        depth: usize,
49    },
50    /// Value was not accepted; ownership is returned.
51    Rejected {
52        /// Unaccepted value.
53        value: T,
54        /// Exact priority from the rejected offer.
55        priority: P,
56        /// Exhaustive reason.
57        reason: PriorityQueueRejection,
58    },
59    /// Highest-priority oldest-tie value was released.
60    Released {
61        /// Queue depth after release.
62        remaining: usize,
63    },
64    /// No retained value was available.
65    Empty,
66}
67
68/// Operations accepted by [`PriorityQueue`].
69pub enum PriorityQueueMessage<T, P, TargetRoute, ReplyRoute> {
70    /// Offer one owned value at an immutable priority.
71    Offer {
72        /// Owned value.
73        value: T,
74        /// Comparable priority; greater values release first.
75        priority: P,
76        /// Typed admission-result recipient.
77        reply_to: ReplyRoute,
78    },
79    /// Release one value to a typed destination.
80    Release {
81        /// Destination.
82        to: TargetRoute,
83        /// Typed release-result recipient.
84        reply_to: ReplyRoute,
85    },
86}
87
88/// Invalid priority-queue definition.
89#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
90pub enum PriorityQueueConfigError {
91    /// A zero-capacity definition can never accept ownership.
92    #[error("priority queue capacity must be positive")]
93    ZeroCapacity,
94}
95
96struct Entry<T, P> {
97    value: T,
98    priority: P,
99    order: u64,
100}
101impl<T, P: Ord> PartialEq for Entry<T, P> {
102    fn eq(&self, other: &Self) -> bool {
103        self.priority == other.priority && self.order == other.order
104    }
105}
106impl<T, P: Ord> Eq for Entry<T, P> {}
107impl<T, P: Ord> PartialOrd for Entry<T, P> {
108    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
109        Some(self.cmp(other))
110    }
111}
112impl<T, P: Ord> Ord for Entry<T, P> {
113    fn cmp(&self, other: &Self) -> Ordering {
114        self.priority
115            .cmp(&other.priority)
116            .then_with(|| other.order.cmp(&self.order))
117    }
118}
119
120/// Bounded stable immutable-priority queue behavior.
121///
122/// Greater priorities release first and equal priorities preserve FIFO by an
123/// explicit insertion token. Full and token-exhausted offers return ownership.
124/// Token exhaustion is a distinct admission phase and never wraps; retained
125/// values remain releasable. Release emits one value and one factual outcome;
126/// empty release emits only `Empty`. Initialization is empty, no actors are
127/// created, and the host never terminates by policy. Capacity, priority order,
128/// and FIFO ties are Bombay policy; mailbox priority and backpressure remain
129/// Communication concerns. The implementation uses `BinaryHeap` because
130/// accepted priorities are immutable; it does not need `priority-queue`.
131/// No transition has a semantic panic condition.
132pub struct PriorityQueue<
133    A: Address,
134    T,
135    P,
136    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
137    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = PriorityQueueOutcome<T, P>>>,
138> {
139    capacity: usize,
140    next: Option<u64>,
141    queued: BinaryHeap<Entry<T, P>>,
142    marker: core::marker::PhantomData<fn() -> (A, TargetRoute, ReplyRoute)>,
143}
144type PriorityActions<A, TargetSends, OutcomeSends> =
145    Actions<A, Never, DeliveryOutcomes<TargetSends, OutcomeSends>, NoBirths>;
146impl<A, T, P, TargetRoute, ReplyRoute> PriorityQueue<A, T, P, TargetRoute, ReplyRoute>
147where
148    A: Address,
149    P: Ord,
150    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
151    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = PriorityQueueOutcome<T, P>>>,
152{
153    /// Construct an empty positive-capacity queue.
154    /// # Errors
155    /// Returns [`PriorityQueueConfigError::ZeroCapacity`] for zero capacity.
156    pub fn new(capacity: usize) -> Result<Self, PriorityQueueConfigError> {
157        if capacity == 0 {
158            return Err(PriorityQueueConfigError::ZeroCapacity);
159        }
160        Ok(Self {
161            capacity,
162            next: Some(0),
163            queued: BinaryHeap::with_capacity(capacity),
164            marker: core::marker::PhantomData,
165        })
166    }
167    /// Return the complete admission phase and depth.
168    #[must_use]
169    pub fn state(&self) -> PriorityQueueState {
170        self.next.map_or(
171            PriorityQueueState::Exhausted {
172                queued: self.queued.len(),
173            },
174            |next| PriorityQueueState::Active {
175                next,
176                queued: self.queued.len(),
177            },
178        )
179    }
180    fn sends(
181        deliveries: TargetRoute::Sends,
182        outcomes: ReplyRoute::Sends,
183    ) -> PriorityActions<A, TargetRoute::Sends, ReplyRoute::Sends> {
184        Actions::send(DeliveryOutcomes {
185            deliveries,
186            outcomes,
187        })
188    }
189    fn offer(
190        &mut self,
191        value: T,
192        priority: P,
193        reply_to: ReplyRoute,
194    ) -> PriorityActions<A, TargetRoute::Sends, ReplyRoute::Sends> {
195        if self.queued.len() == self.capacity {
196            return Self::sends(
197                TargetRoute::Sends::empty(),
198                reply_to.deliver(PriorityQueueOutcome::Rejected {
199                    value,
200                    priority,
201                    reason: PriorityQueueRejection::Full,
202                }),
203            );
204        }
205        let Some(order) = self.next else {
206            return Self::sends(
207                TargetRoute::Sends::empty(),
208                reply_to.deliver(PriorityQueueOutcome::Rejected {
209                    value,
210                    priority,
211                    reason: PriorityQueueRejection::SequenceExhausted,
212                }),
213            );
214        };
215        self.next = order.checked_add(1);
216        self.queued.push(Entry {
217            value,
218            priority,
219            order,
220        });
221        Self::sends(
222            TargetRoute::Sends::empty(),
223            reply_to.deliver(PriorityQueueOutcome::Accepted {
224                depth: self.queued.len(),
225            }),
226        )
227    }
228}
229impl<A, T, P, TargetRoute, ReplyRoute> BehaviorBase
230    for PriorityQueue<A, T, P, TargetRoute, ReplyRoute>
231where
232    A: Address,
233    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
234    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = PriorityQueueOutcome<T, P>>>,
235{
236    type Base = Self;
237    fn base(&self) -> &Self {
238        self
239    }
240}
241impl<A, T, P, TargetRoute, ReplyRoute> behavior::Protocol
242    for PriorityQueue<A, T, P, TargetRoute, ReplyRoute>
243where
244    A: Address,
245    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
246    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = PriorityQueueOutcome<T, P>>>,
247{
248    type Addr = A;
249    type Msg = PriorityQueueMessage<T, P, TargetRoute, ReplyRoute>;
250}
251
252impl<A, T, P, TargetRoute, ReplyRoute> Behavior for PriorityQueue<A, T, P, TargetRoute, ReplyRoute>
253where
254    A: Address,
255    P: Ord,
256    TargetRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = T>>,
257    ReplyRoute: DeliveryRoute<Protocol: Protocol<Addr = A, Msg = PriorityQueueOutcome<T, P>>>,
258    TargetRoute::Sends:
259        behavior::SendsFor<User<A, PriorityQueueMessage<T, P, TargetRoute, ReplyRoute>>>,
260    ReplyRoute::Sends:
261        behavior::SendsFor<User<A, PriorityQueueMessage<T, P, TargetRoute, ReplyRoute>>>,
262{
263    type Protocol = Self;
264    type Event = User<A, behavior::BehaviorMessage<Self>>;
265    type Sends = DeliveryOutcomes<TargetRoute::Sends, ReplyRoute::Sends>;
266    type Ph = Never;
267    type Error = Never;
268    type Birth = NoBirths;
269    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
270        Ok(match event.message {
271            PriorityQueueMessage::Offer {
272                value,
273                priority,
274                reply_to,
275            } => self.offer(value, priority, reply_to),
276            PriorityQueueMessage::Release { to, reply_to } => match self.queued.pop() {
277                None => Self::sends(
278                    TargetRoute::Sends::empty(),
279                    reply_to.deliver(PriorityQueueOutcome::Empty),
280                ),
281                Some(entry) => Self::sends(
282                    to.deliver(entry.value),
283                    reply_to.deliver(PriorityQueueOutcome::Released {
284                        remaining: self.queued.len(),
285                    }),
286                ),
287            },
288        })
289    }
290}
291
292#[cfg(test)]
293mod tests {
294    use super::*;
295    use crate::Activate as _;
296    use behavior::{Delivery, MailAddr, Recipient};
297    struct Target;
298    struct Reply;
299    impl behavior::Protocol for Target {
300        type Addr = MailAddr;
301        type Msg = u8;
302    }
303
304    impl Behavior for Target {
305        type Protocol = Self;
306        type Event = User<MailAddr, u8>;
307        type Sends = Vec<Never>;
308        type Ph = Never;
309        type Error = Never;
310        type Birth = NoBirths;
311        fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
312            Ok(Actions::cont())
313        }
314    }
315    impl behavior::Protocol for Reply {
316        type Addr = MailAddr;
317        type Msg = PriorityQueueOutcome<u8, u8>;
318    }
319
320    impl Behavior for Reply {
321        type Protocol = Self;
322        type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
323        type Sends = Vec<Never>;
324        type Ph = Never;
325        type Error = Never;
326        type Birth = NoBirths;
327        fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
328            Ok(Actions::cont())
329        }
330    }
331    type Subject = PriorityQueue<MailAddr, u8, u8, Recipient<Target>, Recipient<Reply>>;
332    fn reply() -> Recipient<Reply> {
333        Recipient::global(MailAddr(2))
334    }
335    fn offer(s: &mut crate::Active<Subject>, value: u8, priority: u8) {
336        let offered = s
337            .receive(
338                MailAddr(0),
339                PriorityQueueMessage::Offer {
340                    value,
341                    priority,
342                    reply_to: reply(),
343                },
344            )
345            .unwrap();
346        assert!(offered.sends.deliveries.is_empty());
347        assert!(matches!(
348            offered.sends.outcomes.as_slice(),
349            [delivery] if matches!(delivery.message, PriorityQueueOutcome::Accepted { .. })
350        ));
351        assert!(offered.creates.is_empty());
352        assert_eq!(offered.become_, behavior::Step::Continue);
353    }
354    fn release(
355        s: &mut crate::Active<Subject>,
356    ) -> PriorityActions<MailAddr, Vec<Delivery<Target>>, Vec<Delivery<Reply>>> {
357        s.receive(
358            MailAddr(0),
359            PriorityQueueMessage::Release {
360                to: Recipient::global(MailAddr(1)),
361                reply_to: reply(),
362            },
363        )
364        .unwrap()
365    }
366    #[test]
367    fn greater_priority_and_fifo_ties_are_stable() {
368        let mut s = (Subject::new(4).unwrap()).initialize().unwrap().behavior;
369        for pair in [(1, 2), (2, 3), (3, 3), (4, 1)] {
370            offer(&mut s, pair.0, pair.1);
371        }
372        let released = [
373            release(&mut s),
374            release(&mut s),
375            release(&mut s),
376            release(&mut s),
377        ];
378        assert_eq!(
379            released.map(|actions| actions.sends.deliveries[0].message),
380            [2, 3, 1, 4]
381        );
382    }
383    #[test]
384    fn full_and_empty_are_explicit() {
385        let mut s = (Subject::new(1).unwrap()).initialize().unwrap().behavior;
386        offer(&mut s, 1, 0);
387        let rejected = s
388            .receive(
389                MailAddr(0),
390                PriorityQueueMessage::Offer {
391                    value: 2,
392                    priority: 9,
393                    reply_to: reply(),
394                },
395            )
396            .unwrap();
397        assert!(matches!(
398            rejected.sends.outcomes[0].message,
399            PriorityQueueOutcome::Rejected {
400                value: 2,
401                priority: 9,
402                reason: PriorityQueueRejection::Full
403            }
404        ));
405        let released = release(&mut s);
406        assert_eq!(released.sends.deliveries[0].message, 1);
407        assert!(matches!(
408            released.sends.outcomes[0].message,
409            PriorityQueueOutcome::Released { remaining: 0 }
410        ));
411        let empty = release(&mut s);
412        assert!(matches!(
413            empty.sends.outcomes[0].message,
414            PriorityQueueOutcome::Empty
415        ));
416    }
417
418    #[test]
419    fn insertion_sequence_exhaustion_never_wraps_or_consumes_the_next_value() {
420        let mut definition = Subject::new(2).unwrap();
421        definition.next = Some(u64::MAX);
422        let mut subject = (definition).initialize().unwrap().behavior;
423        offer(&mut subject, 1, 0);
424        assert_eq!(subject.state(), PriorityQueueState::Exhausted { queued: 1 });
425        let rejected = subject
426            .receive(
427                MailAddr(0),
428                PriorityQueueMessage::Offer {
429                    value: 2,
430                    priority: 9,
431                    reply_to: reply(),
432                },
433            )
434            .unwrap();
435        assert!(matches!(
436            rejected.sends.outcomes[0].message,
437            PriorityQueueOutcome::Rejected {
438                value: 2,
439                priority: 9,
440                reason: PriorityQueueRejection::SequenceExhausted
441            }
442        ));
443    }
444}