use std::cmp::{Ordering, Reverse};
use std::collections::BinaryHeap;
use crate::types::NanoTime;
#[derive(Debug)]
struct Entry<T> {
time: NanoTime,
seq: u64,
value: T,
}
impl<T> Entry<T> {
fn key(&self) -> (NanoTime, u64) {
(self.time, self.seq)
}
}
impl<T> PartialEq for Entry<T> {
fn eq(&self, other: &Self) -> bool {
self.key() == other.key()
}
}
impl<T> Eq for Entry<T> {}
impl<T> PartialOrd for Entry<T> {
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl<T> Ord for Entry<T> {
fn cmp(&self, other: &Self) -> Ordering {
self.key().cmp(&other.key())
}
}
#[derive(Debug)]
pub(crate) struct TimeQueue<T> {
heap: BinaryHeap<Reverse<Entry<T>>>,
next_seq: u64,
}
impl<T> Default for TimeQueue<T> {
fn default() -> Self {
Self {
heap: BinaryHeap::new(),
next_seq: 0,
}
}
}
impl<T> TimeQueue<T> {
pub fn new() -> Self {
Self::default()
}
pub fn next_time(&self) -> Option<NanoTime> {
self.heap.peek().map(|Reverse(e)| e.time)
}
pub fn is_empty(&self) -> bool {
self.heap.is_empty()
}
pub fn pop(&mut self) -> Option<T> {
self.heap.pop().map(|Reverse(e)| e.value)
}
pub fn pop_if_pending(&mut self, current_time: NanoTime) -> Option<T> {
match self.next_time() {
Some(t) if t <= current_time => self.pop(),
_ => None,
}
}
pub fn clear(&mut self) {
self.heap.clear();
}
}
impl<T: PartialEq> TimeQueue<T> {
pub fn push(&mut self, value: T, time: NanoTime) {
if self
.heap
.iter()
.any(|Reverse(e)| e.time == time && e.value == value)
{
return;
}
let seq = self.next_seq;
self.next_seq += 1;
self.heap.push(Reverse(Entry { time, seq, value }));
}
}
#[cfg(test)]
mod tests {
use super::TimeQueue;
use crate::time::NanoTime;
#[test]
fn identical_pushes_are_deduplicated() {
let mut queue: TimeQueue<u32> = TimeQueue::new();
queue.push(1, NanoTime::new(100));
queue.push(1, NanoTime::new(100));
queue.push(1, NanoTime::new(100));
assert_eq!(queue.pop(), Some(1));
assert!(queue.is_empty());
}
#[test]
fn same_value_distinct_times_all_kept() {
let mut queue: TimeQueue<u32> = TimeQueue::new();
queue.push(1, NanoTime::new(100));
queue.push(1, NanoTime::new(200));
queue.push(1, NanoTime::new(300));
assert_eq!(queue.pop(), Some(1));
assert_eq!(queue.pop(), Some(1));
assert_eq!(queue.pop(), Some(1));
assert!(queue.is_empty());
}
#[test]
fn distinct_values_same_time_pop_in_fifo_order() {
let mut queue: TimeQueue<u32> = TimeQueue::new();
queue.push(1, NanoTime::new(100));
queue.push(2, NanoTime::new(100));
queue.push(3, NanoTime::new(100));
assert_eq!(queue.pop(), Some(1));
assert_eq!(queue.pop(), Some(2));
assert_eq!(queue.pop(), Some(3));
assert!(queue.is_empty());
}
#[test]
fn float_payloads_are_supported_and_deduplicated() {
let mut queue: TimeQueue<f64> = TimeQueue::new();
queue.push(1.5, NanoTime::new(200));
queue.push(1.5, NanoTime::new(200)); queue.push(2.5, NanoTime::new(100));
assert_eq!(queue.pop(), Some(2.5));
assert_eq!(queue.pop(), Some(1.5));
assert!(queue.is_empty());
}
#[test]
fn sorted() {
let mut queue: TimeQueue<u32> = TimeQueue::new();
queue.push(1, NanoTime::new(300));
queue.push(3, NanoTime::new(100));
queue.push(2, NanoTime::new(200));
assert_eq!(queue.pop(), Some(3));
assert_eq!(queue.pop(), Some(2));
assert_eq!(queue.pop(), Some(1));
assert!(queue.is_empty());
}
#[test]
fn pop_if_pending() {
let mut queue: TimeQueue<u32> = TimeQueue::new();
assert_eq!(queue.pop_if_pending(NanoTime::MAX), None);
queue.push(1, NanoTime::new(100));
assert_eq!(queue.pop_if_pending(NanoTime::new(99)), None);
assert_eq!(queue.pop_if_pending(NanoTime::new(100)), Some(1));
assert!(queue.is_empty());
}
#[test]
fn next_time_is_none_when_empty() {
let queue: TimeQueue<u32> = TimeQueue::new();
assert_eq!(queue.next_time(), None);
}
#[test]
fn pop_is_none_when_empty() {
let mut queue: TimeQueue<u32> = TimeQueue::new();
assert_eq!(queue.pop(), None);
}
}