messagebus 0.15.2

MessageBus allows intercommunicate with messages between modules
Documentation
use std::{cmp::Ordering, collections::BinaryHeap};

struct Entry<M> {
    inner: M,
    index: u64,
}

impl<M> PartialOrd for Entry<M> {
    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
        Some(other.index.cmp(&self.index))
    }
}

impl<M> Ord for Entry<M> {
    fn cmp(&self, other: &Self) -> Ordering {
        other.index.cmp(&self.index)
    }
}

impl<M> PartialEq for Entry<M> {
    fn eq(&self, other: &Self) -> bool {
        other.index.eq(&self.index)
    }
}

impl<M> Eq for Entry<M> {}

#[derive(Debug)]
pub enum QueueItem<M> {
    Item(u64, M),
    Drop(u64, M, u64),
}
impl<M> QueueItem<M> {
    fn item((index, inner): (u64, M)) -> QueueItem<M> {
        Self::Item(index, inner)
    }
}

pub(crate) struct ReorderQueueInner<M> {
    cap: usize,
    started: bool,
    recent_index: u64,
    heap: BinaryHeap<Entry<M>>,
}

impl<M> ReorderQueueInner<M> {
    pub fn new(cap: usize, start_index: u64) -> Self {
        Self {
            cap,
            started: false,
            recent_index: start_index,
            heap: BinaryHeap::with_capacity(cap + 1),
        }
    }

    pub fn push(&mut self, index: u64, inner: M) -> Option<QueueItem<M>> {
        if self.started {
            if self.recent_index >= index {
                return Some(QueueItem::Drop(index, inner, self.recent_index));
            }
        } else {
            if index == self.recent_index {
                self.recent_index = index;
                self.started = true;
                return Some(QueueItem::Item(index, inner));
            }
        }

        self.heap.push(Entry { inner, index });

        if self.heap.len() > self.cap {
            self.force_pop().map(QueueItem::item)
        } else {
            self.try_pop().map(QueueItem::item)
        }
    }

    #[inline]
    pub fn try_pop(&mut self) -> Option<(u64, M)> {
        if self.started {
            let e = self.heap.peek()?;

            if e.index == self.recent_index + 1 {
                self.force_pop()
            } else {
                None
            }
        } else {
            None
        }
    }

    #[inline]
    pub fn force_pop(&mut self) -> Option<(u64, M)> {
        let e = self.heap.pop()?;
        self.recent_index = e.index;
        self.started = true;

        Some((e.index, e.inner))
    }
}

#[cfg(test)]
mod tests {
    use crate::{reorder_queue::QueueItem, Message};

    use super::ReorderQueueInner;

    impl Message for i32 {}

    #[test]
    fn test_reordering() {
        let mut queue = ReorderQueueInner::new(8, 0);
        assert!(matches!(queue.push(0, 0), Some(QueueItem::Item(0, 0))));
        assert!(queue.try_pop().is_none());

        assert!(queue.push(3, 3).is_none());
        assert!(queue.try_pop().is_none());

        assert!(queue.push(2, 2).is_none());
        assert!(queue.try_pop().is_none());

        assert!(queue.push(4, 4).is_none());
        assert!(queue.try_pop().is_none());

        assert!(matches!(queue.push(1, 1), Some(QueueItem::Item(1, 1))));
        assert!(matches!(queue.try_pop(), Some((2, 2))));
        assert!(matches!(queue.try_pop(), Some((3, 3))));
        assert!(matches!(queue.try_pop(), Some((4, 4))));
        assert!(matches!(queue.try_pop(), None));
    }

    #[test]
    fn test_overflow() {
        let mut queue = ReorderQueueInner::new(3, 0);
        assert!(matches!(queue.push(0, 0), Some(QueueItem::Item(0, 0))));
        assert!(queue.try_pop().is_none());

        assert!(queue.push(4, 4).is_none());
        assert!(queue.try_pop().is_none());

        assert!(queue.push(2, 2).is_none());
        assert!(queue.try_pop().is_none());

        assert!(queue.push(3, 3).is_none());
        assert!(queue.try_pop().is_none());

        assert!(matches!(queue.push(5, 5), Some(QueueItem::Item(2, 2))));
        assert!(matches!(queue.try_pop(), Some((3, 3))));
        assert!(matches!(queue.try_pop(), Some((4, 4))));
        assert!(matches!(queue.try_pop(), Some((5, 5))));
        assert!(matches!(queue.try_pop(), None));
    }
}