1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
18pub enum PriorityQueueState {
19 Active {
21 next: u64,
23 queued: usize,
25 },
26 Exhausted {
28 queued: usize,
30 },
31}
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
35pub enum PriorityQueueRejection {
36 Full,
38 SequenceExhausted,
40}
41
42#[derive(Debug, PartialEq, Eq)]
44pub enum PriorityQueueOutcome<T, P> {
45 Accepted {
47 depth: usize,
49 },
50 Rejected {
52 value: T,
54 priority: P,
56 reason: PriorityQueueRejection,
58 },
59 Released {
61 remaining: usize,
63 },
64 Empty,
66}
67
68pub enum PriorityQueueMessage<T, P, TargetRoute, ReplyRoute> {
70 Offer {
72 value: T,
74 priority: P,
76 reply_to: ReplyRoute,
78 },
79 Release {
81 to: TargetRoute,
83 reply_to: ReplyRoute,
85 },
86}
87
88#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
90pub enum PriorityQueueConfigError {
91 #[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
120pub 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 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 #[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}