use std::io;
use std::time::Duration;
mod epoll;
mod iocp;
mod iouring;
mod kqueue;
pub use epoll::EpollReactor;
pub use iocp::IocpReactor;
pub use iouring::IoUringReactor;
pub use kqueue::KqueueReactor;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct Token(pub usize);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Interest {
pub readable: bool,
pub writable: bool,
}
impl Interest {
pub const READABLE: Interest = Interest {
readable: true,
writable: false,
};
pub const WRITABLE: Interest = Interest {
readable: false,
writable: true,
};
pub const BOTH: Interest = Interest {
readable: true,
writable: true,
};
}
#[derive(Debug, Clone)]
pub struct Event {
pub token: Token,
pub readable: bool,
pub writable: bool,
pub error: bool,
pub closed: bool,
}
impl Event {
pub fn new(token: Token) -> Self {
Self {
token,
readable: false,
writable: false,
error: false,
closed: false,
}
}
pub fn with_readable(mut self) -> Self {
self.readable = true;
self
}
pub fn with_writable(mut self) -> Self {
self.writable = true;
self
}
pub fn with_error(mut self) -> Self {
self.error = true;
self
}
pub fn with_closed(mut self) -> Self {
self.closed = true;
self
}
}
#[derive(Debug)]
pub struct Completion {
pub token: Token,
pub result: io::Result<usize>,
}
pub trait Reactor: Send + Sync {
fn poll(&mut self, timeout: Option<Duration>) -> io::Result<Vec<Event>>;
fn register(&mut self, fd: RawFd, interest: Interest) -> io::Result<Token>;
fn modify(&mut self, token: Token, interest: Interest) -> io::Result<()>;
fn deregister(&mut self, token: Token) -> io::Result<()>;
fn submit_read(&mut self, _token: Token, _buf: &mut [u8]) -> io::Result<Option<Completion>> {
Ok(None)
}
fn submit_write(&mut self, _token: Token, _buf: &[u8]) -> io::Result<Option<Completion>> {
Ok(None)
}
fn supports_async_io(&self) -> bool {
false
}
fn name(&self) -> &'static str;
}
#[cfg(unix)]
pub type RawFd = std::os::unix::io::RawFd;
#[cfg(windows)]
pub type RawFd = std::os::windows::io::RawSocket;
#[derive(Debug, Clone)]
pub struct ReactorConfig {
pub max_events: usize,
pub prefer_io_uring: bool,
}
impl Default for ReactorConfig {
fn default() -> Self {
Self {
max_events: 1024,
prefer_io_uring: true,
}
}
}
pub fn create_reactor(config: ReactorConfig) -> io::Result<Box<dyn Reactor>> {
#[cfg(target_os = "linux")]
{
if config.prefer_io_uring {
match IoUringReactor::new(config.max_events) {
Ok(reactor) => return Ok(Box::new(reactor)),
Err(_) => {
}
}
}
Ok(Box::new(EpollReactor::new(config.max_events)?))
}
#[cfg(target_os = "macos")]
{
Ok(Box::new(KqueueReactor::new(config.max_events)?))
}
#[cfg(target_os = "windows")]
{
Ok(Box::new(IocpReactor::new(config.max_events)?))
}
#[cfg(not(any(target_os = "linux", target_os = "macos", target_os = "windows")))]
{
Err(io::Error::new(
io::ErrorKind::Unsupported,
"No reactor implementation for this platform",
))
}
}
pub fn create_default_reactor() -> io::Result<Box<dyn Reactor>> {
create_reactor(ReactorConfig::default())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_interest_constants() {
assert!(Interest::READABLE.readable);
assert!(!Interest::READABLE.writable);
assert!(!Interest::WRITABLE.readable);
assert!(Interest::WRITABLE.writable);
assert!(Interest::BOTH.readable);
assert!(Interest::BOTH.writable);
}
#[test]
fn test_event_builder() {
let event = Event::new(Token(42)).with_readable().with_writable();
assert_eq!(event.token, Token(42));
assert!(event.readable);
assert!(event.writable);
assert!(!event.error);
assert!(!event.closed);
}
#[test]
fn test_token_equality() {
assert_eq!(Token(1), Token(1));
assert_ne!(Token(1), Token(2));
}
#[test]
fn test_reactor_config_default() {
let config = ReactorConfig::default();
assert_eq!(config.max_events, 1024);
assert!(config.prefer_io_uring);
}
#[test]
fn test_create_default_reactor() {
let result = create_default_reactor();
assert!(result.is_ok());
let reactor = result.unwrap();
let name = reactor.name();
assert!(!name.is_empty());
}
#[test]
fn test_create_reactor_with_config() {
let config = ReactorConfig {
max_events: 512,
prefer_io_uring: false, };
let result = create_reactor(config);
assert!(result.is_ok());
}
}