use std::{
cmp::Reverse,
collections::HashSet,
time::{Duration, Instant},
};
use priority_queue::PriorityQueue;
use crate::poll::utils::is_bidi;
#[derive(Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Clone, Copy)]
pub enum EventKind {
Send,
Recv,
Connected,
Accept,
Closed,
StreamOpenBidi,
StreamOpenUni,
StreamAccept,
StreamSend,
StreamRecv,
ReadLock,
}
#[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 token: Token,
pub kind: EventKind,
pub is_server: bool,
pub is_error: bool,
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!("insert 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>,
release_timer_threshold: Duration,
) -> 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.checked_duration_since(now).unwrap_or_default() > release_timer_threshold)
{
let (event, _) = self.delayed.pop().unwrap();
readiness.push(event);
continue;
}
return Some(deadline);
}
return None;
}
}
#[derive(Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Clone, Copy)]
pub enum StreamKind {
Uni,
Bidi,
}
impl From<u64> for StreamKind {
fn from(value: u64) -> Self {
if is_bidi(value) {
StreamKind::Bidi
} else {
StreamKind::Uni
}
}
}