pub mod utils;
pub mod route;
pub mod request;
pub mod response;
#[cfg(test)]
#[macro_use]
extern crate serde_derive;
use std::str::FromStr;
use std::io::Result;
use std::io::prelude::*;
use std::net::ToSocketAddrs;
use std::collections::HashMap;
use std::collections::HashSet;
use threadpool::ThreadPool;
use mio::tcp::{TcpListener, TcpStream};
use mio::util::Slab;
use mio::*;
pub use crate::request::*;
pub use crate::response::*;
struct Client {
sock: TcpStream,
token: Token,
events: EventSet,
i_buf: Vec<u8>,
o_buf: Vec<u8>,
}
impl Client {
fn new(sock: TcpStream, token: Token) -> Client {
Client {
sock,
token,
events: EventSet::hup(),
i_buf: Vec::with_capacity(2048),
o_buf: Vec::new(),
}
}
fn receive(&mut self) -> Result<bool> {
let mut bytes_read: usize = 0;
loop {
let mut buf: Vec<u8> = Vec::with_capacity(2048);
match self.sock.try_read_buf(&mut buf) {
Ok(size) => {
match size {
Some(bytes) => {
self.i_buf.extend(buf);
bytes_read += bytes;
},
None => {
self.events.remove(EventSet::readable());
self.events.insert(EventSet::writable());
break;
},
}
},
Err(_) => {
self.events.remove(EventSet::readable());
self.events.insert(EventSet::writable());
break;
},
};
}
Ok(bytes_read > 0)
}
fn send(&mut self) -> Result<bool> {
if self.o_buf.is_empty() {
return Ok(false);
}
while !self.o_buf.is_empty() {
match self.sock.write(&self.o_buf.as_slice()) {
Ok(sz) => {
if sz == self.o_buf.len() {
self.events.remove(EventSet::writable());
break;
} else {
self.o_buf = self.o_buf.split_off(sz);
}
},
Err(_) => {
return Ok(true);
}
}
}
Ok(true)
}
fn register(&mut self, evl: &mut EventLoop<Canteen>) -> Result<()> {
self.events.insert(EventSet::readable());
evl.register(&self.sock, self.token, self.events, PollOpt::edge() | PollOpt::oneshot())
}
fn reregister(&mut self, evl: &mut EventLoop<Canteen>) -> Result<()> {
evl.reregister(&self.sock, self.token, self.events, PollOpt::edge() | PollOpt::oneshot())
}
}
pub struct Canteen {
routes: HashMap<route::RouteDef, route::Route>,
rcache: HashMap<route::RouteDef, route::RouteDef>,
server: Option<TcpListener>,
token: Token,
conns: Slab<Client>,
default: fn(&Request) -> Response,
tpool: ThreadPool,
}
impl Handler for Canteen {
type Timeout = ();
type Message = (Token, Vec<u8>);
fn ready(&mut self, evl: &mut EventLoop<Canteen>, token: Token, events: EventSet) {
if events.is_error() || events.is_hup() {
self.reset_connection(token);
return;
}
if events.is_readable() {
if self.token == token {
let sock = self.accept().unwrap();
if let Some(token) = self.conns.insert_with(|token| Client::new(sock, token)) {
self.get_client(token).register(evl).ok();
}
self.reregister(evl);
} else {
self.readable(evl, token)
.and_then(|_| self.get_client(token)
.reregister(evl)).ok();
}
return;
}
if events.is_writable() {
match self.get_client(token).send() {
Ok(true) => { self.reset_connection(token); },
Ok(false) => { let _ = self.get_client(token).reregister(evl); },
Err(_) => {},
}
}
}
fn notify(&mut self, evl: &mut EventLoop<Canteen>, msg: (Token, Vec<u8>)) {
let (token, output) = msg;
let client = self.get_client(token);
client.o_buf = output;
let _ = client.reregister(evl);
}
}
impl Canteen {
pub fn new() -> Canteen {
Canteen {
routes: HashMap::new(),
rcache: HashMap::new(),
server: None,
token: Token(1),
conns: Slab::new_starting_at(Token(2), 2048),
default: utils::err_404,
tpool: ThreadPool::new(255),
}
}
pub fn bind<A: ToSocketAddrs>(&mut self, addr: A) {
self.server = Some(TcpListener::bind(&addr.to_socket_addrs().unwrap().next().unwrap()).unwrap());
}
pub fn add_route(&mut self, path: &str, mlist: &[Method],
handler: fn(&Request) -> Response) -> &mut Canteen {
let mut methods: HashSet<Method> = HashSet::new();
for m in mlist {
methods.insert(*m);
}
for m in methods {
let rd = route::RouteDef {
pathdef: String::from(path),
method: m,
};
if self.routes.contains_key(&rd) {
panic!("a route handler for {} has already been defined!", path);
}
self.routes.insert(rd, route::Route::new(&path, m, handler));
}
self
}
pub fn set_default(&mut self, handler: fn(&Request) -> Response) -> &mut Canteen {
self.default = handler;
self
}
fn get_client(&mut self, token: Token) -> &mut Client {
self.conns.get_mut(token).unwrap()
}
fn accept(&mut self) -> Result<TcpStream> {
if let Some(ref server) = self.server {
if let Ok(s) = server.accept() {
if let Some((sock, _)) = s {
return Ok(sock);
}
}
}
Err(std::io::Error::new(
std::io::ErrorKind::ConnectionAborted,
"connection aborted prematurely".to_string()
))
}
fn handle_request(&mut self, token: Token, tx: Sender<(Token, Vec<u8>)>, rqstr: &str) {
let mut req = Request::from_str(&rqstr).unwrap();
let mut handler: fn(&Request) -> Response = self.default;
let resolved = route::RouteDef {
pathdef: req.path.clone(),
method: req.method,
};
if self.rcache.contains_key(&resolved) {
let route = &self.routes[&self.rcache[&resolved]];
handler = route.handler;
req.params = route.parse(&req.path);
} else {
for (path, route) in &self.routes {
if route.is_match(&req) {
handler = route.handler;
req.params = route.parse(&req.path);
self.rcache.insert(resolved, (*path).clone());
break;
}
}
}
self.tpool.execute(move || {
let _ = tx.send((token, handler(&req).gen_output()));
});
}
fn readable(&mut self, evl: &mut EventLoop<Canteen>, token: Token) -> Result<bool> {
if let Ok(true) = self.get_client(token).receive() {
let buf = self.get_client(token).i_buf.clone();
if let Ok(rqstr) = String::from_utf8(buf) {
self.handle_request(token, evl.channel(), &rqstr);
} else {
return Ok(false);
}
}
Ok(true)
}
fn reset_connection(&mut self, token: Token) {
self.conns.remove(token);
}
fn register(&mut self, evl: &mut EventLoop<Canteen>) -> Result<()> {
if let Some(ref server) = self.server {
return evl.register(server, self.token, EventSet::readable(), PollOpt::edge() | PollOpt::oneshot());
}
Ok(())
}
fn reregister(&mut self, evl: &mut EventLoop<Canteen>) {
if let Some(ref server) = self.server {
evl.reregister(server, self.token,
EventSet::readable(),
PollOpt::edge() | PollOpt::oneshot()).ok();
}
}
pub fn run(&mut self) {
let mut evl = match EventLoop::new() {
Ok(event_loop) => event_loop,
Err(_) => panic!("unable to initiate event loop"),
};
match self.server {
None => println!("server not bound to an address!"),
Some(_) => {
self.register(&mut evl).ok();
evl.run(self).unwrap();
},
};
}
}
impl Default for Canteen {
fn default() -> Self {
Canteen::new()
}
}