use crate::result::Result;
use std::sync::atomic::{AtomicU32, Ordering};
#[repr(C)]
#[derive(Clone, Copy, Default)]
pub(crate) struct io_uring_sqe {
pub opcode: u8,
pub flags: u8,
pub ioprio: u16,
pub fd: i32,
pub off_addr2: u64,
pub addr_splice_off_in: u64,
pub len: u32,
pub op_flags: u32,
pub user_data: u64,
pub buf_group: u16,
pub personality: u16,
pub splice_fd_in: i32,
pub addr3: u64,
pub __pad2: u64,
}
#[repr(C)]
#[derive(Clone, Copy, Default)]
pub(crate) struct io_uring_cqe {
pub user_data: u64,
pub res: i32,
pub flags: u32,
}
pub struct SubmissionRing {
entries: u32,
#[allow(dead_code)] sq: Vec<AtomicU32>,
#[allow(dead_code)] sqe: Vec<io_uring_sqe>,
pub head: Box<AtomicU32>,
pub tail: Box<AtomicU32>,
}
impl SubmissionRing {
pub fn new(entries: u32) -> Self {
assert!(
entries.is_power_of_two(),
"ring size must be a power of two"
);
let mut sq = Vec::with_capacity(entries as usize);
for _ in 0..entries {
sq.push(AtomicU32::new(0));
}
let mut sqe = Vec::with_capacity(entries as usize);
for _ in 0..entries {
sqe.push(io_uring_sqe::default());
}
Self {
entries,
sq,
sqe,
head: Box::new(AtomicU32::new(0)),
tail: Box::new(AtomicU32::new(0)),
}
}
pub fn entries(&self) -> u32 {
self.entries
}
pub fn free_slots(&self) -> u32 {
let head = self.head.load(Ordering::Acquire);
let tail = self.tail.load(Ordering::Acquire);
self.entries - (tail.wrapping_sub(head) & (self.entries - 1))
}
pub fn has_room(&self) -> bool {
self.free_slots() > 0
}
#[allow(dead_code)] pub(crate) fn sqe_at(&mut self, index: u32) -> &mut io_uring_sqe {
&mut self.sqe[index as usize]
}
pub fn publish(&self, count: u32) -> u32 {
self.tail.fetch_add(count, Ordering::Release)
}
}
pub struct CompletionRing {
entries: u32,
cqes: Vec<io_uring_cqe>,
pub head: Box<AtomicU32>,
pub tail: Box<AtomicU32>,
}
impl CompletionRing {
pub fn new(entries: u32) -> Self {
assert!(
entries.is_power_of_two(),
"ring size must be a power of two"
);
let mut cqes = Vec::with_capacity(entries as usize);
for _ in 0..entries {
cqes.push(io_uring_cqe::default());
}
Self {
entries,
cqes,
head: Box::new(AtomicU32::new(0)),
tail: Box::new(AtomicU32::new(0)),
}
}
pub fn entries(&self) -> u32 {
self.entries
}
pub fn available(&self) -> u32 {
let head = self.head.load(Ordering::Acquire);
let tail = self.tail.load(Ordering::Acquire);
tail.wrapping_sub(head) & (self.entries - 1)
}
pub fn has_completions(&self) -> bool {
self.available() > 0
}
pub fn peek(&self) -> Option<Result> {
if !self.has_completions() {
return None;
}
let head = self.head.load(Ordering::Acquire);
let index = (head & (self.entries - 1)) as usize;
let cqe = &self.cqes[index];
Some(Result::new(cqe.res as i64, cqe.user_data))
}
pub fn consume(&self) -> Option<Result> {
if !self.has_completions() {
return None;
}
let head = self.head.fetch_add(1, Ordering::AcqRel);
let index = (head & (self.entries - 1)) as usize;
let cqe = &self.cqes[index];
Some(Result::new(cqe.res as i64, cqe.user_data))
}
pub fn drain(&self) -> Vec<Result> {
let mut results = Vec::new();
while let Some(r) = self.consume() {
results.push(r);
}
results
}
}