Abstracted client from WebSocket
This commit is contained in:
@@ -0,0 +1,44 @@
|
|||||||
|
use std::{
|
||||||
|
hash::{Hash, Hasher},
|
||||||
|
net::TcpStream,
|
||||||
|
sync::{Arc, Mutex},
|
||||||
|
};
|
||||||
|
|
||||||
|
use tungstenite::{Message, WebSocket, accept};
|
||||||
|
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct Client(Arc<Mutex<WebSocket<TcpStream>>>);
|
||||||
|
|
||||||
|
impl Client {
|
||||||
|
pub fn new_ws(ws: WebSocket<TcpStream>) -> Self {
|
||||||
|
Self(Arc::new(Mutex::new(ws)))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn new_tcp(ws: TcpStream) -> crate::Result<Self> {
|
||||||
|
Ok(Self(Arc::new(Mutex::new(accept(ws)?))))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PartialEq for Client {
|
||||||
|
fn eq(&self, other: &Self) -> bool {
|
||||||
|
Arc::ptr_eq(&self.0, &other.0)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Eq for Client {}
|
||||||
|
|
||||||
|
impl Hash for Client {
|
||||||
|
fn hash<H: Hasher>(&self, state: &mut H) {
|
||||||
|
std::ptr::hash(Arc::as_ptr(&self.0), state)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Client {
|
||||||
|
pub fn read(&self) -> crate::Result<Message> {
|
||||||
|
self.0.lock().unwrap().read().map_err(|e| e.into())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn send(&self, m: Message) -> crate::Result<()> {
|
||||||
|
self.0.lock().unwrap().send(m).map_err(|e| e.into())
|
||||||
|
}
|
||||||
|
}
|
||||||
+22
-21
@@ -1,9 +1,11 @@
|
|||||||
use std::{
|
use std::{
|
||||||
|
collections::HashSet,
|
||||||
net::{TcpListener, TcpStream},
|
net::{TcpListener, TcpStream},
|
||||||
path::{Path, PathBuf},
|
path::{Path, PathBuf},
|
||||||
sync::{Arc, Mutex},
|
sync::{Arc, Mutex},
|
||||||
};
|
};
|
||||||
|
|
||||||
|
pub mod client;
|
||||||
#[cfg(feature = "loader")]
|
#[cfg(feature = "loader")]
|
||||||
pub mod loader;
|
pub mod loader;
|
||||||
pub mod logger;
|
pub mod logger;
|
||||||
@@ -13,9 +15,8 @@ pub mod vfs;
|
|||||||
|
|
||||||
pub use anyhow::Result;
|
pub use anyhow::Result;
|
||||||
pub use tungstenite;
|
pub use tungstenite;
|
||||||
use tungstenite::{WebSocket, accept};
|
|
||||||
|
|
||||||
use crate::plugin::DynPlugin;
|
use crate::{client::Client, plugin::DynPlugin};
|
||||||
pub use once_cell;
|
pub use once_cell;
|
||||||
|
|
||||||
#[derive(serde::Serialize, serde::Deserialize)]
|
#[derive(serde::Serialize, serde::Deserialize)]
|
||||||
@@ -23,11 +24,12 @@ pub struct ServerConfig {
|
|||||||
port: u16,
|
port: u16,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[allow(dead_code)]
|
||||||
pub struct Server {
|
pub struct Server {
|
||||||
root: PathBuf,
|
root: PathBuf,
|
||||||
config: ServerConfig,
|
config: ServerConfig,
|
||||||
plugins: Mutex<Vec<DynPlugin>>,
|
plugins: Mutex<Vec<DynPlugin>>,
|
||||||
clients: Mutex<Vec<Arc<Mutex<WebSocket<TcpStream>>>>>,
|
clients: Mutex<HashSet<Client>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for ServerConfig {
|
impl Default for ServerConfig {
|
||||||
@@ -50,7 +52,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()),
|
clients: Mutex::new(HashSet::new()),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -59,7 +61,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()),
|
clients: Mutex::new(HashSet::new()),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -106,27 +108,29 @@ 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 ws = Arc::new(Mutex::new(accept(stream)?));
|
// Initialize client
|
||||||
|
let client = Client::new_tcp(stream)?;
|
||||||
|
|
||||||
self.clients.lock().unwrap().push(ws.clone());
|
// Insert to the set of all connected clients
|
||||||
|
self.clients.lock().unwrap().insert(client.clone());
|
||||||
|
|
||||||
|
// The main req/res loop
|
||||||
loop {
|
loop {
|
||||||
let req = ws.lock().unwrap().read()?;
|
let req = client.read()?;
|
||||||
|
|
||||||
for plugin in self.plugins.lock().unwrap().iter_mut() {
|
for plugin in self.plugins.lock().unwrap().iter_mut() {
|
||||||
plugin.on_request(&req, self);
|
plugin.on_request(&req, self);
|
||||||
}
|
}
|
||||||
|
|
||||||
if req.is_close() {
|
if req.is_close() {
|
||||||
let mut clients = self.clients.lock().unwrap();
|
self.clients.lock().unwrap().remove(&client);
|
||||||
|
break;
|
||||||
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() {
|
for c in self.clients.lock().unwrap().iter() {
|
||||||
self.wrap_err(&ws, c.lock().unwrap().send(req.clone()))?;
|
if c != &client {
|
||||||
|
self.wrap_err(&client, c.send(req.clone()))?;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -134,16 +138,13 @@ impl Server {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// When there is a error it removes the client
|
/// When there is a error it removes the client
|
||||||
fn wrap_err<T, E>(
|
pub fn wrap_err<T, E>(
|
||||||
self: &Arc<Self>,
|
self: &Arc<Self>,
|
||||||
ws: &Arc<Mutex<WebSocket<TcpStream>>>,
|
client: &Client,
|
||||||
res: std::result::Result<T, E>,
|
res: std::result::Result<T, E>,
|
||||||
) -> std::result::Result<T, E> {
|
) -> std::result::Result<T, E> {
|
||||||
if res.is_err() {
|
if res.is_err() {
|
||||||
let mut clients = self.clients.lock().unwrap();
|
self.clients.lock().unwrap().remove(&client);
|
||||||
if let Some(i) = clients.iter().position(|v| Arc::ptr_eq(v, &ws)) {
|
|
||||||
clients.remove(i);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
res
|
res
|
||||||
|
|||||||
Reference in New Issue
Block a user