use crate::{connection::Connection, routes::RequestProcessor};
use mio::{Events, Interest, Poll, Token};
use std::{
io::ErrorKind,
net::{SocketAddr, TcpListener},
};
use tracing::error;
pub(crate) struct Server {
processor: RequestProcessor,
poll: Poll,
events: Events,
listener: mio::net::TcpListener,
connections: Vec<Connection>,
}
static TOKEN: Token = Token(0);
impl Server {
pub(crate) fn serve(processor: RequestProcessor, addr: SocketAddr) -> Self {
let poll = Poll::new().expect("Cannot create mio Poller");
let events = Events::with_capacity(1024);
let listener = TcpListener::bind(addr).unwrap_or_else(|_| panic!("Could not bind {addr}"));
listener
.set_nonblocking(true)
.expect("Could not set listener to non blocking");
let mut listener = mio::net::TcpListener::from_std(listener);
poll.registry()
.register(
&mut listener,
TOKEN,
Interest::READABLE | Interest::WRITABLE,
)
.expect("Could not register listener to mio Poller");
Self {
processor,
poll,
events,
listener,
connections: vec![],
}
}
pub(crate) fn run(mut self) {
loop {
let mut had_io = false;
had_io |= self.accept();
had_io |= self.read_buffers();
self.handle_requests();
had_io |= self.write_buffers();
self.clean_connections();
if !had_io && let Err(e) = self.poll.poll(&mut self.events, None) {
if e.kind() == ErrorKind::Interrupted {
continue;
}
panic!("{e}");
}
}
}
fn accept(&mut self) -> bool {
let mut accepted = match self.listener.accept() {
Ok(accepted) => accepted,
Err(e) if e.kind() == ErrorKind::WouldBlock => {
return false;
}
Err(e) => {
error!("{e}");
return false;
}
};
self.poll
.registry()
.register(
&mut accepted.0,
TOKEN,
Interest::READABLE | Interest::WRITABLE,
)
.expect("Could not register socket to mio Poller");
self.connections.push(accepted.into());
true
}
fn write_buffers(&mut self) -> bool {
let mut network_activity = false;
for c in &mut self.connections {
if c.write_buffer() {
network_activity = true;
}
}
network_activity
}
fn read_buffers(&mut self) -> bool {
let mut network_activity = false;
for c in &mut self.connections {
if c.read_buffer() {
network_activity = true;
}
}
network_activity
}
fn handle_requests(&mut self) {
for c in &mut self.connections {
while let Some(req) = c.pop_request() {
c.send_response(self.processor.process(req));
}
}
}
fn clean_connections(&mut self) {
self.connections.retain_mut(|c| {
if c.is_closed() {
c.unregister(self.poll.registry());
false
} else {
true
}
});
}
}