use mio::{Poll, Token, Waker};
use std::io;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{
mpsc::{self, Receiver, Sender},
Arc,
};
#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct EventId(pub usize);
impl From<usize> for EventId {
fn from(id: usize) -> Self {
EventId(id)
}
}
impl From<EventId> for usize {
fn from(event_id: EventId) -> Self {
event_id.0
}
}
pub struct WakeQSender<T> {
tx: Sender<T>,
waker: Arc<Waker>,
}
impl<T> WakeQSender<T> {
pub fn send_event(&self, event: T) -> io::Result<()> {
self.tx.send(event).map_err(|err| {
io::Error::new(
io::ErrorKind::Other,
format!("Failed to send message: {}", err),
)
})?;
self.waker.wake()
}
}
pub struct WakeQ<T> {
rx: Receiver<T>,
sender: WakeQSender<T>,
}
impl<T> WakeQ<T> {
pub fn new(poll: &Poll, token: Token) -> io::Result<Self> {
let (tx, rx) = mpsc::channel();
let waker = Waker::new(poll.registry(), token)?;
let sender = WakeQSender {
tx,
waker: Arc::new(waker),
};
Ok(WakeQ { rx, sender })
}
pub fn get_sender(&self) -> WakeQSender<T> {
self.sender.clone()
}
pub fn iter_pending_events(&self) -> impl Iterator<Item = T> + '_ {
std::iter::from_fn(|| self.rx.try_recv().ok())
}
}
impl<T> Clone for WakeQSender<T> {
fn clone(&self) -> Self {
WakeQSender {
tx: self.tx.clone(),
waker: Arc::clone(&self.waker),
}
}
}
#[derive(Clone)]
pub struct EventQSender {
events: Arc<Vec<AtomicUsize>>,
waker: Arc<Waker>,
}
impl EventQSender {
pub fn trigger_event(&self, event_id: EventId) -> io::Result<()> {
let event_id = usize::from(event_id);
let (chunk_idx, bit_idx) = Self::event_position(event_id);
self.events[chunk_idx].fetch_or(1 << bit_idx, Ordering::SeqCst);
self.waker.wake()
}
fn event_position(event_id: usize) -> (usize, u32) {
let chunk_idx = event_id / usize::BITS as usize;
let bit_idx = (event_id % usize::BITS as usize) as u32;
(chunk_idx, bit_idx)
}
}
pub struct EventQ {
events: Arc<Vec<AtomicUsize>>,
sender: EventQSender,
}
impl EventQ {
pub fn new(poll: &Poll, token: Token, num_events: usize) -> io::Result<Self> {
let num_chunks = (num_events + usize::BITS as usize - 1) / usize::BITS as usize;
let events = Arc::new((0..num_chunks).map(|_| AtomicUsize::new(0)).collect());
let waker = Arc::new(Waker::new(poll.registry(), token)?);
let sender = EventQSender {
events: Arc::clone(&events),
waker,
};
Ok(EventQ { events, sender })
}
pub fn get_sender(&self) -> EventQSender {
self.sender.clone()
}
pub fn triggered_events(&self) -> EventQIterator {
let triggered_events = self
.events
.iter()
.map(|chunk| chunk.swap(0, Ordering::SeqCst))
.collect();
EventQIterator {
events: triggered_events,
chunk_idx: 0,
bit_idx: 0,
}
}
}
pub struct EventQIterator {
events: Vec<usize>,
chunk_idx: usize,
bit_idx: u32,
}
impl Iterator for EventQIterator {
type Item = EventId;
fn next(&mut self) -> Option<Self::Item> {
while self.chunk_idx < self.events.len() {
let chunk = self.events[self.chunk_idx];
while self.bit_idx < usize::BITS {
let bit = 1 << self.bit_idx;
let bit_position = self.chunk_idx * usize::BITS as usize + self.bit_idx as usize;
self.bit_idx += 1;
if chunk & bit != 0 {
return Some(EventId(bit_position));
}
}
self.chunk_idx += 1;
self.bit_idx = 0;
}
None
}
}