mise-server 0.1.11

MIcro SErvice
Documentation
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
            }
        });
    }
}