use std::cell::{Cell, RefCell};
use std::collections::VecDeque;
use std::io;
use std::rc::Rc;
use std::task::Poll;
use std::time::Instant;
use io_uring::{IoUring, opcode, squeue};
use crate::metrics::Counters;
use crate::park::Unpark;
use crate::{timer, udp};
pub(crate) type Task = Box<dyn FnMut(&kio::Waiter) -> Poll<()>>;
#[derive(Clone, Copy)]
pub(crate) struct Cqe {
pub user_data: u64,
pub result: i32,
pub flags: u32,
}
pub(crate) enum Op {
Recv {
sock: Rc<udp::SockShared>,
one: Option<Box<udp::OneshotRecv>>,
},
Send(udp::SendOp),
FutexWait,
Cancel,
}
pub(crate) struct Shared {
pub ring: RefCell<IoUring>,
pub ops: RefCell<slab::Slab<Op>>,
pub timers: Rc<RefCell<timer::Heap>>,
pub spawns: RefCell<Vec<Task>>,
pub unpark: std::sync::Arc<Unpark>,
pub metrics: std::sync::Arc<Counters>,
pub next_bgid: Cell<u16>,
pub stopped: Cell<bool>,
pub spill: RefCell<VecDeque<Cqe>>,
}
impl Shared {
pub fn gone_error() -> io::Error {
io::Error::new(io::ErrorKind::NotConnected, "the worker was dropped")
}
pub fn push(&self, entry: &squeue::Entry) -> io::Result<()> {
self.push_inner(entry, None)
}
fn push_until(&self, entry: &squeue::Entry, deadline: Instant) -> io::Result<()> {
self.push_inner(entry, Some(deadline))
}
fn push_inner(&self, entry: &squeue::Entry, deadline: Option<Instant>) -> io::Result<()> {
let mut ring = self.ring.borrow_mut();
loop {
if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
return Err(io::Error::new(
io::ErrorKind::TimedOut,
"io_uring submission deadline elapsed",
));
}
{
let mut sq = ring.submission();
if unsafe { sq.push(entry) }.is_ok() {
return Ok(());
}
}
self.metrics.enters.add(1);
let submitted = ring.submit();
if let Ok(count) = &submitted {
self.metrics.submissions.add(*count as u64);
}
if let Err(err) = submitted {
if deadline.is_some() && err.raw_os_error() == Some(libc::EINTR) {
continue;
}
if err.raw_os_error() != Some(libc::EBUSY) {
return Err(err);
}
let mut spill = self.spill.borrow_mut();
let before = spill.len();
spill.extend(ring.completion().map(|entry| Cqe {
user_data: entry.user_data(),
result: entry.result(),
flags: entry.flags(),
}));
if spill.len() == before {
return Err(err);
}
}
}
}
pub fn cancel(&self, target: u64) -> io::Result<()> {
self.cancel_inner(target, None)
}
pub fn cancel_until(&self, target: u64, deadline: Instant) -> io::Result<()> {
self.cancel_inner(target, Some(deadline))
}
fn cancel_inner(&self, target: u64, deadline: Option<Instant>) -> io::Result<()> {
let key = self.insert(Op::Cancel);
let entry = opcode::AsyncCancel::new(target).build().user_data(key);
let result = match deadline {
Some(deadline) => self.push_until(&entry, deadline),
None => self.push(&entry),
};
if let Err(err) = result {
self.ops.borrow_mut().remove(key as usize);
return Err(err);
}
Ok(())
}
pub fn insert(&self, op: Op) -> u64 {
self.ops.borrow_mut().insert(op) as u64
}
}