use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Debug, Clone)]
pub struct IoUringConfig {
pub ring_size: u32,
pub sqpoll: bool,
pub sqpoll_idle_ms: u32,
pub iopoll: bool,
pub single_issuer: bool,
pub defer_taskrun: bool,
pub fixed_buffers: usize,
pub buffer_size: usize,
pub buffer_ring: bool,
pub buffer_ring_entries: u32,
}
impl Default for IoUringConfig {
fn default() -> Self {
Self {
ring_size: 4096,
sqpoll: false,
sqpoll_idle_ms: 1000,
iopoll: false,
single_issuer: true,
defer_taskrun: true,
fixed_buffers: 1024,
buffer_size: 16384, buffer_ring: true,
buffer_ring_entries: 4096,
}
}
}
impl IoUringConfig {
pub fn builder() -> IoUringConfigBuilder {
IoUringConfigBuilder::default()
}
pub fn high_performance() -> Self {
Self {
ring_size: 8192,
sqpoll: true,
sqpoll_idle_ms: 2000,
iopoll: false,
single_issuer: true,
defer_taskrun: true,
fixed_buffers: 2048,
buffer_size: 32768, buffer_ring: true,
buffer_ring_entries: 8192,
}
}
pub fn balanced() -> Self {
Self::default()
}
pub fn low_resource() -> Self {
Self {
ring_size: 1024,
sqpoll: false,
sqpoll_idle_ms: 500,
iopoll: false,
single_issuer: true,
defer_taskrun: true,
fixed_buffers: 256,
buffer_size: 8192, buffer_ring: false,
buffer_ring_entries: 1024,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct IoUringConfigBuilder {
config: IoUringConfig,
}
impl IoUringConfigBuilder {
pub fn ring_size(mut self, size: u32) -> Self {
self.config.ring_size = size.next_power_of_two();
self
}
pub fn sqpoll(mut self, enable: bool) -> Self {
self.config.sqpoll = enable;
self
}
pub fn sqpoll_idle_ms(mut self, ms: u32) -> Self {
self.config.sqpoll_idle_ms = ms;
self
}
pub fn iopoll(mut self, enable: bool) -> Self {
self.config.iopoll = enable;
self
}
pub fn single_issuer(mut self, enable: bool) -> Self {
self.config.single_issuer = enable;
self
}
pub fn defer_taskrun(mut self, enable: bool) -> Self {
self.config.defer_taskrun = enable;
self
}
pub fn fixed_buffers(mut self, count: usize) -> Self {
self.config.fixed_buffers = count;
self
}
pub fn buffer_size(mut self, size: usize) -> Self {
self.config.buffer_size = size;
self
}
pub fn buffer_ring(mut self, enable: bool) -> Self {
self.config.buffer_ring = enable;
self
}
pub fn build(self) -> IoUringConfig {
self.config
}
}
#[derive(Debug, Default)]
pub struct IoUringStats {
submissions: AtomicU64,
completions: AtomicU64,
sq_full: AtomicU64,
cq_overflow: AtomicU64,
bytes_read: AtomicU64,
bytes_written: AtomicU64,
accepts: AtomicU64,
errors: AtomicU64,
ring_utilization: AtomicU64,
}
impl IoUringStats {
pub fn new() -> Self {
Self::default()
}
#[inline]
pub fn record_submission(&self) {
self.submissions.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_completion(&self) {
self.completions.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_sq_full(&self) {
self.sq_full.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_cq_overflow(&self) {
self.cq_overflow.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_read(&self, bytes: u64) {
self.bytes_read.fetch_add(bytes, Ordering::Relaxed);
}
#[inline]
pub fn record_write(&self, bytes: u64) {
self.bytes_written.fetch_add(bytes, Ordering::Relaxed);
}
#[inline]
pub fn record_accept(&self) {
self.accepts.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_error(&self) {
self.errors.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn update_utilization(&self, percent: u64) {
self.ring_utilization.store(percent, Ordering::Relaxed);
}
pub fn submissions(&self) -> u64 {
self.submissions.load(Ordering::Relaxed)
}
pub fn completions(&self) -> u64 {
self.completions.load(Ordering::Relaxed)
}
pub fn sq_full(&self) -> u64 {
self.sq_full.load(Ordering::Relaxed)
}
pub fn cq_overflow(&self) -> u64 {
self.cq_overflow.load(Ordering::Relaxed)
}
pub fn bytes_read(&self) -> u64 {
self.bytes_read.load(Ordering::Relaxed)
}
pub fn bytes_written(&self) -> u64 {
self.bytes_written.load(Ordering::Relaxed)
}
pub fn accepts(&self) -> u64 {
self.accepts.load(Ordering::Relaxed)
}
pub fn errors(&self) -> u64 {
self.errors.load(Ordering::Relaxed)
}
pub fn ring_utilization(&self) -> u64 {
self.ring_utilization.load(Ordering::Relaxed)
}
pub fn pending(&self) -> u64 {
self.submissions().saturating_sub(self.completions())
}
}
#[cfg(all(target_os = "linux", feature = "io-uring"))]
pub fn is_available() -> bool {
::io_uring::IoUring::new(2).is_ok()
}
#[cfg(all(target_os = "linux", not(feature = "io-uring")))]
pub fn is_available() -> bool {
if let Ok(version) = std::fs::read_to_string("/proc/version")
&& let Some(ver) = parse_kernel_version(&version)
{
return ver >= (5, 1);
}
false
}
#[cfg(not(target_os = "linux"))]
pub fn is_available() -> bool {
false
}
#[cfg(target_os = "linux")]
#[cfg_attr(feature = "io-uring", allow(dead_code))]
fn parse_kernel_version(version_str: &str) -> Option<(u32, u32)> {
let parts: Vec<&str> = version_str.split_whitespace().collect();
if parts.len() >= 3 && parts[0] == "Linux" && parts[1] == "version" {
let ver_parts: Vec<&str> = parts[2].split('.').collect();
if ver_parts.len() >= 2 {
let major = ver_parts[0].parse().ok()?;
let minor_str = ver_parts[1].split('-').next()?;
let minor = minor_str.parse().ok()?;
return Some((major, minor));
}
}
None
}
#[derive(Debug, Clone)]
pub struct IoUringFeatures {
pub basic: bool,
pub sqpoll: bool,
pub buffer_ring: bool,
pub multishot_accept: bool,
pub zerocopy: bool,
pub fixed_files: bool,
}
impl IoUringFeatures {
#[cfg(all(target_os = "linux", feature = "io-uring"))]
pub fn detect() -> Self {
use ::io_uring::{IoUring, Probe, opcode};
let Ok(ring) = IoUring::new(2) else {
return Self::unsupported();
};
let mut probe = Probe::new();
let probed = ring.submitter().register_probe(&mut probe).is_ok();
let supported = |code: u8| probed && probe.is_supported(code);
let sqpoll_ring: std::io::Result<IoUring> = IoUring::builder().setup_sqpoll(100).build(2);
Self {
basic: true,
sqpoll: sqpoll_ring.is_ok(),
buffer_ring: supported(opcode::Socket::CODE),
multishot_accept: supported(opcode::Socket::CODE),
zerocopy: supported(opcode::SendZc::CODE),
fixed_files: supported(opcode::ReadFixed::CODE),
}
}
#[cfg(all(target_os = "linux", not(feature = "io-uring")))]
pub fn detect() -> Self {
let version = std::fs::read_to_string("/proc/version")
.ok()
.and_then(|v| parse_kernel_version(&v))
.unwrap_or((0, 0));
Self {
basic: version >= (5, 1),
sqpoll: version >= (5, 11),
buffer_ring: version >= (5, 19),
multishot_accept: version >= (5, 19),
zerocopy: version >= (6, 0),
fixed_files: version >= (5, 1),
}
}
#[cfg(not(target_os = "linux"))]
pub fn detect() -> Self {
Self::unsupported()
}
#[cfg(any(not(target_os = "linux"), feature = "io-uring"))]
fn unsupported() -> Self {
Self {
basic: false,
sqpoll: false,
buffer_ring: false,
multishot_accept: false,
zerocopy: false,
fixed_files: false,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum IoBackend {
#[default]
Epoll,
IoUring,
Auto,
}
impl IoBackend {
pub fn resolve(self) -> Self {
match self {
Self::Auto => {
if is_available() {
Self::IoUring
} else {
Self::Epoll
}
}
other => other,
}
}
pub fn is_io_uring(self) -> bool {
matches!(self.resolve(), Self::IoUring)
}
}
#[cfg(all(target_os = "linux", feature = "io-uring"))]
pub struct IoUringRuntime {
ring: std::sync::Mutex<::io_uring::IoUring>,
stats: IoUringStats,
}
#[cfg(all(target_os = "linux", feature = "io-uring"))]
impl std::fmt::Debug for IoUringRuntime {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("IoUringRuntime")
.field("stats", &self.stats)
.finish_non_exhaustive()
}
}
#[cfg(all(target_os = "linux", feature = "io-uring"))]
impl IoUringRuntime {
pub fn is_available() -> bool {
is_available()
}
pub fn new(config: IoUringConfig) -> std::io::Result<Self> {
let entries = config.ring_size.max(1).next_power_of_two().min(32768);
let build = |sqpoll: bool| {
let mut builder = ::io_uring::IoUring::builder();
if config.iopoll {
builder.setup_iopoll();
}
if sqpoll {
builder.setup_sqpoll(config.sqpoll_idle_ms);
}
builder.build(entries)
};
let ring = match build(config.sqpoll) {
Ok(ring) => ring,
Err(_) if config.sqpoll => build(false)?,
Err(e) => return Err(e),
};
Ok(Self {
ring: std::sync::Mutex::new(ring),
stats: IoUringStats::new(),
})
}
pub fn stats(&self) -> &IoUringStats {
&self.stats
}
pub fn read_at(
&self,
fd: std::os::fd::BorrowedFd<'_>,
buf: &mut [u8],
offset: u64,
) -> std::io::Result<usize> {
use std::os::fd::AsRawFd;
let len = u32::try_from(buf.len()).unwrap_or(u32::MAX);
let sqe = ::io_uring::opcode::Read::new(
::io_uring::types::Fd(fd.as_raw_fd()),
buf.as_mut_ptr(),
len,
)
.offset(offset)
.build();
let n = unsafe { self.submit_one(sqe) }?;
self.stats.record_read(n as u64);
Ok(n)
}
pub fn write_at(
&self,
fd: std::os::fd::BorrowedFd<'_>,
buf: &[u8],
offset: u64,
) -> std::io::Result<usize> {
use std::os::fd::AsRawFd;
let len = u32::try_from(buf.len()).unwrap_or(u32::MAX);
let sqe = ::io_uring::opcode::Write::new(
::io_uring::types::Fd(fd.as_raw_fd()),
buf.as_ptr(),
len,
)
.offset(offset)
.build();
let n = unsafe { self.submit_one(sqe) }?;
self.stats.record_write(n as u64);
Ok(n)
}
unsafe fn submit_one(&self, sqe: ::io_uring::squeue::Entry) -> std::io::Result<usize> {
let mut ring = self.ring.lock().unwrap();
unsafe {
let mut sq = ring.submission();
if sq.push(&sqe).is_err() {
self.stats.record_sq_full();
drop(sq);
ring.submit()?;
let mut sq = ring.submission();
sq.push(&sqe)
.map_err(|_| std::io::Error::other("io_uring submission queue full"))?;
}
}
self.stats.record_submission();
loop {
match ring.submit_and_wait(1) {
Ok(_) => break,
Err(e) if e.raw_os_error() == Some(libc::EINTR) => continue,
Err(e) => {
self.stats.record_error();
return Err(e);
}
}
}
let cqe = ring
.completion()
.next()
.ok_or_else(|| std::io::Error::other("io_uring completion queue empty after wait"))?;
self.stats.record_completion();
let res = cqe.result();
if res < 0 {
self.stats.record_error();
Err(std::io::Error::from_raw_os_error(-res))
} else {
Ok(res as usize)
}
}
}
#[derive(Debug)]
pub struct BufferPool {
buffers: Vec<std::sync::Mutex<Vec<u8>>>,
free_list: std::sync::Mutex<Vec<usize>>,
buffer_size: usize,
stats: BufferPoolStats,
}
#[derive(Debug, Default)]
struct BufferPoolStats {
allocations: AtomicU64,
deallocations: AtomicU64,
pool_misses: AtomicU64,
}
impl BufferPool {
pub fn new(count: usize, buffer_size: usize) -> Self {
let mut buffers = Vec::with_capacity(count);
let mut free_list = Vec::with_capacity(count);
for i in 0..count {
buffers.push(std::sync::Mutex::new(vec![0u8; buffer_size]));
free_list.push(i);
}
Self {
buffers,
free_list: std::sync::Mutex::new(free_list),
buffer_size,
stats: BufferPoolStats::default(),
}
}
pub fn acquire(&self) -> Option<BufferGuard<'_>> {
let idx = self.free_list.lock().unwrap().pop();
if let Some(idx) = idx {
self.stats.allocations.fetch_add(1, Ordering::Relaxed);
let buf = self.buffers[idx].lock().unwrap();
Some(BufferGuard {
pool: self,
index: idx,
buf,
})
} else {
self.stats.pool_misses.fetch_add(1, Ordering::Relaxed);
None
}
}
pub fn buffer_size(&self) -> usize {
self.buffer_size
}
pub fn capacity(&self) -> usize {
self.buffers.len()
}
pub fn available(&self) -> usize {
self.free_list.lock().unwrap().len()
}
}
#[derive(Debug)]
pub struct BufferGuard<'a> {
pool: &'a BufferPool,
index: usize,
buf: std::sync::MutexGuard<'a, Vec<u8>>,
}
impl BufferGuard<'_> {
pub fn index(&self) -> usize {
self.index
}
}
impl std::ops::Deref for BufferGuard<'_> {
type Target = [u8];
fn deref(&self) -> &[u8] {
&self.buf
}
}
impl std::ops::DerefMut for BufferGuard<'_> {
fn deref_mut(&mut self) -> &mut [u8] {
&mut self.buf
}
}
impl Drop for BufferGuard<'_> {
fn drop(&mut self) {
self.pool.free_list.lock().unwrap().push(self.index);
self.pool
.stats
.deallocations
.fetch_add(1, Ordering::Relaxed);
}
}
#[derive(Debug, Clone)]
pub struct TcpOptions {
pub nodelay: bool,
pub reuseaddr: bool,
pub reuseport: bool,
pub keepalive_secs: Option<u32>,
pub send_buffer: Option<usize>,
pub recv_buffer: Option<usize>,
pub backlog: u32,
}
impl Default for TcpOptions {
fn default() -> Self {
Self {
nodelay: true,
reuseaddr: true,
reuseport: true,
keepalive_secs: Some(60),
send_buffer: Some(65536),
recv_buffer: Some(65536),
backlog: 1024,
}
}
}
impl TcpOptions {
pub fn high_performance() -> Self {
Self {
nodelay: true,
reuseaddr: true,
reuseport: true,
keepalive_secs: Some(120),
send_buffer: Some(262144), recv_buffer: Some(262144),
backlog: 4096,
}
}
pub fn low_latency() -> Self {
Self {
nodelay: true,
reuseaddr: true,
reuseport: false,
keepalive_secs: Some(30),
send_buffer: Some(32768),
recv_buffer: Some(32768),
backlog: 512,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IoOp {
Accept,
Read,
Write,
Close,
Connect,
SendZC,
RecvBuf,
Poll,
Timeout,
Cancel,
Link,
Nop,
}
impl IoOp {
pub fn opcode(self) -> u8 {
match self {
Self::Accept => 13,
Self::Read => 22,
Self::Write => 23,
Self::Close => 19,
Self::Connect => 16,
Self::SendZC => 52,
Self::RecvBuf => 58,
Self::Poll => 6,
Self::Timeout => 11,
Self::Cancel => 14,
Self::Link => 255, Self::Nop => 0,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_config_builder() {
let config = IoUringConfig::builder()
.ring_size(2048)
.sqpoll(true)
.buffer_size(32768)
.build();
assert_eq!(config.ring_size, 2048);
assert!(config.sqpoll);
assert_eq!(config.buffer_size, 32768);
}
#[test]
fn test_high_performance_config() {
let config = IoUringConfig::high_performance();
assert_eq!(config.ring_size, 8192);
assert!(config.sqpoll);
}
#[test]
fn test_stats() {
let stats = IoUringStats::new();
stats.record_submission();
stats.record_submission();
stats.record_completion();
assert_eq!(stats.submissions(), 2);
assert_eq!(stats.completions(), 1);
assert_eq!(stats.pending(), 1);
stats.record_read(1024);
stats.record_write(2048);
assert_eq!(stats.bytes_read(), 1024);
assert_eq!(stats.bytes_written(), 2048);
}
#[test]
fn test_io_backend() {
let backend = IoBackend::Auto;
let resolved = backend.resolve();
assert!(matches!(resolved, IoBackend::Epoll | IoBackend::IoUring));
}
#[test]
fn test_buffer_pool() {
let pool = BufferPool::new(10, 1024);
assert_eq!(pool.capacity(), 10);
assert_eq!(pool.available(), 10);
assert_eq!(pool.buffer_size(), 1024);
let mut buf = pool.acquire().unwrap();
assert_eq!(buf.len(), 1024);
assert_eq!(pool.available(), 9);
buf[0] = 42;
drop(buf);
assert_eq!(pool.available(), 10);
}
#[test]
fn test_buffer_pool_exclusive_buffers() {
let pool = BufferPool::new(2, 64);
let a = pool.acquire().unwrap();
let b = pool.acquire().unwrap();
assert_ne!(a.index(), b.index());
assert_eq!(pool.available(), 0);
assert!(pool.acquire().is_none());
drop(a);
drop(b);
assert_eq!(pool.available(), 2);
let c = pool.acquire().unwrap();
assert_eq!(c.len(), 64);
}
#[test]
fn test_tcp_options() {
let opts = TcpOptions::high_performance();
assert!(opts.nodelay);
assert!(opts.reuseport);
assert_eq!(opts.backlog, 4096);
}
#[test]
fn test_features_detect() {
let features = IoUringFeatures::detect();
let _ = features.basic;
}
#[cfg(target_os = "linux")]
#[test]
fn test_parse_kernel_version() {
let version = "Linux version 5.15.0-generic (buildd@lcy02-amd64-086)";
let parsed = parse_kernel_version(version);
assert_eq!(parsed, Some((5, 15)));
let version2 = "Linux version 6.1.0-18-amd64 (debian-kernel@lists.debian.org)";
let parsed2 = parse_kernel_version(version2);
assert_eq!(parsed2, Some((6, 1)));
}
#[test]
fn test_io_op_opcodes() {
assert_eq!(IoOp::Nop.opcode(), 0);
assert_eq!(IoOp::Accept.opcode(), 13);
assert_eq!(IoOp::Read.opcode(), 22);
}
#[cfg(all(target_os = "linux", feature = "io-uring"))]
mod runtime_tests {
use super::super::*;
use std::io::Read as _;
use std::os::fd::AsFd;
fn runtime_or_skip(config: IoUringConfig) -> Option<IoUringRuntime> {
match IoUringRuntime::new(config) {
Ok(rt) => Some(rt),
Err(e) => {
eprintln!("skipping io_uring test: kernel refused ring creation: {e}");
None
}
}
}
#[test]
fn test_write_read_round_trip() {
let config = IoUringConfig::builder().ring_size(64).build();
let Some(runtime) = runtime_or_skip(config) else {
return;
};
let path = std::env::temp_dir().join(format!(
"armature-io-uring-roundtrip-{}",
std::process::id()
));
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(true)
.open(&path)
.expect("create temp file");
let data = b"hello from io_uring: the quick brown fox";
let written = runtime
.write_at(file.as_fd(), data, 0)
.expect("write_at failed");
assert_eq!(written, data.len());
let mut buf = vec![0u8; data.len()];
let read = runtime
.read_at(file.as_fd(), &mut buf, 0)
.expect("read_at failed");
assert_eq!(read, data.len());
assert_eq!(&buf[..], &data[..]);
let mut verify = Vec::new();
let mut reopened = std::fs::File::open(&path).expect("reopen temp file");
reopened.read_to_end(&mut verify).expect("std read");
assert_eq!(&verify[..], &data[..]);
assert_eq!(runtime.stats().submissions(), 2);
assert_eq!(runtime.stats().completions(), 2);
assert_eq!(runtime.stats().pending(), 0);
assert_eq!(runtime.stats().errors(), 0);
assert_eq!(runtime.stats().bytes_written(), data.len() as u64);
assert_eq!(runtime.stats().bytes_read(), data.len() as u64);
drop(file);
drop(reopened);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_read_at_eof_returns_zero() {
let Some(runtime) = runtime_or_skip(IoUringConfig::low_resource()) else {
return;
};
let path =
std::env::temp_dir().join(format!("armature-io-uring-eof-{}", std::process::id()));
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(true)
.open(&path)
.expect("create temp file");
let mut buf = [0u8; 16];
let read = runtime
.read_at(file.as_fd(), &mut buf, 0)
.expect("read_at failed");
assert_eq!(read, 0, "empty file should read 0 bytes (EOF)");
drop(file);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_error_result_recorded_in_stats() {
let Some(runtime) = runtime_or_skip(IoUringConfig::default()) else {
return;
};
let file = std::fs::File::open("/dev/null").expect("open /dev/null");
let err = runtime
.write_at(file.as_fd(), b"nope", 0)
.expect_err("write to read-only fd should fail");
assert_eq!(err.raw_os_error(), Some(libc::EBADF));
assert_eq!(runtime.stats().errors(), 1);
assert_eq!(runtime.stats().submissions(), 1);
assert_eq!(runtime.stats().completions(), 1);
assert_eq!(runtime.stats().bytes_written(), 0);
}
#[test]
fn test_is_available_probes() {
let can_build = IoUringRuntime::new(IoUringConfig::low_resource()).is_ok();
assert_eq!(IoUringRuntime::is_available(), can_build);
}
#[test]
fn test_features_detect_probed() {
let features = IoUringFeatures::detect();
if !features.basic {
eprintln!("skipping: io_uring not available for feature probing");
return;
}
assert!(features.fixed_files);
}
#[test]
fn test_sqpoll_config_falls_back() {
let config = IoUringConfig::builder().ring_size(32).sqpoll(true).build();
let Some(runtime) = runtime_or_skip(config) else {
return;
};
let _ = runtime.stats();
}
}
}