use crate::server::handler::Handler;
use crate::bucket::bucket_list::BucketList;
use crate::bootstrap::file::{save, load};
use crate::server::router_utils::open_port;
use crate::server::threadpool::ThreadPool;
use crate::error::FileError;
use crate::node::Node;
use std::net::{TcpListener, TcpStream, Ipv4Addr};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use std::thread;
const MAX_BUCKETS: usize = 10;
const BUCKET_SIZE: usize = 10;
pub struct Server {
thread_pool: Arc<Mutex<ThreadPool>>, bucket_list: Arc<Mutex<BucketList>>, port: Option<u16>, }
impl Server {
pub fn new(local_ip: Ipv4Addr) -> Self {
let mut rng = rand::thread_rng();
let listener;
let local_port;
loop {
let p = 1024;
if let Ok(l) = TcpListener::bind(format!("{}:{}", local_ip, p)) {
listener = l;
local_port = p;
break;
}
}
let port = match open_port(local_ip, local_port) {
Ok(p) => Some(p),
Err(e) => {
println!("Error opening port: {}", e);
None
}
};
println!("Bound to {}:{}", local_ip, local_port);
let bucket_list = Arc::new(Mutex::new(BucketList::new(MAX_BUCKETS, BUCKET_SIZE)));
let bucket_list_thread = Arc::clone(&bucket_list);
let thread_pool = Arc::new(Mutex::new(ThreadPool::new(4)));
let thread_pool_thread = Arc::clone(&thread_pool);
println!("Created bucket list and threadpool");
thread::spawn(move || {
for s in listener.incoming() {
if let Ok(stream) = s {
if stream.peer_addr().unwrap().ip() == local_ip {
break;
}
let bl = bucket_list_thread.clone();
thread_pool_thread.lock().unwrap().execute(move || {
if let Err(e) = stream.set_read_timeout(Some(Duration::from_secs(10))) {
println!("Failed to set read timeout for tcpstream: {:?}", e);
}
let mut handler = Handler::new(stream, bl);
handler.start();
});
}
};
});
println!("Server launched ok");
Server {
thread_pool,
bucket_list,
port,
}
}
pub fn add_node(&mut self, node: Node) {
let mut self_bucket_list = self.bucket_list.lock().unwrap();
self_bucket_list.add_node(&node).unwrap();
}
pub fn queue_stream(&self, stream: TcpStream) {
let bl = Arc::clone(&self.bucket_list);
self.thread_pool.lock().unwrap().execute(move || {
println!("Accepted connection with {}", stream.peer_addr().unwrap());
if let Err(e) = stream.set_read_timeout(Some(Duration::from_secs(10))) {
println!("Failed to set read timeout for tcpstream: {:?}", e);
}
let mut handler = Handler::new(stream, bl);
handler.start();
});
}
pub fn port(&self) -> Option<u16> {
self.port
}
pub fn save(&self, path: &str) {
let node_list = self.bucket_list.lock().unwrap().node_list();
let _ = save(path, &node_list);
}
pub fn load(&self, path: &str) -> Result<(), FileError> {
let file_bucket_list = load(path)?;
let mut self_bucket_list = self.bucket_list.lock().unwrap();
for node in file_bucket_list.iter() {
self_bucket_list.add_node(node).unwrap();
}
Ok(())
}
}