pub mod mio_engine;
#[cfg(all(target_os = "linux", feature = "io-uring"))]
pub mod uring;
use std::io;
use std::net::SocketAddr;
use std::time::Duration;
use crate::buffer::BufferPool;
use crate::token::Token;
macro_rules! tracing_slow_warn {
($($arg:tt)*) => {
eprintln!("[vane:warn] {}", format_args!($($arg)*))
};
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Poll {
Done(u32),
Pending,
}
#[derive(Debug)]
pub struct Cqe {
pub token: Token,
pub result: io::Result<u32>,
}
pub trait Engine {
fn kind(&self) -> &'static str;
fn add_listener(&mut self, fd: std::os::fd::RawFd, token: Token) -> io::Result<()>;
fn add_stream(&mut self, fd: std::os::fd::RawFd, token: Token) -> io::Result<()>;
fn read(&mut self, token: Token, fd: std::os::fd::RawFd, slot: u32) -> io::Result<Poll>;
fn write(
&mut self,
token: Token,
fd: std::os::fd::RawFd,
slot: u32,
len: usize,
offset: usize,
) -> io::Result<Poll>;
fn connect(
&mut self,
token: Token,
addr: std::net::SocketAddr,
) -> io::Result<(std::os::fd::RawFd, Poll)>;
fn connect_unix(
&mut self,
token: Token,
path: &std::path::Path,
) -> io::Result<(std::os::fd::RawFd, Poll)>;
fn accept(
&mut self,
lfd: std::os::fd::RawFd,
ltoken: Token,
) -> io::Result<Option<(std::os::fd::RawFd, SocketAddr)>>;
fn splice_pump(&mut self, a: Token, afd: i32, b: Token, bfd: i32) -> io::Result<()>;
fn remove(&mut self, fd: std::os::fd::RawFd);
fn poll(&mut self, timeout: Option<Duration>, out: &mut Vec<Cqe>) -> io::Result<()>;
fn take_accepted(&mut self, fd: std::os::fd::RawFd) -> Option<SocketAddr>;
}
pub fn create_engine(
entries: u32,
buffers: Option<&BufferPool>,
sqpoll: bool,
) -> io::Result<Box<dyn Engine>> {
#[cfg(all(target_os = "linux", feature = "io-uring"))]
{
match uring::UringEngine::new(entries, buffers, sqpoll) {
Ok(e) => return Ok(Box::new(e)),
Err(err) => {
tracing_slow_warn!("io_uring unavailable ({}), falling back to mio", err);
}
}
}
let _ = (entries, sqpoll);
let fallback = buffers.ok_or_else(|| {
io::Error::new(
io::ErrorKind::InvalidInput,
"mio fallback requires a buffer pool",
)
})?;
Ok(Box::new(mio_engine::MioEngine::new(fallback)?))
}
#[cfg(test)]
mod create_tests {
use super::*;
#[test]
fn creates_engine_with_pool() {
let pool = BufferPool::new(64, crate::buffer::DEFAULT_BUF_SIZE).expect("pool");
let engine = create_engine(64, Some(&pool), false).expect("engine");
let _dyn: Box<dyn Engine> = engine;
}
#[test]
fn mio_fallback_requires_pool() {
let res = create_engine(8, None, false);
if let Err(e) = res {
assert_eq!(
e.kind(),
io::ErrorKind::InvalidInput,
"fallback error must be InvalidInput"
);
}
}
#[test]
fn zero_entries_still_builds_mio() {
let pool = BufferPool::new(64, crate::buffer::DEFAULT_BUF_SIZE).expect("pool");
let _ = create_engine(0, Some(&pool), false);
}
}