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,
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 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 ring_utilization(&self) -> u64 {
self.ring_utilization.load(Ordering::Relaxed)
}
pub fn pending(&self) -> u64 {
self.submissions().saturating_sub(self.completions())
}
}
#[cfg(target_os = "linux")]
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")]
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(target_os = "linux")]
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 {
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)
}
}
#[derive(Debug)]
pub struct BufferPool {
buffers: Vec<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(vec![0u8; buffer_size]);
free_list.push(i);
}
Self {
buffers,
free_list: std::sync::Mutex::new(free_list),
buffer_size,
stats: BufferPoolStats::default(),
}
}
#[allow(clippy::mut_from_ref)] pub fn acquire(&self) -> Option<(usize, &mut [u8])> {
let mut free = self.free_list.lock().unwrap();
if let Some(idx) = free.pop() {
self.stats.allocations.fetch_add(1, Ordering::Relaxed);
let buf = unsafe {
let ptr = self.buffers.as_ptr().add(idx) as *mut Vec<u8>;
(*ptr).as_mut_slice()
};
Some((idx, buf))
} else {
self.stats.pool_misses.fetch_add(1, Ordering::Relaxed);
None
}
}
pub fn release(&self, idx: usize) {
if idx < self.buffers.len() {
let mut free = self.free_list.lock().unwrap();
free.push(idx);
self.stats.deallocations.fetch_add(1, Ordering::Relaxed);
}
}
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, 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 (idx, buf) = pool.acquire().unwrap();
assert_eq!(buf.len(), 1024);
assert_eq!(pool.available(), 9);
pool.release(idx);
assert_eq!(pool.available(), 10);
}
#[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);
}
}