Skip to main content

alopex_chirps/buffer/
priority_queue.rs

1use crate::MessageProfile;
2use std::collections::VecDeque;
3
4/// FIFO queues ordered by the profile priority Control > Durable > Ephemeral.
5#[derive(Debug)]
6pub struct PriorityQueue<T> {
7    control: VecDeque<T>,
8    durable: VecDeque<T>,
9    ephemeral: VecDeque<T>,
10}
11
12impl<T> Default for PriorityQueue<T> {
13    fn default() -> Self {
14        Self {
15            control: VecDeque::new(),
16            durable: VecDeque::new(),
17            ephemeral: VecDeque::new(),
18        }
19    }
20}
21
22impl<T> PriorityQueue<T> {
23    pub fn push(&mut self, profile: MessageProfile, value: T) {
24        match profile {
25            MessageProfile::Control => self.control.push_back(value),
26            MessageProfile::Durable => self.durable.push_back(value),
27            MessageProfile::Ephemeral => self.ephemeral.push_back(value),
28        }
29    }
30
31    pub fn pop(&mut self) -> Option<(MessageProfile, T)> {
32        self.control
33            .pop_front()
34            .map(|value| (MessageProfile::Control, value))
35            .or_else(|| {
36                self.durable
37                    .pop_front()
38                    .map(|value| (MessageProfile::Durable, value))
39            })
40            .or_else(|| {
41                self.ephemeral
42                    .pop_front()
43                    .map(|value| (MessageProfile::Ephemeral, value))
44            })
45    }
46
47    pub fn len(&self) -> usize {
48        self.control.len() + self.durable.len() + self.ephemeral.len()
49    }
50
51    pub fn is_empty(&self) -> bool {
52        self.control.is_empty() && self.durable.is_empty() && self.ephemeral.is_empty()
53    }
54}