macro_rules! resultify {
($ret:expr) => {{
let ret = $ret;
match ret >= 0 {
true => Ok(ret as _),
false => Err(std::io::Error::from_raw_os_error(-ret)),
}
}};
}
use crate::uring_sys::uring_sys;
use std::io;
use std::mem::MaybeUninit;
use std::ptr::{self, NonNull};
use std::time::Duration;
pub(crate) use crate::iou::cqe::{CompletionQueue, CompletionQueueEvent};
pub(crate) use crate::iou::registrar::Registrar;
pub(crate) use crate::iou::sqe::{
FallocateFlags, FsyncFlags, OpenMode, SockAddrStorage, StatxFlags, StatxMode, SubmissionFlags,
SubmissionQueue, SubmissionQueueEvent,
};
pub(crate) use nix::poll::PollFlags;
pub(crate) use nix::sys::socket::{
AlgAddr, InetAddr, LinkAddr, NetlinkAddr, SockAddr, SockFlag, UnixAddr, VsockAddr,
};
bitflags::bitflags! {
pub(crate) struct SetupFlags: u32 {
const IOPOLL = 1 << 0;
const SQPOLL = 1 << 1;
const SQ_AFF = 1 << 2;
}
}
pub(crate) struct IoUring {
pub(crate) ring: uring_sys::io_uring,
}
impl IoUring {
pub(crate) fn new(entries: u32) -> io::Result<IoUring> {
IoUring::new_with_flags(entries, SetupFlags::empty())
}
pub(crate) fn new_with_flags(entries: u32, flags: SetupFlags) -> io::Result<IoUring> {
unsafe {
let mut ring = MaybeUninit::uninit();
let _: i32 = resultify! {
uring_sys::io_uring_queue_init(entries as _, ring.as_mut_ptr(), flags.bits() as _)
}?;
Ok(IoUring {
ring: ring.assume_init(),
})
}
}
pub(crate) fn sq(&mut self) -> SubmissionQueue<'_> {
SubmissionQueue::new(&*self)
}
pub(crate) fn cq(&mut self) -> CompletionQueue<'_> {
CompletionQueue::new(&*self)
}
pub(crate) fn registrar(&self) -> Registrar<'_> {
Registrar::new(self)
}
pub(crate) fn queues(&mut self) -> (SubmissionQueue<'_>, CompletionQueue<'_>, Registrar<'_>) {
(
SubmissionQueue::new(&*self),
CompletionQueue::new(&*self),
Registrar::new(&*self),
)
}
pub(crate) fn next_sqe(&mut self) -> Option<SubmissionQueueEvent<'_>> {
unsafe {
let sqe = uring_sys::io_uring_get_sqe(&mut self.ring);
if sqe != ptr::null_mut() {
let mut sqe = SubmissionQueueEvent::new(&mut *sqe);
sqe.clear();
Some(sqe)
} else {
None
}
}
}
pub(crate) fn submit_sqes(&mut self) -> io::Result<usize> {
self.sq().submit()
}
pub(crate) fn submit_sqes_and_wait(&mut self, wait_for: u32) -> io::Result<usize> {
self.sq().submit_and_wait(wait_for)
}
pub(crate) fn submit_sqes_and_wait_with_timeout(
&mut self,
wait_for: u32,
duration: Duration,
) -> io::Result<usize> {
self.sq().submit_and_wait_with_timeout(wait_for, duration)
}
pub(crate) fn peek_for_cqe(&mut self) -> Option<CompletionQueueEvent> {
unsafe {
let mut cqe = MaybeUninit::uninit();
let count = uring_sys::io_uring_peek_batch_cqe(&mut self.ring, cqe.as_mut_ptr(), 1);
if count > 0 {
Some(CompletionQueueEvent::new(
NonNull::from(&self.ring),
&mut *cqe.assume_init(),
))
} else {
None
}
}
}
pub(crate) fn wait_for_cqe(&mut self) -> io::Result<CompletionQueueEvent> {
self.inner_wait_for_cqes(1, ptr::null())
}
pub(crate) fn wait_for_cqe_with_timeout(
&mut self,
duration: Duration,
) -> io::Result<CompletionQueueEvent> {
let ts = uring_sys::__kernel_timespec {
tv_sec: duration.as_secs() as _,
tv_nsec: duration.subsec_nanos() as _,
};
self.inner_wait_for_cqes(1, &ts)
}
pub(crate) fn wait_for_cqes(&mut self, count: usize) -> io::Result<CompletionQueueEvent> {
self.inner_wait_for_cqes(count as _, ptr::null())
}
pub(crate) fn wait_for_cqes_with_timeout(
&mut self,
count: usize,
duration: Duration,
) -> io::Result<CompletionQueueEvent> {
let ts = uring_sys::__kernel_timespec {
tv_sec: duration.as_secs() as _,
tv_nsec: duration.subsec_nanos() as _,
};
self.inner_wait_for_cqes(count as _, &ts)
}
fn inner_wait_for_cqes(
&mut self,
count: u32,
ts: *const uring_sys::__kernel_timespec,
) -> io::Result<CompletionQueueEvent> {
unsafe {
let mut cqe = MaybeUninit::uninit();
let _: i32 = resultify!(uring_sys::io_uring_wait_cqes(
&mut self.ring,
cqe.as_mut_ptr(),
count,
ts,
ptr::null(),
))?;
Ok(CompletionQueueEvent::new(
NonNull::from(&self.ring),
&mut *cqe.assume_init(),
))
}
}
pub(crate) fn raw(&self) -> &uring_sys::io_uring {
&self.ring
}
pub(crate) fn raw_mut(&mut self) -> &mut uring_sys::io_uring {
&mut self.ring
}
}
impl Drop for IoUring {
fn drop(&mut self) {
unsafe { uring_sys::io_uring_queue_exit(&mut self.ring) };
}
}
unsafe impl Send for IoUring {}
unsafe impl Sync for IoUring {}
#[cfg(test)]
mod tests {
#[test]
fn test_resultify() {
let side_effect = |i, effect: &mut _| -> i32 {
*effect += 1;
return i;
};
let mut calls = 0;
let ret: Result<i32, _> = resultify!(side_effect(0, &mut calls));
assert!(match ret {
Ok(0) => true,
_ => false,
});
assert_eq!(calls, 1);
calls = 0;
let ret: Result<i32, _> = resultify!(side_effect(1, &mut calls));
assert!(match ret {
Ok(1) => true,
_ => false,
});
assert_eq!(calls, 1);
calls = 0;
let ret: Result<i32, _> = resultify!(side_effect(-1, &mut calls));
assert!(match ret {
Err(e) if e.raw_os_error() == Some(1) => true,
_ => false,
});
assert_eq!(calls, 1);
}
}