Simple echo server
This commit is contained in:
+45
-5
@@ -12,7 +12,8 @@ pub mod plugin;
|
|||||||
pub mod vfs;
|
pub mod vfs;
|
||||||
|
|
||||||
pub use anyhow::Result;
|
pub use anyhow::Result;
|
||||||
use tungstenite::accept;
|
pub use tungstenite;
|
||||||
|
use tungstenite::{WebSocket, accept};
|
||||||
|
|
||||||
use crate::plugin::DynPlugin;
|
use crate::plugin::DynPlugin;
|
||||||
pub use once_cell;
|
pub use once_cell;
|
||||||
@@ -23,9 +24,10 @@ pub struct ServerConfig {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub struct Server {
|
pub struct Server {
|
||||||
plugins: Mutex<Vec<DynPlugin>>,
|
|
||||||
root: PathBuf,
|
root: PathBuf,
|
||||||
config: ServerConfig,
|
config: ServerConfig,
|
||||||
|
plugins: Mutex<Vec<DynPlugin>>,
|
||||||
|
clients: Mutex<Vec<Arc<Mutex<WebSocket<TcpStream>>>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for ServerConfig {
|
impl Default for ServerConfig {
|
||||||
@@ -48,6 +50,7 @@ impl Server {
|
|||||||
plugins: Mutex::new(Vec::new()),
|
plugins: Mutex::new(Vec::new()),
|
||||||
root: root.to_path_buf(),
|
root: root.to_path_buf(),
|
||||||
config: ServerConfig::default(),
|
config: ServerConfig::default(),
|
||||||
|
clients: Mutex::new(Vec::new()),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -56,6 +59,7 @@ impl Server {
|
|||||||
plugins: Mutex::new(Vec::new()),
|
plugins: Mutex::new(Vec::new()),
|
||||||
root: root.to_path_buf(),
|
root: root.to_path_buf(),
|
||||||
config,
|
config,
|
||||||
|
clients: Mutex::new(Vec::new()),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -102,10 +106,46 @@ impl Server {
|
|||||||
|
|
||||||
fn handle_client(self: &Arc<Self>, stream: TcpStream) -> anyhow::Result<()> {
|
fn handle_client(self: &Arc<Self>, stream: TcpStream) -> anyhow::Result<()> {
|
||||||
Self::LOGGER.info(format!("New connection: {}", stream.peer_addr()?));
|
Self::LOGGER.info(format!("New connection: {}", stream.peer_addr()?));
|
||||||
let mut ws = accept(stream)?;
|
let ws = Arc::new(Mutex::new(accept(stream)?));
|
||||||
|
|
||||||
|
self.clients.lock().unwrap().push(ws.clone());
|
||||||
|
|
||||||
|
loop {
|
||||||
|
let req = ws.lock().unwrap().read()?;
|
||||||
|
|
||||||
|
for plugin in self.plugins.lock().unwrap().iter_mut() {
|
||||||
|
plugin.on_request(&req, self);
|
||||||
|
}
|
||||||
|
|
||||||
|
if req.is_close() {
|
||||||
|
let mut clients = self.clients.lock().unwrap();
|
||||||
|
|
||||||
|
if let Some(i) = clients.iter().position(|v| Arc::ptr_eq(v, &ws)) {
|
||||||
|
clients.remove(i);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for c in self.clients.lock().unwrap().iter_mut() {
|
||||||
|
self.wrap_err(&ws, c.lock().unwrap().send(req.clone()))?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
let msg = ws.read()?;
|
|
||||||
ws.send(msg)?;
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// When there is a error it removes the client
|
||||||
|
fn wrap_err<T, E>(
|
||||||
|
self: &Arc<Self>,
|
||||||
|
ws: &Arc<Mutex<WebSocket<TcpStream>>>,
|
||||||
|
res: std::result::Result<T, E>,
|
||||||
|
) -> std::result::Result<T, E> {
|
||||||
|
if res.is_err() {
|
||||||
|
let mut clients = self.clients.lock().unwrap();
|
||||||
|
if let Some(i) = clients.iter().position(|v| Arc::ptr_eq(v, &ws)) {
|
||||||
|
clients.remove(i);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
res
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,9 +1,15 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
use tungstenite::Message;
|
||||||
|
|
||||||
use crate::Server;
|
use crate::Server;
|
||||||
|
|
||||||
pub type DynPlugin = Box<dyn Plugin + Send + Sync>;
|
pub type DynPlugin = Box<dyn Plugin + Send + Sync>;
|
||||||
|
|
||||||
pub trait Plugin {
|
pub trait Plugin {
|
||||||
fn init(&mut self, server: &Arc<Server>);
|
fn init(&mut self, server: &Arc<Server>);
|
||||||
|
#[allow(unused_variables)]
|
||||||
|
fn on_request(&mut self, msg: &Message, server: &Arc<Server>) -> bool {
|
||||||
|
false
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user