use std::{cmp::Reverse, collections::HashSet, time::Instant};
use priority_queue::PriorityQueue;
#[derive(Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Clone, Copy)]
pub enum EventKind {
Send,
Recv,
Connected,
Accept,
Closed,
StreamOpen,
StreamAccept,
StreamSend,
StreamRecv,
}
#[derive(Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Clone, Copy)]
pub struct Token(pub u32);
#[derive(Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Clone, Copy)]
pub struct Event {
pub kind: EventKind,
pub is_server: bool,
pub is_error: bool,
pub token: Token,
pub stream_id: u64,
}
#[derive(Default)]
pub struct Readiness {
pub(crate) delayed: PriorityQueue<Event, Reverse<Instant>>,
pub(crate) events: HashSet<Event>,
}
impl Readiness {
pub fn insert(&mut self, event: Event, delay_to: Option<Instant>) {
log::trace!("readiness, event={:?}, delay={:?}", event, delay_to);
if let Some(delay_to) = delay_to {
self.events.remove(&event);
self.delayed.push(event, Reverse(delay_to));
} else {
self.delayed.remove(&event);
self.events.insert(event);
}
}
pub fn remove(&mut self, event: Event) {
self.events.remove(&event);
self.delayed.remove(&event);
}
pub fn poll(&mut self, readiness: &mut Vec<Event>) -> Option<Instant> {
assert!(
readiness.is_empty(),
"The input readiness Vec<_> is not empty."
);
readiness.extend(self.events.drain());
if self.delayed.is_empty() {
return None;
}
let now = Instant::now();
while let Some(deadline) = self.delayed.peek().map(|(_, instant)| instant.0) {
if !(deadline > now) {
let (event, _) = self.delayed.pop().unwrap();
readiness.push(event);
continue;
}
return Some(deadline);
}
return None;
}
}
#[cfg(test)]
mod tests {
use std::vec;
use super::*;
#[test]
fn test_deduplication() {
let mut readiness = Readiness::default();
readiness.insert(
Event {
kind: EventKind::Send,
is_server: false,
is_error: false,
token: Token(0),
stream_id: 0,
},
None,
);
readiness.insert(
Event {
kind: EventKind::Send,
is_server: false,
is_error: false,
token: Token(0),
stream_id: 0,
},
None,
);
let mut events = vec![];
assert!(readiness.poll(&mut events).is_none());
assert_eq!(
events,
vec![Event {
kind: EventKind::Send,
is_server: false,
is_error: false,
token: Token(0),
stream_id: 0,
}]
);
}
#[test]
fn test_timeout() {
let delay_to = Instant::now();
let mut readiness = Readiness::default();
readiness.insert(
Event {
kind: EventKind::Send,
is_server: false,
is_error: false,
token: Token(0),
stream_id: 0,
},
Some(delay_to),
);
readiness.insert(
Event {
kind: EventKind::Send,
is_server: false,
is_error: false,
token: Token(0),
stream_id: 0,
},
Some(delay_to),
);
let mut events = vec![];
assert!(readiness.poll(&mut events).is_none());
assert_eq!(
events,
vec![Event {
kind: EventKind::Send,
is_server: false,
is_error: false,
token: Token(0),
stream_id: 0,
}]
);
}
}