From 6a9c27a76dc78896e6b81f9db56969433461bc4d Mon Sep 17 00:00:00 2001 From: Meisterlala <6453306+Meisterlala@users.noreply.github.com> Date: Tue, 26 Dec 2023 20:43:02 +0100 Subject: [PATCH] Simple Loopback --- .vscode/launch.json | 7 +--- client/src/main.rs | 2 +- server/Cargo.toml | 8 ++++ server/src/connection.rs | 82 ++++++++++++++++++++++++++++++++++++++++ server/src/lib.rs | 2 + server/src/main.rs | 29 +++++++------- server/src/websocket.rs | 55 +++++++++++++++++++++++++++ 7 files changed, 163 insertions(+), 22 deletions(-) create mode 100644 server/src/connection.rs create mode 100644 server/src/lib.rs create mode 100644 server/src/websocket.rs diff --git a/.vscode/launch.json b/.vscode/launch.json index b9f1c77..4f83626 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -16,13 +16,8 @@ "kind": "bin" } }, - "stdio": [ - null, - null, - "${workspaceFolder}/server.log", - ], "env": { - "RUST_LOG": "chat_server=debug", + "RUST_LOG": "debug", "RUST_BACKTRACE": "1", }, "cwd": "${workspaceFolder}", diff --git a/client/src/main.rs b/client/src/main.rs index ad97b65..e91909c 100644 --- a/client/src/main.rs +++ b/client/src/main.rs @@ -8,7 +8,7 @@ async fn main() { // Setup env_logger to write to stderr env_logger::init(); - let app = Application::new("wss://echo.websocket.events"); + let app = Application::new("ws://127.0.0.1:9001"); app.run().await.unwrap(); info!("Exiting"); diff --git a/server/Cargo.toml b/server/Cargo.toml index 33774a2..7fe2ea6 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -9,3 +9,11 @@ tokio = { version = "1.35.1", features = ["full"] } # Async runtime tokio-tungstenite = { version = "0.21.0", features = [ "native-tls", ] } # Async WebSocket + + +# Error handling +anyhow = "1.0.76" + +# Logging +log = "0.4.20" +env_logger = "0.10.1" diff --git a/server/src/connection.rs b/server/src/connection.rs new file mode 100644 index 0000000..d2d5ac9 --- /dev/null +++ b/server/src/connection.rs @@ -0,0 +1,82 @@ +use futures_util::{SinkExt, StreamExt}; +use log::{debug, error}; +use tokio::net::TcpStream; + +pub struct Connection { + pub read: tokio::sync::mpsc::UnboundedReceiver, + pub write: tokio::sync::mpsc::UnboundedSender, + pub task: tokio::task::JoinHandle<()>, +} + +impl Connection { + pub fn new(stream: tokio_tungstenite::WebSocketStream) -> Self { + let (tx_read, rx_read) = tokio::sync::mpsc::unbounded_channel(); + let (tx_write, mut rx_write) = tokio::sync::mpsc::unbounded_channel(); + + let t = tokio::spawn(async move { + let (mut ws_write, mut ws_read) = stream.split(); + + loop { + tokio::select! { + Some(msg) = rx_write.recv() => { + debug!("Sending message: {}", msg); + ws_write.send(tokio_tungstenite::tungstenite::Message::Text(msg)).await.unwrap(); + } + Some(msg) = ws_read.next() => { + match msg { + Ok(msg) => { + let msg = msg.into_text().expect("Failed to convert message to text"); + debug!("Recieved message: {}", msg); + tx_read.send(msg).unwrap(); + } + Err(e) => { + error!("Error reading from websocket: {}", e); + break; + } + } + } + } + } + }); + + Self { + read: rx_read, + write: tx_write, + task: t, + } + } + + pub fn loopback(stream: tokio_tungstenite::WebSocketStream) -> Self { + let task = tokio::spawn(async move { + let (mut ws_write, mut ws_read) = stream.split(); + while let Some(Ok(msg)) = ws_read.next().await { + if !msg.is_close() && !msg.is_pong() && !msg.is_empty() { + debug!("Resending: {}", msg); + ws_write.send(msg).await.unwrap(); + } + } + debug!("Loopback closed") + }); + + let (w, r) = tokio::sync::mpsc::unbounded_channel(); + Self { + read: r, + write: w, + task, + } + } + + pub fn send(&mut self, msg: String) -> anyhow::Result<()> { + self.write.send(msg)?; + Ok(()) + } + + pub async fn recieve(&mut self) -> anyhow::Result { + let msg = self + .read + .recv() + .await + .ok_or_else(|| anyhow::anyhow!("Failed to recieve message"))?; + Ok(msg) + } +} diff --git a/server/src/lib.rs b/server/src/lib.rs new file mode 100644 index 0000000..68b217d --- /dev/null +++ b/server/src/lib.rs @@ -0,0 +1,2 @@ +pub mod websocket; +pub mod connection; \ No newline at end of file diff --git a/server/src/main.rs b/server/src/main.rs index 27c398d..9a5a419 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -1,20 +1,19 @@ use std::net::TcpListener; -/// A WebSocket echo server -fn main() { - let server = TcpListener::bind("127.0.0.1:9001").unwrap(); - /* for stream in server.incoming() { - spawn (move || { - let mut websocket = accept(stream.unwrap()).unwrap(); - loop { - let msg = websocket.read().unwrap(); +use chat_server::websocket::Websocket; +use futures_util::FutureExt; +use log::info; +#[tokio::main] +async fn main() { + env_logger::init(); - // We do not want to send back ping/pong messages. - if msg.is_binary() || msg.is_text() { - websocket.send(msg).unwrap(); - } - } - }); - } */ + let ws = Websocket::new("127.0.0.1:9001"); + + info!("Strted"); + tokio::select! { + _ = tokio::signal::ctrl_c() => {info!("Ctrl-C recieved")}, + r = ws.serve().fuse() => {info!("Websocket task ended with: {:?}", r)}, + } + info!("Exiting"); } diff --git a/server/src/websocket.rs b/server/src/websocket.rs new file mode 100644 index 0000000..e46866b --- /dev/null +++ b/server/src/websocket.rs @@ -0,0 +1,55 @@ +use std::sync::{Arc, Mutex}; + +use log::*; +use tokio::net::TcpListener; + +use crate::connection::Connection; +pub struct Websocket { + adress: String, + connections: Arc>>, +} + +impl Websocket { + pub fn new(adress: &str) -> Self { + Self { + adress: String::from(adress), + connections: Arc::new(Mutex::new(Vec::new())), + } + } + + pub async fn serve(&self) { + let listener = TcpListener::bind(&self.adress) + .await + .unwrap_or_else(|_| panic!("Failed to bind to adress: {}", &self.adress)); + + info!("Listening on: {}", &self.adress); + + let c = self.connections.clone(); + debug!("TCP listener started"); + while let Ok((stream, _)) = listener.accept().await { + info!( + "New TCP connection: {:?}", + stream + .peer_addr() + .map_or_else(|_| "Unknown".to_owned(), |a| a.to_string()) + ); + + let c2 = c.clone(); + tokio::spawn(async move { + match tokio_tungstenite::accept_async(stream).await { + Ok(ws_stream) => { + info!( + "New WebSocket connection: {:?}", + ws_stream.get_ref().peer_addr() + ); + c2.lock().unwrap().push(Connection::loopback(ws_stream)); + } + Err(e) => { + error!("Error during the websocket handshake occurred: {}", e); + } + } + }); + } + debug!("TCP listener ended"); + } +} -- 2.51.2