From 0919ec30d3c9eeff3db23a5c7c86688068d8fe9f Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sat, 13 Sep 2025 10:35:54 +0200 Subject: [PATCH 01/11] Base protocol --- client.js | 18 ++++++++++++++++++ src/client.rs | 31 ++++++++++++++++++++++++++----- src/lib.rs | 40 +++++++++++++++++++++++++++------------- src/plugin.rs | 9 +++++---- src/types.rs | 40 ++++++++++++++++++++++++++++++++++++++++ 5 files changed, 116 insertions(+), 22 deletions(-) create mode 100644 client.js create mode 100644 src/types.rs diff --git a/client.js b/client.js new file mode 100644 index 0000000..571abf9 --- /dev/null +++ b/client.js @@ -0,0 +1,18 @@ +setInterval(() => { + const ws = new WebSocket('ws://localhost:7080'); + + function sendMessage(message) { + ws.send(JSON.stringify({ type: 'send_message', params: message })); + } + + ws.onopen = () => { + console.log('WebSocket connection established'); + sendMessage('Hello, Server!'); + }; + + ws.onmessage = (event) => { + const message = JSON.parse(event.data); + console.log('Received:', message); + ws.close(); + }; +}, 100) \ No newline at end of file diff --git a/src/client.rs b/src/client.rs index 73a8aeb..fe8cf3f 100644 --- a/src/client.rs +++ b/src/client.rs @@ -4,7 +4,10 @@ use std::{ sync::{Arc, Mutex}, }; -use tungstenite::{Message, WebSocket, accept}; +use anyhow::Error; +use tungstenite::{Message, Utf8Bytes, WebSocket, accept}; + +use crate::types::{FromClient, ToClient, WsMessage}; #[derive(Clone)] pub struct Client(Arc>>); @@ -34,11 +37,29 @@ impl Hash for Client { } impl Client { - pub fn read(&self) -> crate::Result { - self.0.lock().unwrap().read().map_err(|e| e.into()) + pub fn read(&self) -> crate::Result>> { + match self.0.lock().unwrap().read()? { + Message::Text(t) => { + let v = t.to_string(); + match serde_json::from_str(&v) { + Ok(f) => Ok(Some(WsMessage::FromClient(f))), + Err(_) => Ok(Some(WsMessage::String(v))), + } + } + + Message::Binary(b) => Ok(Some(WsMessage::Binary(b))), + + Message::Close(_) => Ok(None), + + m => Err(Error::msg(format!("Invalid websocket format: {m}"))), + } } - pub fn send(&self, m: Message) -> crate::Result<()> { - self.0.lock().unwrap().send(m).map_err(|e| e.into()) + pub fn send(&self, m: ToClient) -> crate::Result<()> { + self.0 + .lock() + .unwrap() + .send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?))) + .map_err(|e| e.into()) } } diff --git a/src/lib.rs b/src/lib.rs index e411d95..a75419c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -11,12 +11,17 @@ pub mod loader; pub mod logger; pub mod macros; pub mod plugin; +pub mod types; pub mod vfs; pub use anyhow::Result; pub use tungstenite; -use crate::{client::Client, plugin::DynPlugin}; +use crate::{ + client::Client, + plugin::DynPlugin, + types::{ToClient, data}, +}; pub use once_cell; #[derive(serde::Serialize, serde::Deserialize)] @@ -116,20 +121,29 @@ impl Server { // The main req/res loop loop { - let req = client.read()?; + match client.read()? { + Some(req) => { + println!("Request: {:?}", req); + for plugin in self.plugins.lock().unwrap().iter_mut() { + plugin.on_request(&req, self); + } - for plugin in self.plugins.lock().unwrap().iter_mut() { - plugin.on_request(&req, self); - } + for c in self.clients.lock().unwrap().iter() { + if c == &client { + self.wrap_err( + &client, + c.send(ToClient::Message(data::Message { + from: format!("Server"), + contents: format!("Hello"), + })), + )?; + } + } + } - if req.is_close() { - self.clients.lock().unwrap().remove(&client); - break; - } - - for c in self.clients.lock().unwrap().iter() { - if c != &client { - self.wrap_err(&client, c.send(req.clone()))?; + None => { + self.clients.lock().unwrap().remove(&client); + break; } } } diff --git a/src/plugin.rs b/src/plugin.rs index 3766b22..ba72677 100644 --- a/src/plugin.rs +++ b/src/plugin.rs @@ -1,15 +1,16 @@ use std::sync::Arc; -use tungstenite::Message; - -use crate::Server; +use crate::{ + Server, + types::{FromClient, WsMessage}, +}; pub type DynPlugin = Box; pub trait Plugin { fn init(&mut self, server: &Arc); #[allow(unused_variables)] - fn on_request(&mut self, msg: &Message, server: &Arc) -> bool { + fn on_request(&mut self, msg: &WsMessage, server: &Arc) -> bool { false } } diff --git a/src/types.rs b/src/types.rs new file mode 100644 index 0000000..ee0e128 --- /dev/null +++ b/src/types.rs @@ -0,0 +1,40 @@ +use serde::{Deserialize, Serialize}; +use tungstenite::Bytes; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "type", content = "params", rename_all = "snake_case")] +pub enum FromClient { + SendMessage(String), +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "type", content = "params", rename_all = "snake_case")] +pub enum FromServer {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "type", content = "params", rename_all = "snake_case")] +pub enum ToClient { + Message(data::Message), +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "type", content = "params", rename_all = "snake_case")] +pub enum ToServer {} + +/// WebSocket for client messages +#[derive(Debug, Clone)] +pub enum WsMessage Deserialize<'de>> { + FromClient(T), + Binary(Bytes), + String(String), +} + +pub mod data { + use serde::{Deserialize, Serialize}; + + #[derive(Debug, Clone, Serialize, Deserialize)] + pub struct Message { + pub from: String, + pub contents: String, + } +} From 642fe6cfcbc69a38c08679f4d96e614f103765ff Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sat, 13 Sep 2025 11:45:44 +0200 Subject: [PATCH 02/11] Updated protocol --- Cargo.toml | 1 + cli/Cargo.lock | 232 ++++++++++++++++++++++++++++++++++++++++++++++++- client.js | 5 +- src/client.rs | 6 +- src/lib.rs | 45 ++++++---- src/plugin.rs | 4 +- src/types.rs | 75 +++++++++++++--- 7 files changed, 334 insertions(+), 34 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 427fd0c..9dfe812 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,6 +5,7 @@ edition = "2024" [dependencies] anyhow = "1.0.99" +chrono = "0.4.42" libloading = { version = "0.8.8", optional = true } once_cell = "1.21.3" serde = { version = "1.0.219", features = ["serde_derive"] } diff --git a/cli/Cargo.lock b/cli/Cargo.lock index 1a42a52..64d6923 100644 --- a/cli/Cargo.lock +++ b/cli/Cargo.lock @@ -2,12 +2,27 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "android_system_properties" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "819e7219dbd41043ac279b19830f2efc897156490d7fd6ea916720117ee66311" +dependencies = [ + "libc", +] + [[package]] name = "anyhow" version = "1.0.99" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b0674a1ddeecb70197781e945de4b3b8ffb61fa939a5597bcf48503737663100" +[[package]] +name = "autocfg" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" + [[package]] name = "block-buffer" version = "0.10.4" @@ -17,18 +32,47 @@ dependencies = [ "generic-array", ] +[[package]] +name = "bumpalo" +version = "3.19.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "46c5e41b57b8bba42a04676d81cb89e9ee8e859a1a66f80a5a72e1cb76b34d43" + [[package]] name = "bytes" version = "1.10.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d71b6127be86fdcfddb610f7182ac57211d4b18a3e9c82eb2d17662f2227ad6a" +[[package]] +name = "cc" +version = "1.2.37" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "65193589c6404eb80b450d618eaf9a2cafaaafd57ecce47370519ef674a7bd44" +dependencies = [ + "find-msvc-tools", + "shlex", +] + [[package]] name = "cfg-if" version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2fd1289c04a9ea8cb22300a459a72a385d7c73d3259e2ed7dcb2af674838cfa9" +[[package]] +name = "chrono" +version = "0.4.42" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "145052bdd345b87320e369255277e3fb5152762ad123a901ef5c262dd38fe8d2" +dependencies = [ + "iana-time-zone", + "js-sys", + "num-traits", + "wasm-bindgen", + "windows-link 0.2.0", +] + [[package]] name = "cli" version = "0.1.0" @@ -36,6 +80,12 @@ dependencies = [ "voxa-server", ] +[[package]] +name = "core-foundation-sys" +version = "0.8.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" + [[package]] name = "cpufeatures" version = "0.2.17" @@ -71,6 +121,12 @@ dependencies = [ "crypto-common", ] +[[package]] +name = "find-msvc-tools" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7fd99930f64d146689264c637b5af2f0233a933bef0d8570e2526bf9e083192d" + [[package]] name = "fnv" version = "1.0.7" @@ -116,12 +172,46 @@ version = "1.10.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" +[[package]] +name = "iana-time-zone" +version = "0.1.64" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "33e57f83510bb73707521ebaffa789ec8caf86f9657cad665b092b581d40e9fb" +dependencies = [ + "android_system_properties", + "core-foundation-sys", + "iana-time-zone-haiku", + "js-sys", + "log", + "wasm-bindgen", + "windows-core", +] + +[[package]] +name = "iana-time-zone-haiku" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f31827a206f56af32e590ba56d5d2d085f558508192593743f16b2306495269f" +dependencies = [ + "cc", +] + [[package]] name = "itoa" version = "1.0.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4a5f13b858c8d314ee3e8f639011f7ccefe71f97f96e50151fb991f267928e2c" +[[package]] +name = "js-sys" +version = "0.3.78" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c0b063578492ceec17683ef2f8c5e89121fbd0b172cbc280635ab7567db2738" +dependencies = [ + "once_cell", + "wasm-bindgen", +] + [[package]] name = "libc" version = "0.2.175" @@ -150,6 +240,15 @@ version = "2.7.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a282da65faaf38286cf3be983213fcf1d2e2a58700e808f83f4ea9a4804bc0" +[[package]] +name = "num-traits" +version = "0.2.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "071dfc062690e90b734c0b2273ce72ad0ffa95f0c74596bc250dcfd960262841" +dependencies = [ + "autocfg", +] + [[package]] name = "once_cell" version = "1.21.3" @@ -218,6 +317,12 @@ dependencies = [ "getrandom", ] +[[package]] +name = "rustversion" +version = "1.0.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" + [[package]] name = "ryu" version = "1.0.20" @@ -267,6 +372,12 @@ dependencies = [ "digest", ] +[[package]] +name = "shlex" +version = "1.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0fda2ff0d084019ba4d7c6f371c95d8fd75ce3524c3cb8fb653a3023f6323e64" + [[package]] name = "syn" version = "2.0.106" @@ -344,6 +455,7 @@ name = "voxa-server" version = "0.1.0" dependencies = [ "anyhow", + "chrono", "libloading", "once_cell", "serde", @@ -360,19 +472,137 @@ dependencies = [ "wit-bindgen", ] +[[package]] +name = "wasm-bindgen" +version = "0.2.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7e14915cadd45b529bb8d1f343c4ed0ac1de926144b746e2710f9cd05df6603b" +dependencies = [ + "cfg-if", + "once_cell", + "rustversion", + "wasm-bindgen-macro", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-backend" +version = "0.2.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e28d1ba982ca7923fd01448d5c30c6864d0a14109560296a162f80f305fb93bb" +dependencies = [ + "bumpalo", + "log", + "proc-macro2", + "quote", + "syn", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-macro" +version = "0.2.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7c3d463ae3eff775b0c45df9da45d68837702ac35af998361e2c84e7c5ec1b0d" +dependencies = [ + "quote", + "wasm-bindgen-macro-support", +] + +[[package]] +name = "wasm-bindgen-macro-support" +version = "0.2.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7bb4ce89b08211f923caf51d527662b75bdc9c9c7aab40f86dcb9fb85ac552aa" +dependencies = [ + "proc-macro2", + "quote", + "syn", + "wasm-bindgen-backend", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-shared" +version = "0.2.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f143854a3b13752c6950862c906306adb27c7e839f7414cec8fea35beab624c1" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "windows-core" +version = "0.62.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "57fe7168f7de578d2d8a05b07fd61870d2e73b4020e9f49aa00da8471723497c" +dependencies = [ + "windows-implement", + "windows-interface", + "windows-link 0.2.0", + "windows-result", + "windows-strings", +] + +[[package]] +name = "windows-implement" +version = "0.60.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a47fddd13af08290e67f4acabf4b459f647552718f683a7b415d290ac744a836" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "windows-interface" +version = "0.59.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bd9211b69f8dcdfa817bfd14bf1c97c9188afa36f4750130fcdf3f400eca9fa8" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "windows-link" version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5e6ad25900d524eaabdbbb96d20b4311e1e7ae1699af4fb28c17ae66c80d798a" +[[package]] +name = "windows-link" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "45e46c0661abb7180e7b9c281db115305d49ca1709ab8242adf09666d2173c65" + +[[package]] +name = "windows-result" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7084dcc306f89883455a206237404d3eaf961e5bd7e0f312f7c91f57eb44167f" +dependencies = [ + "windows-link 0.2.0", +] + +[[package]] +name = "windows-strings" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7218c655a553b0bed4426cf54b20d7ba363ef543b52d515b3e48d7fd55318dda" +dependencies = [ + "windows-link 0.2.0", +] + [[package]] name = "windows-targets" version = "0.53.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d5fe6031c4041849d7c496a8ded650796e7b6ecc19df1a431c1a363342e5dc91" dependencies = [ - "windows-link", + "windows-link 0.1.3", "windows_aarch64_gnullvm", "windows_aarch64_msvc", "windows_i686_gnu", diff --git a/client.js b/client.js index 571abf9..44c6df7 100644 --- a/client.js +++ b/client.js @@ -2,7 +2,10 @@ setInterval(() => { const ws = new WebSocket('ws://localhost:7080'); function sendMessage(message) { - ws.send(JSON.stringify({ type: 'send_message', params: message })); + ws.send(JSON.stringify({ type: 'send_message', params: { + channel_id: 'Hello', + contents: message + }})); } ws.onopen = () => { diff --git a/src/client.rs b/src/client.rs index fe8cf3f..9a00b8c 100644 --- a/src/client.rs +++ b/src/client.rs @@ -7,7 +7,7 @@ use std::{ use anyhow::Error; use tungstenite::{Message, Utf8Bytes, WebSocket, accept}; -use crate::types::{FromClient, ToClient, WsMessage}; +use crate::types::{ClientMessage, ServerMessage, WsMessage}; #[derive(Clone)] pub struct Client(Arc>>); @@ -37,7 +37,7 @@ impl Hash for Client { } impl Client { - pub fn read(&self) -> crate::Result>> { + pub fn read(&self) -> crate::Result>> { match self.0.lock().unwrap().read()? { Message::Text(t) => { let v = t.to_string(); @@ -55,7 +55,7 @@ impl Client { } } - pub fn send(&self, m: ToClient) -> crate::Result<()> { + pub fn send(&self, m: ServerMessage) -> crate::Result<()> { self.0 .lock() .unwrap() diff --git a/src/lib.rs b/src/lib.rs index a75419c..8063e98 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -20,7 +20,7 @@ pub use tungstenite; use crate::{ client::Client, plugin::DynPlugin, - types::{ToClient, data}, + types::{ClientMessage, WsMessage}, }; pub use once_cell; @@ -122,23 +122,38 @@ impl Server { // The main req/res loop loop { match client.read()? { - Some(req) => { - println!("Request: {:?}", req); - for plugin in self.plugins.lock().unwrap().iter_mut() { - plugin.on_request(&req, self); + Some(WsMessage::FromClient(req)) => match req { + ClientMessage::SendMessage { + channel_id, + contents, + } => { + Self::LOGGER.info(format!("SendMessage to {channel_id}: {contents}")); } - for c in self.clients.lock().unwrap().iter() { - if c == &client { - self.wrap_err( - &client, - c.send(ToClient::Message(data::Message { - from: format!("Server"), - contents: format!("Hello"), - })), - )?; - } + ClientMessage::EditMessage { + channel_id, + message_id, + new_contents, + } => { + Self::LOGGER.info(format!( + "EditMessage {message_id} in {channel_id}: {new_contents}" + )); } + + ClientMessage::DeleteMessage { + channel_id, + message_id, + } => { + Self::LOGGER.info(format!("DeleteMessage {message_id} in {channel_id}")); + } + }, + + Some(WsMessage::Binary(b)) => { + Self::LOGGER.info(format!("Binary message: {b:?}")); + } + + Some(WsMessage::String(s)) => { + Self::LOGGER.info(format!("String message: {s}")); } None => { diff --git a/src/plugin.rs b/src/plugin.rs index ba72677..55495df 100644 --- a/src/plugin.rs +++ b/src/plugin.rs @@ -2,7 +2,7 @@ use std::sync::Arc; use crate::{ Server, - types::{FromClient, WsMessage}, + types::{ClientMessage, WsMessage}, }; pub type DynPlugin = Box; @@ -10,7 +10,7 @@ pub type DynPlugin = Box; pub trait Plugin { fn init(&mut self, server: &Arc); #[allow(unused_variables)] - fn on_request(&mut self, msg: &WsMessage, server: &Arc) -> bool { + fn on_request(&mut self, msg: &WsMessage, server: &Arc) -> bool { false } } diff --git a/src/types.rs b/src/types.rs index ee0e128..616dc8a 100644 --- a/src/types.rs +++ b/src/types.rs @@ -1,27 +1,60 @@ use serde::{Deserialize, Serialize}; use tungstenite::Bytes; +/// Messages sent *from the client* (user’s app) to the server #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(tag = "type", content = "params", rename_all = "snake_case")] -pub enum FromClient { - SendMessage(String), +pub enum ClientMessage { + /// Send a message to a channel + SendMessage { + channel_id: String, + contents: String, + }, + + /// Edit a message (if allowed) + EditMessage { + channel_id: String, + message_id: String, + new_contents: String, + }, + + /// Delete a message (if allowed) + DeleteMessage { + channel_id: String, + message_id: String, + }, } +/// Messages sent *from the server* to the client #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(tag = "type", content = "params", rename_all = "snake_case")] -pub enum FromServer {} +pub enum ServerMessage { + /// Successful authentication + Authenticated { user_id: String }, -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(tag = "type", content = "params", rename_all = "snake_case")] -pub enum ToClient { - Message(data::Message), + /// Error responses + Error { message: String }, + + /// A new message in a channel + MessageCreate(data::Message), + + /// A message was edited + MessageUpdate(data::Message), + + /// A message was deleted + MessageDelete { + channel_id: String, + message_id: String, + }, + + /// Presence updates + PresenceUpdate { user_id: String, status: String }, + + /// Typing indicator + Typing { user_id: String, channel_id: String }, } -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(tag = "type", content = "params", rename_all = "snake_case")] -pub enum ToServer {} - -/// WebSocket for client messages +/// WebSocket wrapper #[derive(Debug, Clone)] pub enum WsMessage Deserialize<'de>> { FromClient(T), @@ -29,12 +62,30 @@ pub enum WsMessage Deserialize<'de>> { String(String), } +/// Shared data structures pub mod data { use serde::{Deserialize, Serialize}; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Message { + pub id: String, + pub channel_id: u8, pub from: String, pub contents: String, + pub timestamp: i64, + } + + #[derive(Debug, Clone, Serialize, Deserialize)] + pub struct Channel { + pub id: usize, + pub name: String, + pub kind: ChannelKind, + } + + #[derive(Debug, Clone, Serialize, Deserialize)] + #[serde(rename_all = "snake_case")] + pub enum ChannelKind { + Text, + Voice, } } From 3613d0bfdc7e82f09e97acba64d412cb77605f1c Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sat, 13 Sep 2025 12:37:18 +0200 Subject: [PATCH 03/11] Added database --- .gitignore | 3 +- Cargo.toml | 1 + cli/Cargo.lock | 85 +++++++++++++++++++++++++++++++++++++++++++++++++ src/database.rs | 83 +++++++++++++++++++++++++++++++++++++++++++++++ src/lib.rs | 25 +++++++++++---- src/types.rs | 2 +- 6 files changed, 190 insertions(+), 9 deletions(-) create mode 100644 src/database.rs diff --git a/.gitignore b/.gitignore index 598a1f3..b0bc9c8 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,5 @@ target/ plugins/ config.json -/Cargo.lock \ No newline at end of file +/Cargo.lock +*.db \ No newline at end of file diff --git a/Cargo.toml b/Cargo.toml index 9dfe812..bd412a9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,6 +8,7 @@ anyhow = "1.0.99" chrono = "0.4.42" libloading = { version = "0.8.8", optional = true } once_cell = "1.21.3" +rusqlite = "0.37.0" serde = { version = "1.0.219", features = ["serde_derive"] } serde_json = "1.0.143" tungstenite = "0.27.0" diff --git a/cli/Cargo.lock b/cli/Cargo.lock index 64d6923..d83c247 100644 --- a/cli/Cargo.lock +++ b/cli/Cargo.lock @@ -23,6 +23,12 @@ version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" +[[package]] +name = "bitflags" +version = "2.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2261d10cca569e4643e526d8dc2e62e433cc8aba21ab764233731f8d369bf394" + [[package]] name = "block-buffer" version = "0.10.4" @@ -121,6 +127,18 @@ dependencies = [ "crypto-common", ] +[[package]] +name = "fallible-iterator" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649" + +[[package]] +name = "fallible-streaming-iterator" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a" + [[package]] name = "find-msvc-tools" version = "0.1.1" @@ -133,6 +151,12 @@ version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" +[[package]] +name = "foldhash" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" + [[package]] name = "generic-array" version = "0.14.7" @@ -155,6 +179,24 @@ dependencies = [ "wasi", ] +[[package]] +name = "hashbrown" +version = "0.15.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" +dependencies = [ + "foldhash", +] + +[[package]] +name = "hashlink" +version = "0.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7382cf6263419f2d8df38c55d7da83da5c18aef87fc7a7fc1fb1e344edfe14c1" +dependencies = [ + "hashbrown", +] + [[package]] name = "http" version = "1.3.1" @@ -228,6 +270,16 @@ dependencies = [ "windows-targets", ] +[[package]] +name = "libsqlite3-sys" +version = "0.35.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "133c182a6a2c87864fe97778797e46c7e999672690dc9fa3ee8e241aa4a9c13f" +dependencies = [ + "pkg-config", + "vcpkg", +] + [[package]] name = "log" version = "0.4.27" @@ -255,6 +307,12 @@ version = "1.21.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "42f5e15c9953c5e4ccceeb2e7382a716482c34515315f7b03532b8b4e8393d2d" +[[package]] +name = "pkg-config" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" + [[package]] name = "ppv-lite86" version = "0.2.21" @@ -317,6 +375,20 @@ dependencies = [ "getrandom", ] +[[package]] +name = "rusqlite" +version = "0.37.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "165ca6e57b20e1351573e3729b958bc62f0e48025386970b6e4d29e7a7e71f3f" +dependencies = [ + "bitflags", + "fallible-iterator", + "fallible-streaming-iterator", + "hashlink", + "libsqlite3-sys", + "smallvec", +] + [[package]] name = "rustversion" version = "1.0.22" @@ -378,6 +450,12 @@ version = "1.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0fda2ff0d084019ba4d7c6f371c95d8fd75ce3524c3cb8fb653a3023f6323e64" +[[package]] +name = "smallvec" +version = "1.15.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "67b1b7a3b5fe4f1376887184045fcf45c69e92af734b7aaddc05fb777b6fbd03" + [[package]] name = "syn" version = "2.0.106" @@ -444,6 +522,12 @@ version = "0.7.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" +[[package]] +name = "vcpkg" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" + [[package]] name = "version_check" version = "0.9.5" @@ -458,6 +542,7 @@ dependencies = [ "chrono", "libloading", "once_cell", + "rusqlite", "serde", "serde_json", "tungstenite", diff --git a/src/database.rs b/src/database.rs new file mode 100644 index 0000000..9c956f5 --- /dev/null +++ b/src/database.rs @@ -0,0 +1,83 @@ +use crate::ServerConfig; +use rusqlite::{Connection, Result, params}; + +pub struct Database { + pub messages_db: MessagesDb, +} + +impl Database { + pub fn new(config: &ServerConfig) -> Option { + let conn = Connection::open("main.db").ok()?; + let messages_db = MessagesDb(conn); + messages_db.init(config)?; + Some(Self { messages_db }) + } +} + +unsafe impl Send for Database {} +unsafe impl Sync for Database {} + +pub struct MessagesDb(pub Connection); + +impl MessagesDb { + pub fn init(&self, _config: &ServerConfig) -> Option { + self.0 + .execute( + "CREATE TABLE IF NOT EXISTS chat ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + channel_id TEXT NOT NULL, + user_id TEXT NOT NULL, + contents TEXT NOT NULL, + timestamp INTEGER NOT NULL + )", + [], + ) + .ok() + } + + /// Insert a message into the DB + pub fn insert( + &self, + channel_id: &str, + user_id: &str, + contents: &str, + timestamp: i64, + ) -> Result { + self.0.execute( + "INSERT INTO chat (channel_id, user_id, contents, timestamp) + VALUES (?1, ?2, ?3, ?4)", + params![channel_id, user_id, contents, timestamp], + ) + } + + /// Fetch the latest N messages for a channel + pub fn fetch_recent( + &self, + channel_id: &str, + limit: usize, + ) -> Result> { + let mut stmt = self.0.prepare( + "SELECT id, channel_id, user_id, contents, timestamp + FROM chat + WHERE channel_id = ?1 + ORDER BY timestamp DESC + LIMIT ?2", + )?; + + let rows = stmt.query_map(params![channel_id, limit], |row| { + Ok(( + row.get::<_, i64>(0)?, // id + row.get::<_, String>(1)?, // channel_id + row.get::<_, String>(2)?, // user_id + row.get::<_, String>(3)?, // contents + row.get::<_, i64>(4)?, // timestamp + )) + })?; + + let mut results = Vec::new(); + for row in rows { + results.push(row?); + } + Ok(results) + } +} diff --git a/src/lib.rs b/src/lib.rs index 8063e98..4e6620b 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,6 +6,7 @@ use std::{ }; pub mod client; +pub mod database; #[cfg(feature = "loader")] pub mod loader; pub mod logger; @@ -27,6 +28,7 @@ pub use once_cell; #[derive(serde::Serialize, serde::Deserialize)] pub struct ServerConfig { port: u16, + channels: Vec, } #[allow(dead_code)] @@ -35,11 +37,15 @@ pub struct Server { config: ServerConfig, plugins: Mutex>, clients: Mutex>, + pub db: database::Database, } impl Default for ServerConfig { fn default() -> Self { - Self { port: 7080 } + Self { + port: 7080, + channels: Vec::new(), + } } } @@ -53,16 +59,12 @@ impl Server { logger!(LOGGER "Server"); pub fn new(root: &Path) -> Arc { - Arc::new(Self { - plugins: Mutex::new(Vec::new()), - root: root.to_path_buf(), - config: ServerConfig::default(), - clients: Mutex::new(HashSet::new()), - }) + Self::new_config(root, ServerConfig::default()) } pub fn new_config(root: &Path, config: ServerConfig) -> Arc { Arc::new(Self { + db: database::Database::new(&config).unwrap(), plugins: Mutex::new(Vec::new()), root: root.to_path_buf(), config, @@ -128,6 +130,15 @@ impl Server { contents, } => { Self::LOGGER.info(format!("SendMessage to {channel_id}: {contents}")); + self.wrap_err( + &client, + self.db.messages_db.insert( + &channel_id, + "idk", + &contents, + chrono::Utc::now().timestamp(), + ), + )?; } ClientMessage::EditMessage { diff --git a/src/types.rs b/src/types.rs index 616dc8a..af1bb0d 100644 --- a/src/types.rs +++ b/src/types.rs @@ -77,7 +77,7 @@ pub mod data { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Channel { - pub id: usize, + pub id: String, pub name: String, pub kind: ChannelKind, } From a1b54fda6bbf3f2a3ecabe0f745aa28cef6b30d3 Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sat, 13 Sep 2025 17:59:12 +0200 Subject: [PATCH 04/11] Small stuff --- client.js | 51 ++++++++++++++++++++++++++++---------------- src/client.rs | 14 ++++++++++--- src/database.rs | 56 +++++++++++++++++++++++++++++-------------------- src/lib.rs | 9 ++++++-- src/types.rs | 17 +++++++++++---- 5 files changed, 97 insertions(+), 50 deletions(-) diff --git a/client.js b/client.js index 44c6df7..727bda0 100644 --- a/client.js +++ b/client.js @@ -1,21 +1,36 @@ -setInterval(() => { - const ws = new WebSocket('ws://localhost:7080'); +import readline from 'readline'; - function sendMessage(message) { - ws.send(JSON.stringify({ type: 'send_message', params: { - channel_id: 'Hello', - contents: message - }})); - } +const rl = readline.createInterface({ + input: process.stdin, + output: process.stdout +}); - ws.onopen = () => { - console.log('WebSocket connection established'); - sendMessage('Hello, Server!'); - }; +const ws = new WebSocket('ws://localhost:7080'); - ws.onmessage = (event) => { - const message = JSON.parse(event.data); - console.log('Received:', message); - ws.close(); - }; -}, 100) \ No newline at end of file +ws.onopen = () => { + console.log('WebSocket connection established'); +}; + +ws.onmessage = (event) => { + const message = JSON.parse(event.data); + console.log('Received:', message); + // ws.close(); +}; + +function sendMessage(message) { + ws.send(JSON.stringify({ type: 'send_message', params: { + channel_id: 'Hello', + contents: message + }})); +} + +function ask(question) { + return new Promise((resolve) => { + rl.question(question, resolve); + }); +} + +while (true) { + const message = await ask(''); + sendMessage(message); +} \ No newline at end of file diff --git a/src/client.rs b/src/client.rs index 9a00b8c..7155d42 100644 --- a/src/client.rs +++ b/src/client.rs @@ -7,7 +7,7 @@ use std::{ use anyhow::Error; use tungstenite::{Message, Utf8Bytes, WebSocket, accept}; -use crate::types::{ClientMessage, ServerMessage, WsMessage}; +use crate::types::{ClientMessage, ServerMessage, WsMessage, data::ResponseError}; #[derive(Clone)] pub struct Client(Arc>>); @@ -18,7 +18,7 @@ impl Client { } pub fn new_tcp(ws: TcpStream) -> crate::Result { - Ok(Self(Arc::new(Mutex::new(accept(ws)?)))) + Ok(Self::new_ws(accept(ws)?)) } } @@ -42,7 +42,7 @@ impl Client { Message::Text(t) => { let v = t.to_string(); match serde_json::from_str(&v) { - Ok(f) => Ok(Some(WsMessage::FromClient(f))), + Ok(f) => Ok(Some(WsMessage::Message(f))), Err(_) => Ok(Some(WsMessage::String(v))), } } @@ -62,4 +62,12 @@ impl Client { .send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?))) .map_err(|e| e.into()) } + + pub fn send_err(&self, m: ResponseError) -> crate::Result<()> { + self.0 + .lock() + .unwrap() + .send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?))) + .map_err(|e| e.into()) + } } diff --git a/src/database.rs b/src/database.rs index 9c956f5..b76bdac 100644 --- a/src/database.rs +++ b/src/database.rs @@ -1,4 +1,4 @@ -use crate::ServerConfig; +use crate::{ServerConfig, types::data::Message}; use rusqlite::{Connection, Result, params}; pub struct Database { @@ -42,29 +42,33 @@ impl MessagesDb { user_id: &str, contents: &str, timestamp: i64, - ) -> Result { + ) -> Result { self.0.execute( "INSERT INTO chat (channel_id, user_id, contents, timestamp) - VALUES (?1, ?2, ?3, ?4)", + VALUES (?1, ?2, ?3, ?4)", params![channel_id, user_id, contents, timestamp], - ) - } - - /// Fetch the latest N messages for a channel - pub fn fetch_recent( - &self, - channel_id: &str, - limit: usize, - ) -> Result> { - let mut stmt = self.0.prepare( - "SELECT id, channel_id, user_id, contents, timestamp - FROM chat - WHERE channel_id = ?1 - ORDER BY timestamp DESC - LIMIT ?2", )?; - let rows = stmt.query_map(params![channel_id, limit], |row| { + let id = self.0.last_insert_rowid(); + + Ok(Message { + id, + channel_id: channel_id.to_string(), + from: user_id.to_string(), + contents: contents.to_string(), + timestamp, + }) + } + + /// Get a message by its ID + pub fn get_by_id(&self, message_id: usize) -> Result> { + let mut stmt = self.0.prepare( + "SELECT id, channel_id, user_id, contents, timestamp + FROM chat + WHERE id = ?1", + )?; + + let mut rows = stmt.query_map(params![message_id], |row| { Ok(( row.get::<_, i64>(0)?, // id row.get::<_, String>(1)?, // channel_id @@ -74,10 +78,16 @@ impl MessagesDb { )) })?; - let mut results = Vec::new(); - for row in rows { - results.push(row?); + if let Some(row) = rows.next() { + let (id, channel_id, user_id, contents, timestamp) = row?; + return Ok(Some(Message { + id, + channel_id, + from: user_id, + contents, + timestamp, + })); } - Ok(results) + Ok(None) } } diff --git a/src/lib.rs b/src/lib.rs index 4e6620b..0863ed5 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -124,13 +124,13 @@ impl Server { // The main req/res loop loop { match client.read()? { - Some(WsMessage::FromClient(req)) => match req { + Some(WsMessage::Message(req)) => match req { ClientMessage::SendMessage { channel_id, contents, } => { Self::LOGGER.info(format!("SendMessage to {channel_id}: {contents}")); - self.wrap_err( + let msg = self.wrap_err( &client, self.db.messages_db.insert( &channel_id, @@ -139,6 +139,11 @@ impl Server { chrono::Utc::now().timestamp(), ), )?; + + self.wrap_err( + &client, + client.send(types::ServerMessage::MessageCreate(msg)), + )?; } ClientMessage::EditMessage { diff --git a/src/types.rs b/src/types.rs index af1bb0d..9694079 100644 --- a/src/types.rs +++ b/src/types.rs @@ -33,7 +33,7 @@ pub enum ServerMessage { Authenticated { user_id: String }, /// Error responses - Error { message: String }, + Error(), /// A new message in a channel MessageCreate(data::Message), @@ -57,7 +57,7 @@ pub enum ServerMessage { /// WebSocket wrapper #[derive(Debug, Clone)] pub enum WsMessage Deserialize<'de>> { - FromClient(T), + Message(T), Binary(Bytes), String(String), } @@ -68,8 +68,8 @@ pub mod data { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Message { - pub id: String, - pub channel_id: u8, + pub id: i64, + pub channel_id: String, pub from: String, pub contents: String, pub timestamp: i64, @@ -88,4 +88,13 @@ pub mod data { Text, Voice, } + + #[derive(Debug, Clone, Serialize, Deserialize)] + #[serde(tag = "error", rename_all = "snake_case")] + pub enum ResponseError { + InvalidRequest { message: String }, + Unauthorized { message: String }, + NotFound { message: String }, + InternalError { message: String }, + } } From e1585bedab15c3d3afe6a15366a35341d90ee076 Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sun, 14 Sep 2025 11:08:02 +0200 Subject: [PATCH 05/11] Custom websocket --- Cargo.toml | 2 + cli/Cargo.lock | 8 +++ src/client.rs | 180 +++++++++++++++++++++++++++++++++++++++---------- src/lib.rs | 2 +- src/types.rs | 3 +- 5 files changed, 157 insertions(+), 38 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index bd412a9..fc91b01 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,12 +5,14 @@ edition = "2024" [dependencies] anyhow = "1.0.99" +base64 = "0.22.1" chrono = "0.4.42" libloading = { version = "0.8.8", optional = true } once_cell = "1.21.3" rusqlite = "0.37.0" serde = { version = "1.0.219", features = ["serde_derive"] } serde_json = "1.0.143" +sha1 = "0.10.6" tungstenite = "0.27.0" [features] diff --git a/cli/Cargo.lock b/cli/Cargo.lock index d83c247..76abc67 100644 --- a/cli/Cargo.lock +++ b/cli/Cargo.lock @@ -23,6 +23,12 @@ version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" +[[package]] +name = "base64" +version = "0.22.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" + [[package]] name = "bitflags" version = "2.9.4" @@ -539,12 +545,14 @@ name = "voxa-server" version = "0.1.0" dependencies = [ "anyhow", + "base64", "chrono", "libloading", "once_cell", "rusqlite", "serde", "serde_json", + "sha1", "tungstenite", ] diff --git a/src/client.rs b/src/client.rs index 7155d42..d20988f 100644 --- a/src/client.rs +++ b/src/client.rs @@ -1,30 +1,77 @@ use std::{ hash::{Hash, Hasher}, + io::{Read, Write}, net::TcpStream, - sync::{Arc, Mutex}, }; -use anyhow::Error; -use tungstenite::{Message, Utf8Bytes, WebSocket, accept}; +use anyhow::anyhow; +use serde::Serialize; -use crate::types::{ClientMessage, ServerMessage, WsMessage, data::ResponseError}; +use crate::types::{ClientMessage, WsMessage}; -#[derive(Clone)] -pub struct Client(Arc>>); +pub mod handshake { + use base64::Engine; + use base64::engine::general_purpose::STANDARD as Base64; + use sha1::{Digest, Sha1}; + use std::io::{Read, Write}; + use std::net::TcpStream; + + const WS_GUID: &str = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; + + pub fn handle_websocket_handshake(stream: &mut TcpStream) -> std::io::Result<()> { + let mut buffer = [0; 1024]; + let size = stream.read(&mut buffer)?; + let request = String::from_utf8_lossy(&buffer[..size]); + + let key_line = request + .lines() + .find(|line| line.to_lowercase().starts_with("sec-websocket-key")) + .ok_or_else(|| { + std::io::Error::new(std::io::ErrorKind::InvalidData, "Missing Sec-WebSocket-Key") + })?; + + let key = key_line.splitn(2, ':').nth(1).unwrap().trim(); + + let mut hasher = Sha1::new(); + hasher.update(key.as_bytes()); + hasher.update(WS_GUID.as_bytes()); + let hash = hasher.finalize(); + + let accept_key = Base64.encode(hash); + + let response = format!( + "HTTP/1.1 101 Switching Protocols\r\n\ + Upgrade: websocket\r\n\ + Connection: Upgrade\r\n\ + Sec-WebSocket-Accept: {}\r\n\r\n", + accept_key + ); + + stream.write_all(response.as_bytes())?; + stream.flush()?; + + Ok(()) + } +} + +pub struct Client(TcpStream); impl Client { - pub fn new_ws(ws: WebSocket) -> Self { - Self(Arc::new(Mutex::new(ws))) + pub fn new(mut stream: TcpStream) -> crate::Result { + handshake::handle_websocket_handshake(&mut stream)?; + Ok(Client(stream)) } +} - pub fn new_tcp(ws: TcpStream) -> crate::Result { - Ok(Self::new_ws(accept(ws)?)) +impl Clone for Client { + fn clone(&self) -> Self { + Client(self.0.try_clone().expect("failed to clone TcpStream")) } } impl PartialEq for Client { fn eq(&self, other: &Self) -> bool { - Arc::ptr_eq(&self.0, &other.0) + self.0.peer_addr().unwrap() == other.0.peer_addr().unwrap() } } @@ -32,42 +79,105 @@ impl Eq for Client {} impl Hash for Client { fn hash(&self, state: &mut H) { - std::ptr::hash(Arc::as_ptr(&self.0), state) + self.0.peer_addr().unwrap().hash(state); } } impl Client { + /// Read a full WebSocket message, handling fragmentation (FIN) pub fn read(&self) -> crate::Result>> { - match self.0.lock().unwrap().read()? { - Message::Text(t) => { - let v = t.to_string(); - match serde_json::from_str(&v) { - Ok(f) => Ok(Some(WsMessage::Message(f))), - Err(_) => Ok(Some(WsMessage::String(v))), + let mut stream = &self.0; + let mut message_payload = Vec::new(); + let mut final_frame = false; + + while !final_frame { + let mut header = [0u8; 2]; + if stream.read_exact(&mut header).is_err() { + return Ok(None); // connection closed + } + + let fin = header[0] & 0x80 != 0; + let opcode = header[0] & 0x0F; + let masked = header[1] & 0x80 != 0; + let mut payload_len = (header[1] & 0x7F) as u64; + + // Extended payload lengths + if payload_len == 126 { + let mut ext_len = [0u8; 2]; + stream.read_exact(&mut ext_len)?; + payload_len = u16::from_be_bytes(ext_len) as u64; + } else if payload_len == 127 { + let mut ext_len = [0u8; 8]; + stream.read_exact(&mut ext_len)?; + payload_len = u64::from_be_bytes(ext_len); + } + + // Mask key (client → server) + let mut mask = [0u8; 4]; + if masked { + stream.read_exact(&mut mask)?; + } + + // Read payload + let mut payload = vec![0u8; payload_len as usize]; + stream.read_exact(&mut payload)?; + + if masked { + for i in 0..payload.len() { + payload[i] ^= mask[i % 4]; } } - Message::Binary(b) => Ok(Some(WsMessage::Binary(b))), + match opcode { + 0x0 | 0x1 | 0x2 => { + // Continuation / Text / Binary + message_payload.extend(payload); + } + 0x8 => return Ok(None), // Close + 0x9 => continue, // Ping → ignore + 0xA => continue, // Pong → ignore + _ => return Err(anyhow!("Unsupported WebSocket opcode: {}", opcode).into()), + } - Message::Close(_) => Ok(None), - - m => Err(Error::msg(format!("Invalid websocket format: {m}"))), + final_frame = fin; } + + // Try parsing JSON into ClientMessage + let message = match String::from_utf8(message_payload.clone()) { + Ok(text) => match serde_json::from_str(&text) { + Ok(msg) => WsMessage::Message(msg), + Err(_) => WsMessage::String(text), + }, + Err(_) => WsMessage::Binary(message_payload), + }; + + Ok(Some(message)) } - pub fn send(&self, m: ServerMessage) -> crate::Result<()> { - self.0 - .lock() - .unwrap() - .send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?))) - .map_err(|e| e.into()) - } + /// Send a JSON-serializable object as a WebSocket text frame + pub fn send(&self, m: T) -> crate::Result<()> { + let payload = serde_json::to_string(&m)?; + let payload_bytes = payload.as_bytes(); - pub fn send_err(&self, m: ResponseError) -> crate::Result<()> { - self.0 - .lock() - .unwrap() - .send(Message::Text(Utf8Bytes::from(serde_json::to_string(&m)?))) - .map_err(|e| e.into()) + let mut stream = self.0.try_clone()?; + let mut header = Vec::new(); + header.push(0x81); // FIN=1, opcode=0x1 (text) + + let len = payload_bytes.len(); + if len < 126 { + header.push(len as u8); + } else if len <= 65535 { + header.push(126); + header.extend_from_slice(&(len as u16).to_be_bytes()); + } else { + header.push(127); + header.extend_from_slice(&(len as u64).to_be_bytes()); + } + + stream.write_all(&header)?; + stream.write_all(payload_bytes)?; + stream.flush()?; + + Ok(()) } } diff --git a/src/lib.rs b/src/lib.rs index 0863ed5..120b5a7 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -116,7 +116,7 @@ impl Server { fn handle_client(self: &Arc, stream: TcpStream) -> anyhow::Result<()> { Self::LOGGER.info(format!("New connection: {}", stream.peer_addr()?)); // Initialize client - let client = Client::new_tcp(stream)?; + let client = Client::new(stream)?; // Insert to the set of all connected clients self.clients.lock().unwrap().insert(client.clone()); diff --git a/src/types.rs b/src/types.rs index 9694079..b4f0855 100644 --- a/src/types.rs +++ b/src/types.rs @@ -1,5 +1,4 @@ use serde::{Deserialize, Serialize}; -use tungstenite::Bytes; /// Messages sent *from the client* (user’s app) to the server #[derive(Debug, Clone, Serialize, Deserialize)] @@ -58,7 +57,7 @@ pub enum ServerMessage { #[derive(Debug, Clone)] pub enum WsMessage Deserialize<'de>> { Message(T), - Binary(Bytes), + Binary(Vec), String(String), } From 59e96fe05e9c65efc1d474297d825dcfb6de9f8e Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sun, 14 Sep 2025 11:12:47 +0200 Subject: [PATCH 06/11] Working broadcast --- src/lib.rs | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/src/lib.rs b/src/lib.rs index 120b5a7..340b558 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -140,10 +140,12 @@ impl Server { ), )?; - self.wrap_err( - &client, - client.send(types::ServerMessage::MessageCreate(msg)), - )?; + for c in self.clients.lock().unwrap().iter() { + self.wrap_err( + &c, + c.send(types::ServerMessage::MessageCreate(msg.clone())), + )?; + } } ClientMessage::EditMessage { From a4ce0c05890906b2111580225483d6956e193bcf Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sun, 14 Sep 2025 14:20:36 +0200 Subject: [PATCH 07/11] Improved structure --- Cargo.toml | 1 - cli/Cargo.lock | 170 ---------------- cli/src/main.rs | 2 +- client.js | 2 +- src/lib.rs | 32 +-- src/macros.rs | 14 +- src/types.rs | 18 +- src/{ => utils}/client.rs | 0 src/{ => utils}/database.rs | 0 src/{ => utils}/loader.rs | 2 +- src/{ => utils}/logger.rs | 0 src/utils/mod.rs | 7 + src/{ => utils}/plugin.rs | 8 +- src/{ => utils}/vfs.rs | 0 test-plugin/Cargo.lock | 391 +++++++++++++++++++++++++----------- test-plugin/src/lib.rs | 29 ++- test-plugin/src/main.rs | 2 +- 17 files changed, 357 insertions(+), 321 deletions(-) rename src/{ => utils}/client.rs (100%) rename src/{ => utils}/database.rs (100%) rename src/{ => utils}/loader.rs (95%) rename src/{ => utils}/logger.rs (100%) create mode 100644 src/utils/mod.rs rename src/{ => utils}/plugin.rs (59%) rename src/{ => utils}/vfs.rs (100%) diff --git a/Cargo.toml b/Cargo.toml index fc91b01..12fe364 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -13,7 +13,6 @@ rusqlite = "0.37.0" serde = { version = "1.0.219", features = ["serde_derive"] } serde_json = "1.0.143" sha1 = "0.10.6" -tungstenite = "0.27.0" [features] default = [] diff --git a/cli/Cargo.lock b/cli/Cargo.lock index 76abc67..30944fb 100644 --- a/cli/Cargo.lock +++ b/cli/Cargo.lock @@ -50,12 +50,6 @@ version = "3.19.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "46c5e41b57b8bba42a04676d81cb89e9ee8e859a1a66f80a5a72e1cb76b34d43" -[[package]] -name = "bytes" -version = "1.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d71b6127be86fdcfddb610f7182ac57211d4b18a3e9c82eb2d17662f2227ad6a" - [[package]] name = "cc" version = "1.2.37" @@ -117,12 +111,6 @@ dependencies = [ "typenum", ] -[[package]] -name = "data-encoding" -version = "2.9.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2a2330da5de22e8a3cb63252ce2abb30116bf5265e89c0e01bc17015ce30a476" - [[package]] name = "digest" version = "0.10.7" @@ -151,12 +139,6 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7fd99930f64d146689264c637b5af2f0233a933bef0d8570e2526bf9e083192d" -[[package]] -name = "fnv" -version = "1.0.7" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" - [[package]] name = "foldhash" version = "0.1.5" @@ -173,18 +155,6 @@ dependencies = [ "version_check", ] -[[package]] -name = "getrandom" -version = "0.3.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26145e563e54f2cadc477553f1ec5ee650b00862f0a58bcd12cbdc5f0ea2d2f4" -dependencies = [ - "cfg-if", - "libc", - "r-efi", - "wasi", -] - [[package]] name = "hashbrown" version = "0.15.5" @@ -203,23 +173,6 @@ dependencies = [ "hashbrown", ] -[[package]] -name = "http" -version = "1.3.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f4a85d31aea989eead29a3aaf9e1115a180df8282431156e533de47660892565" -dependencies = [ - "bytes", - "fnv", - "itoa", -] - -[[package]] -name = "httparse" -version = "1.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" - [[package]] name = "iana-time-zone" version = "0.1.64" @@ -319,15 +272,6 @@ version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" -[[package]] -name = "ppv-lite86" -version = "0.2.21" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" -dependencies = [ - "zerocopy", -] - [[package]] name = "proc-macro2" version = "1.0.101" @@ -346,41 +290,6 @@ dependencies = [ "proc-macro2", ] -[[package]] -name = "r-efi" -version = "5.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" - -[[package]] -name = "rand" -version = "0.9.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6db2770f06117d490610c7488547d543617b21bfa07796d7a12f6f1bd53850d1" -dependencies = [ - "rand_chacha", - "rand_core", -] - -[[package]] -name = "rand_chacha" -version = "0.9.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" -dependencies = [ - "ppv-lite86", - "rand_core", -] - -[[package]] -name = "rand_core" -version = "0.9.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" -dependencies = [ - "getrandom", -] - [[package]] name = "rusqlite" version = "0.37.0" @@ -473,43 +382,6 @@ dependencies = [ "unicode-ident", ] -[[package]] -name = "thiserror" -version = "2.0.16" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3467d614147380f2e4e374161426ff399c91084acd2363eaf549172b3d5e60c0" -dependencies = [ - "thiserror-impl", -] - -[[package]] -name = "thiserror-impl" -version = "2.0.16" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6c5e1be1c48b9172ee610da68fd9cd2770e7a4056cb3fc98710ee6906f0c7960" -dependencies = [ - "proc-macro2", - "quote", - "syn", -] - -[[package]] -name = "tungstenite" -version = "0.27.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "eadc29d668c91fcc564941132e17b28a7ceb2f3ebf0b9dae3e03fd7a6748eb0d" -dependencies = [ - "bytes", - "data-encoding", - "http", - "httparse", - "log", - "rand", - "sha1", - "thiserror", - "utf-8", -] - [[package]] name = "typenum" version = "1.18.0" @@ -522,12 +394,6 @@ version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5a5f39404a5da50712a4c1eecf25e90dd62b613502b7e925fd4e4d19b5c96512" -[[package]] -name = "utf-8" -version = "0.7.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" - [[package]] name = "vcpkg" version = "0.2.15" @@ -553,16 +419,6 @@ dependencies = [ "serde", "serde_json", "sha1", - "tungstenite", -] - -[[package]] -name = "wasi" -version = "0.14.3+wasi-0.2.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a51ae83037bdd272a9e28ce236db8c07016dd0d50c27038b3f407533c030c95" -dependencies = [ - "wit-bindgen", ] [[package]] @@ -753,29 +609,3 @@ name = "windows_x86_64_msvc" version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "271414315aff87387382ec3d271b52d7ae78726f5d44ac98b4f4030c91880486" - -[[package]] -name = "wit-bindgen" -version = "0.45.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "052283831dbae3d879dc7f51f3d92703a316ca49f91540417d38591826127814" - -[[package]] -name = "zerocopy" -version = "0.8.26" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1039dd0d3c310cf05de012d8a39ff557cb0d23087fd44cad61df08fc31907a2f" -dependencies = [ - "zerocopy-derive", -] - -[[package]] -name = "zerocopy-derive" -version = "0.8.26" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9ecf5b4cc5364572d7f4c329661bcc82724222973f2cab6f050a4e5c22f75181" -dependencies = [ - "proc-macro2", - "quote", - "syn", -] diff --git a/cli/src/main.rs b/cli/src/main.rs index f87dec1..a5c5dc3 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -1,6 +1,6 @@ use std::path::PathBuf; -use voxa_server::{ServerConfig, vfs}; +use voxa_server::{ServerConfig, utils::vfs}; fn main() -> voxa_server::Result<()> { let root = PathBuf::from(""); diff --git a/client.js b/client.js index 727bda0..e1b8742 100644 --- a/client.js +++ b/client.js @@ -19,7 +19,7 @@ ws.onmessage = (event) => { function sendMessage(message) { ws.send(JSON.stringify({ type: 'send_message', params: { - channel_id: 'Hello', + channel_id: 'general', contents: message }})); } diff --git a/src/lib.rs b/src/lib.rs index 340b558..734d2fd 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -5,23 +5,16 @@ use std::{ sync::{Arc, Mutex}, }; -pub mod client; -pub mod database; -#[cfg(feature = "loader")] -pub mod loader; -pub mod logger; pub mod macros; -pub mod plugin; pub mod types; -pub mod vfs; +pub mod utils; pub use anyhow::Result; -pub use tungstenite; use crate::{ - client::Client, - plugin::DynPlugin, types::{ClientMessage, WsMessage}, + utils::client::Client, + utils::plugin::DynPlugin, }; pub use once_cell; @@ -37,7 +30,7 @@ pub struct Server { config: ServerConfig, plugins: Mutex>, clients: Mutex>, - pub db: database::Database, + pub db: utils::database::Database, } impl Default for ServerConfig { @@ -64,7 +57,7 @@ impl Server { pub fn new_config(root: &Path, config: ServerConfig) -> Arc { Arc::new(Self { - db: database::Database::new(&config).unwrap(), + db: utils::database::Database::new(&config).unwrap(), plugins: Mutex::new(Vec::new()), root: root.to_path_buf(), config, @@ -75,7 +68,7 @@ impl Server { pub fn run(self: &Arc) -> Result<()> { // Load plugins #[cfg(feature = "loader")] - loader::load_plugins( + utils::loader::load_plugins( &mut *self.plugins.lock().unwrap(), &self.root.join("./plugins"), )?; @@ -122,8 +115,17 @@ impl Server { self.clients.lock().unwrap().insert(client.clone()); // The main req/res loop - loop { - match client.read()? { + 'outer: loop { + let req = client.read()?; + if let Some(r) = &req { + for p in self.plugins.lock().unwrap().iter_mut() { + if p.on_request(r, &client, self) { + continue 'outer; + } + } + } + + match req { Some(WsMessage::Message(req)) => match req { ClientMessage::SendMessage { channel_id, diff --git a/src/macros.rs b/src/macros.rs index 0a8fb8c..fd1d31e 100644 --- a/src/macros.rs +++ b/src/macros.rs @@ -2,7 +2,7 @@ macro_rules! export_plugin { ($p:expr) => { #[unsafe(no_mangle)] - pub extern "C" fn load_plugin() -> $crate::plugin::DynPlugin { + pub extern "C" fn load_plugin() -> $crate::utils::plugin::DynPlugin { $p } }; @@ -11,20 +11,20 @@ macro_rules! export_plugin { #[macro_export] macro_rules! logger { (const $i:ident $name:expr) => { - const $i: $crate::once_cell::sync::Lazy<$crate::logger::Logger> = - $crate::once_cell::sync::Lazy::new(|| $crate::logger::Logger::new($name)); + pub const $i: $crate::once_cell::sync::Lazy<$crate::utils::logger::Logger> = + $crate::once_cell::sync::Lazy::new(|| $crate::utils::logger::Logger::new($name)); }; ($i:ident $name:expr) => { - const $i: $crate::once_cell::sync::Lazy<$crate::logger::Logger> = - $crate::once_cell::sync::Lazy::new(|| $crate::logger::Logger::new($name)); + pub const $i: $crate::once_cell::sync::Lazy<$crate::utils::logger::Logger> = + $crate::once_cell::sync::Lazy::new(|| $crate::utils::logger::Logger::new($name)); }; (const $name:expr) => { - $crate::once_cell::sync::Lazy::new(|| $crate::logger::Logger::new($name)) + $crate::once_cell::sync::Lazy::new(|| $crate::utils::logger::Logger::new($name)) }; ($name:expr) => { - $crate::logger::Logger::new($name) + $crate::logger::utils::Logger::new($name) }; } diff --git a/src/types.rs b/src/types.rs index b4f0855..2fd8636 100644 --- a/src/types.rs +++ b/src/types.rs @@ -29,11 +29,17 @@ pub enum ClientMessage { #[serde(tag = "type", content = "params", rename_all = "snake_case")] pub enum ServerMessage { /// Successful authentication - Authenticated { user_id: String }, + Authenticated { + user_id: String, + }, /// Error responses Error(), + TempMessage { + message: String, + }, + /// A new message in a channel MessageCreate(data::Message), @@ -47,10 +53,16 @@ pub enum ServerMessage { }, /// Presence updates - PresenceUpdate { user_id: String, status: String }, + PresenceUpdate { + user_id: String, + status: String, + }, /// Typing indicator - Typing { user_id: String, channel_id: String }, + Typing { + user_id: String, + channel_id: String, + }, } /// WebSocket wrapper diff --git a/src/client.rs b/src/utils/client.rs similarity index 100% rename from src/client.rs rename to src/utils/client.rs diff --git a/src/database.rs b/src/utils/database.rs similarity index 100% rename from src/database.rs rename to src/utils/database.rs diff --git a/src/loader.rs b/src/utils/loader.rs similarity index 95% rename from src/loader.rs rename to src/utils/loader.rs index 9e356cd..70d52ac 100644 --- a/src/loader.rs +++ b/src/utils/loader.rs @@ -2,7 +2,7 @@ use std::path::Path; use libloading::{Library, Symbol}; -use crate::{logger, plugin::DynPlugin, vfs}; +use crate::{logger, utils::plugin::DynPlugin, utils::vfs}; logger! { const LOGGER "Loader" diff --git a/src/logger.rs b/src/utils/logger.rs similarity index 100% rename from src/logger.rs rename to src/utils/logger.rs diff --git a/src/utils/mod.rs b/src/utils/mod.rs new file mode 100644 index 0000000..d7279fe --- /dev/null +++ b/src/utils/mod.rs @@ -0,0 +1,7 @@ +pub mod client; +pub mod database; +#[cfg(feature = "loader")] +pub mod loader; +pub mod logger; +pub mod plugin; +pub mod vfs; diff --git a/src/plugin.rs b/src/utils/plugin.rs similarity index 59% rename from src/plugin.rs rename to src/utils/plugin.rs index 55495df..033e173 100644 --- a/src/plugin.rs +++ b/src/utils/plugin.rs @@ -3,6 +3,7 @@ use std::sync::Arc; use crate::{ Server, types::{ClientMessage, WsMessage}, + utils::client::Client, }; pub type DynPlugin = Box; @@ -10,7 +11,12 @@ pub type DynPlugin = Box; pub trait Plugin { fn init(&mut self, server: &Arc); #[allow(unused_variables)] - fn on_request(&mut self, msg: &WsMessage, server: &Arc) -> bool { + fn on_request( + &mut self, + req: &WsMessage, + client: &Client, + server: &Arc, + ) -> bool { false } } diff --git a/src/vfs.rs b/src/utils/vfs.rs similarity index 100% rename from src/vfs.rs rename to src/utils/vfs.rs diff --git a/test-plugin/Cargo.lock b/test-plugin/Cargo.lock index 4bff9b1..679c798 100644 --- a/test-plugin/Cargo.lock +++ b/test-plugin/Cargo.lock @@ -2,12 +2,39 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "android_system_properties" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "819e7219dbd41043ac279b19830f2efc897156490d7fd6ea916720117ee66311" +dependencies = [ + "libc", +] + [[package]] name = "anyhow" version = "1.0.99" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b0674a1ddeecb70197781e945de4b3b8ffb61fa939a5597bcf48503737663100" +[[package]] +name = "autocfg" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" + +[[package]] +name = "base64" +version = "0.22.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" + +[[package]] +name = "bitflags" +version = "2.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2261d10cca569e4643e526d8dc2e62e433cc8aba21ab764233731f8d369bf394" + [[package]] name = "block-buffer" version = "0.10.4" @@ -18,10 +45,20 @@ dependencies = [ ] [[package]] -name = "bytes" -version = "1.10.1" +name = "bumpalo" +version = "3.19.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d71b6127be86fdcfddb610f7182ac57211d4b18a3e9c82eb2d17662f2227ad6a" +checksum = "46c5e41b57b8bba42a04676d81cb89e9ee8e859a1a66f80a5a72e1cb76b34d43" + +[[package]] +name = "cc" +version = "1.2.37" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "65193589c6404eb80b450d618eaf9a2cafaaafd57ecce47370519ef674a7bd44" +dependencies = [ + "find-msvc-tools", + "shlex", +] [[package]] name = "cfg-if" @@ -29,6 +66,25 @@ version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2fd1289c04a9ea8cb22300a459a72a385d7c73d3259e2ed7dcb2af674838cfa9" +[[package]] +name = "chrono" +version = "0.4.42" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "145052bdd345b87320e369255277e3fb5152762ad123a901ef5c262dd38fe8d2" +dependencies = [ + "iana-time-zone", + "js-sys", + "num-traits", + "wasm-bindgen", + "windows-link", +] + +[[package]] +name = "core-foundation-sys" +version = "0.8.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" + [[package]] name = "cpufeatures" version = "0.2.17" @@ -48,12 +104,6 @@ dependencies = [ "typenum", ] -[[package]] -name = "data-encoding" -version = "2.9.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2a2330da5de22e8a3cb63252ce2abb30116bf5265e89c0e01bc17015ce30a476" - [[package]] name = "digest" version = "0.10.7" @@ -65,10 +115,28 @@ dependencies = [ ] [[package]] -name = "fnv" -version = "1.0.7" +name = "fallible-iterator" +version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" +checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649" + +[[package]] +name = "fallible-streaming-iterator" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a" + +[[package]] +name = "find-msvc-tools" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7fd99930f64d146689264c637b5af2f0233a933bef0d8570e2526bf9e083192d" + +[[package]] +name = "foldhash" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" [[package]] name = "generic-array" @@ -81,33 +149,46 @@ dependencies = [ ] [[package]] -name = "getrandom" -version = "0.3.3" +name = "hashbrown" +version = "0.15.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26145e563e54f2cadc477553f1ec5ee650b00862f0a58bcd12cbdc5f0ea2d2f4" +checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" dependencies = [ - "cfg-if", - "libc", - "r-efi", - "wasi", + "foldhash", ] [[package]] -name = "http" -version = "1.3.1" +name = "hashlink" +version = "0.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f4a85d31aea989eead29a3aaf9e1115a180df8282431156e533de47660892565" +checksum = "7382cf6263419f2d8df38c55d7da83da5c18aef87fc7a7fc1fb1e344edfe14c1" dependencies = [ - "bytes", - "fnv", - "itoa", + "hashbrown", ] [[package]] -name = "httparse" -version = "1.10.1" +name = "iana-time-zone" +version = "0.1.64" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" +checksum = "33e57f83510bb73707521ebaffa789ec8caf86f9657cad665b092b581d40e9fb" +dependencies = [ + "android_system_properties", + "core-foundation-sys", + "iana-time-zone-haiku", + "js-sys", + "log", + "wasm-bindgen", + "windows-core", +] + +[[package]] +name = "iana-time-zone-haiku" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f31827a206f56af32e590ba56d5d2d085f558508192593743f16b2306495269f" +dependencies = [ + "cc", +] [[package]] name = "itoa" @@ -115,12 +196,32 @@ version = "1.0.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4a5f13b858c8d314ee3e8f639011f7ccefe71f97f96e50151fb991f267928e2c" +[[package]] +name = "js-sys" +version = "0.3.78" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c0b063578492ceec17683ef2f8c5e89121fbd0b172cbc280635ab7567db2738" +dependencies = [ + "once_cell", + "wasm-bindgen", +] + [[package]] name = "libc" version = "0.2.175" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6a82ae493e598baaea5209805c49bbf2ea7de956d50d7da0da1164f9c6d28543" +[[package]] +name = "libsqlite3-sys" +version = "0.35.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "133c182a6a2c87864fe97778797e46c7e999672690dc9fa3ee8e241aa4a9c13f" +dependencies = [ + "pkg-config", + "vcpkg", +] + [[package]] name = "log" version = "0.4.27" @@ -133,6 +234,15 @@ version = "2.7.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a282da65faaf38286cf3be983213fcf1d2e2a58700e808f83f4ea9a4804bc0" +[[package]] +name = "num-traits" +version = "0.2.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "071dfc062690e90b734c0b2273ce72ad0ffa95f0c74596bc250dcfd960262841" +dependencies = [ + "autocfg", +] + [[package]] name = "once_cell" version = "1.21.3" @@ -140,13 +250,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "42f5e15c9953c5e4ccceeb2e7382a716482c34515315f7b03532b8b4e8393d2d" [[package]] -name = "ppv-lite86" -version = "0.2.21" +name = "pkg-config" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" -dependencies = [ - "zerocopy", -] +checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" [[package]] name = "proc-macro2" @@ -167,39 +274,24 @@ dependencies = [ ] [[package]] -name = "r-efi" -version = "5.3.0" +name = "rusqlite" +version = "0.37.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" - -[[package]] -name = "rand" -version = "0.9.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6db2770f06117d490610c7488547d543617b21bfa07796d7a12f6f1bd53850d1" +checksum = "165ca6e57b20e1351573e3729b958bc62f0e48025386970b6e4d29e7a7e71f3f" dependencies = [ - "rand_chacha", - "rand_core", + "bitflags", + "fallible-iterator", + "fallible-streaming-iterator", + "hashlink", + "libsqlite3-sys", + "smallvec", ] [[package]] -name = "rand_chacha" -version = "0.9.0" +name = "rustversion" +version = "1.0.22" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" -dependencies = [ - "ppv-lite86", - "rand_core", -] - -[[package]] -name = "rand_core" -version = "0.9.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" -dependencies = [ - "getrandom", -] +checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" [[package]] name = "ryu" @@ -250,6 +342,18 @@ dependencies = [ "digest", ] +[[package]] +name = "shlex" +version = "1.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0fda2ff0d084019ba4d7c6f371c95d8fd75ce3524c3cb8fb653a3023f6323e64" + +[[package]] +name = "smallvec" +version = "1.15.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "67b1b7a3b5fe4f1376887184045fcf45c69e92af734b7aaddc05fb777b6fbd03" + [[package]] name = "syn" version = "2.0.106" @@ -268,43 +372,6 @@ dependencies = [ "voxa-server", ] -[[package]] -name = "thiserror" -version = "2.0.16" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3467d614147380f2e4e374161426ff399c91084acd2363eaf549172b3d5e60c0" -dependencies = [ - "thiserror-impl", -] - -[[package]] -name = "thiserror-impl" -version = "2.0.16" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6c5e1be1c48b9172ee610da68fd9cd2770e7a4056cb3fc98710ee6906f0c7960" -dependencies = [ - "proc-macro2", - "quote", - "syn", -] - -[[package]] -name = "tungstenite" -version = "0.27.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "eadc29d668c91fcc564941132e17b28a7ceb2f3ebf0b9dae3e03fd7a6748eb0d" -dependencies = [ - "bytes", - "data-encoding", - "http", - "httparse", - "log", - "rand", - "sha1", - "thiserror", - "utf-8", -] - [[package]] name = "typenum" version = "1.18.0" @@ -318,10 +385,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5a5f39404a5da50712a4c1eecf25e90dd62b613502b7e925fd4e4d19b5c96512" [[package]] -name = "utf-8" -version = "0.7.6" +name = "vcpkg" +version = "0.2.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" [[package]] name = "version_check" @@ -334,43 +401,129 @@ name = "voxa-server" version = "0.1.0" dependencies = [ "anyhow", + "base64", + "chrono", "once_cell", + "rusqlite", "serde", "serde_json", - "tungstenite", + "sha1", ] [[package]] -name = "wasi" -version = "0.14.3+wasi-0.2.4" +name = "wasm-bindgen" +version = "0.2.101" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a51ae83037bdd272a9e28ce236db8c07016dd0d50c27038b3f407533c030c95" +checksum = "7e14915cadd45b529bb8d1f343c4ed0ac1de926144b746e2710f9cd05df6603b" dependencies = [ - "wit-bindgen", + "cfg-if", + "once_cell", + "rustversion", + "wasm-bindgen-macro", + "wasm-bindgen-shared", ] [[package]] -name = "wit-bindgen" -version = "0.45.0" +name = "wasm-bindgen-backend" +version = "0.2.101" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "052283831dbae3d879dc7f51f3d92703a316ca49f91540417d38591826127814" - -[[package]] -name = "zerocopy" -version = "0.8.26" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1039dd0d3c310cf05de012d8a39ff557cb0d23087fd44cad61df08fc31907a2f" +checksum = "e28d1ba982ca7923fd01448d5c30c6864d0a14109560296a162f80f305fb93bb" dependencies = [ - "zerocopy-derive", + "bumpalo", + "log", + "proc-macro2", + "quote", + "syn", + "wasm-bindgen-shared", ] [[package]] -name = "zerocopy-derive" -version = "0.8.26" +name = "wasm-bindgen-macro" +version = "0.2.101" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9ecf5b4cc5364572d7f4c329661bcc82724222973f2cab6f050a4e5c22f75181" +checksum = "7c3d463ae3eff775b0c45df9da45d68837702ac35af998361e2c84e7c5ec1b0d" +dependencies = [ + "quote", + "wasm-bindgen-macro-support", +] + +[[package]] +name = "wasm-bindgen-macro-support" +version = "0.2.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7bb4ce89b08211f923caf51d527662b75bdc9c9c7aab40f86dcb9fb85ac552aa" +dependencies = [ + "proc-macro2", + "quote", + "syn", + "wasm-bindgen-backend", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-shared" +version = "0.2.101" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f143854a3b13752c6950862c906306adb27c7e839f7414cec8fea35beab624c1" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "windows-core" +version = "0.62.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "57fe7168f7de578d2d8a05b07fd61870d2e73b4020e9f49aa00da8471723497c" +dependencies = [ + "windows-implement", + "windows-interface", + "windows-link", + "windows-result", + "windows-strings", +] + +[[package]] +name = "windows-implement" +version = "0.60.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a47fddd13af08290e67f4acabf4b459f647552718f683a7b415d290ac744a836" dependencies = [ "proc-macro2", "quote", "syn", ] + +[[package]] +name = "windows-interface" +version = "0.59.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bd9211b69f8dcdfa817bfd14bf1c97c9188afa36f4750130fcdf3f400eca9fa8" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "windows-link" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "45e46c0661abb7180e7b9c281db115305d49ca1709ab8242adf09666d2173c65" + +[[package]] +name = "windows-result" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7084dcc306f89883455a206237404d3eaf961e5bd7e0f312f7c91f57eb44167f" +dependencies = [ + "windows-link", +] + +[[package]] +name = "windows-strings" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7218c655a553b0bed4426cf54b20d7ba363ef543b52d515b3e48d7fd55318dda" +dependencies = [ + "windows-link", +] diff --git a/test-plugin/src/lib.rs b/test-plugin/src/lib.rs index 8b53dcf..4f0ea34 100644 --- a/test-plugin/src/lib.rs +++ b/test-plugin/src/lib.rs @@ -1,5 +1,5 @@ use std::sync::Arc; -use voxa_server::{Server, export_plugin, logger, plugin::Plugin}; +use voxa_server::{Server, export_plugin, logger, utils::plugin::Plugin}; logger! { const LOGGER "My Plugin" @@ -12,6 +12,33 @@ impl Plugin for MyPlugin { fn init(&mut self, _server: &Arc) { LOGGER.info("MyPlugin initialized!"); } + + fn on_request( + &mut self, + msg: &voxa_server::types::WsMessage, + client: &voxa_server::utils::client::Client, + _server: &Arc, + ) -> bool { + LOGGER.info(&format!("Received message: {:?}", msg)); + match msg { + voxa_server::types::WsMessage::Message( + voxa_server::types::ClientMessage::SendMessage { contents, .. }, + ) => { + if contents == "ping" { + LOGGER.info("Pong!"); + client + .send(voxa_server::types::ServerMessage::TempMessage { + message: "pong".to_string(), + }) + .unwrap(); + true + } else { + false + } + } + _ => false, + } + } } export_plugin!(Box::new(MyPlugin)); diff --git a/test-plugin/src/main.rs b/test-plugin/src/main.rs index 69acdbf..c0a9e5e 100644 --- a/test-plugin/src/main.rs +++ b/test-plugin/src/main.rs @@ -1,6 +1,6 @@ use std::path::PathBuf; -use voxa_server::{ServerConfig, vfs}; +use voxa_server::{ServerConfig, utils::vfs}; fn main() -> voxa_server::Result<()> { let root = PathBuf::from("./"); From 46b7a60dd430bb7e340fceb5715f78f3efe19b56 Mon Sep 17 00:00:00 2001 From: Leo dev Date: Sun, 14 Sep 2025 15:48:29 +0200 Subject: [PATCH 08/11] Updated client.rs with better prod --- src/utils/client.rs | 381 +++++++++++++++++++++++++++++++------------- 1 file changed, 266 insertions(+), 115 deletions(-) diff --git a/src/utils/client.rs b/src/utils/client.rs index d20988f..4f1331f 100644 --- a/src/utils/client.rs +++ b/src/utils/client.rs @@ -1,7 +1,8 @@ use std::{ hash::{Hash, Hasher}, - io::{Read, Write}, + io::{self, Read, Write}, net::TcpStream, + time::Duration, }; use anyhow::anyhow; @@ -13,54 +14,303 @@ pub mod handshake { use base64::Engine; use base64::engine::general_purpose::STANDARD as Base64; use sha1::{Digest, Sha1}; - use std::io::{Read, Write}; + use std::collections::HashMap; + use std::io::{BufRead, BufReader, Write}; use std::net::TcpStream; const WS_GUID: &str = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; pub fn handle_websocket_handshake(stream: &mut TcpStream) -> std::io::Result<()> { - let mut buffer = [0; 1024]; - let size = stream.read(&mut buffer)?; - let request = String::from_utf8_lossy(&buffer[..size]); + let mut reader = BufReader::new(stream.try_clone()?); + let mut request_line = String::new(); + reader.read_line(&mut request_line)?; - let key_line = request - .lines() - .find(|line| line.to_lowercase().starts_with("sec-websocket-key")) - .ok_or_else(|| { - std::io::Error::new(std::io::ErrorKind::InvalidData, "Missing Sec-WebSocket-Key") - })?; + if !request_line.starts_with("GET") { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "Invalid HTTP method", + )); + } - let key = key_line.splitn(2, ':').nth(1).unwrap().trim(); + let mut headers = HashMap::new(); + let mut line = String::new(); + loop { + line.clear(); + let bytes = reader.read_line(&mut line)?; + if bytes == 0 || line == "\r\n" { + break; + } + if let Some((k, v)) = line.split_once(':') { + headers.insert(k.trim().to_lowercase(), v.trim().to_string()); + } + } + let key = headers.get("sec-websocket-key").ok_or_else(|| { + std::io::Error::new(std::io::ErrorKind::InvalidData, "Missing Sec-WebSocket-Key") + })?; + + if headers.get("upgrade").map(|v| v.to_lowercase()) != Some("websocket".to_string()) { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "Missing or invalid Upgrade header", + )); + } + + if !headers + .get("connection") + .map(|v| v.to_lowercase().contains("upgrade")) + .unwrap_or(false) + { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "Missing or invalid Connection header", + )); + } + + // Optional: validate Sec-WebSocket-Version == 13 (most common) + if let Some(ver) = headers.get("sec-websocket-version") { + if ver.trim() != "13" { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "Unsupported Sec-WebSocket-Version", + )); + } + } + + // Compute accept key let mut hasher = Sha1::new(); hasher.update(key.as_bytes()); hasher.update(WS_GUID.as_bytes()); let hash = hasher.finalize(); - let accept_key = Base64.encode(hash); + // Note: include Sec-WebSocket-Protocol handling if you support subprotocols let response = format!( "HTTP/1.1 101 Switching Protocols\r\n\ Upgrade: websocket\r\n\ Connection: Upgrade\r\n\ - Sec-WebSocket-Accept: {}\r\n\r\n", + Sec-WebSocket-Accept: {}\r\n\ + Sec-WebSocket-Version: 13\r\n\r\n", accept_key ); stream.write_all(response.as_bytes())?; stream.flush()?; - Ok(()) } } -pub struct Client(TcpStream); +pub struct Client(pub TcpStream); impl Client { + /// Create a client with no timeouts pub fn new(mut stream: TcpStream) -> crate::Result { handshake::handle_websocket_handshake(&mut stream)?; Ok(Client(stream)) } + + /// Create a client and set read/write timeouts (useful in prod) + pub fn with_timeouts( + mut stream: TcpStream, + read_timeout: Option, + write_timeout: Option, + ) -> crate::Result { + if let Some(t) = read_timeout { + stream.set_read_timeout(Some(t))?; + } + if let Some(t) = write_timeout { + stream.set_write_timeout(Some(t))?; + } + handshake::handle_websocket_handshake(&mut stream)?; + Ok(Client(stream)) + } + + /// Send a close frame and flush. `code` is a WebSocket close code (e.g., 1000 normal). + pub fn send_close(&self, code: u16, reason: &str) -> crate::Result<()> { + let mut stream = self.0.try_clone()?; + + // control frames must be <= 125 bytes + let mut payload = Vec::new(); + payload.extend_from_slice(&code.to_be_bytes()); + payload.extend_from_slice(reason.as_bytes()); + if payload.len() > 125 { + return Err(anyhow!("close reason too long").into()); + } + + let mut frame = Vec::with_capacity(2 + payload.len()); + frame.push(0x88); // FIN=1, opcode=0x8 (Close) + frame.push(payload.len() as u8); // server->client MUST NOT mask + frame.extend_from_slice(&payload); + + stream.write_all(&frame)?; + stream.flush()?; + Ok(()) + } + + /// Send a pong with given payload (control frames must be <=125) + fn send_pong(&self, payload: &[u8]) -> crate::Result<()> { + let mut stream = self.0.try_clone()?; + + if payload.len() > 125 { + return Err(anyhow!("pong payload too long").into()); + } + let mut frame = Vec::with_capacity(2 + payload.len()); + frame.push(0x8A); // FIN=1, opcode=0xA (Pong) + frame.push(payload.len() as u8); + frame.extend_from_slice(payload); + stream.write_all(&frame)?; + stream.flush()?; + Ok(()) + } + + /// Send a text/binary frame (server->client must NOT mask) + pub fn send(&self, m: T) -> crate::Result<()> { + let mut stream = self.0.try_clone()?; + + let payload = serde_json::to_string(&m)?; + let payload_bytes = payload.as_bytes(); + let len = payload_bytes.len(); + + let mut header = Vec::new(); + header.push(0x81); // FIN=1, opcode=0x1 (text) + + if len < 126 { + header.push(len as u8); + } else if len <= 65535 { + header.push(126); + header.extend_from_slice(&(len as u16).to_be_bytes()); + } else { + header.push(127); + header.extend_from_slice(&(len as u64).to_be_bytes()); + } + + stream.write_all(&header)?; + stream.write_all(payload_bytes)?; + stream.flush()?; + Ok(()) + } + + /// Read a full WebSocket message, handling fragmentation and control frames. + /// + /// Returns: + /// - Ok(Some(WsMessage)) on an application message (text/binary) + /// - Ok(None) if the connection should be closed (close received / read EOF) + /// - Err on protocol or IO errors. + pub fn read(&self) -> crate::Result>> { + let mut stream = self.0.try_clone()?; + + let mut message_payload = Vec::new(); + + loop { + // read the 2-byte header + let mut header = [0u8; 2]; + if let Err(e) = stream.read_exact(&mut header) { + if e.kind() == io::ErrorKind::UnexpectedEof || e.kind() == io::ErrorKind::BrokenPipe + { + return Ok(None); // treat EOF as closed + } + return Err(e.into()); + } + + let fin = header[0] & 0x80 != 0; + let opcode = header[0] & 0x0F; + let masked = header[1] & 0x80 != 0; + let mut payload_len = (header[1] & 0x7F) as u64; + + // Extended payload lengths + if payload_len == 126 { + let mut ext_len = [0u8; 2]; + stream.read_exact(&mut ext_len)?; + payload_len = u16::from_be_bytes(ext_len) as u64; + } else if payload_len == 127 { + let mut ext_len = [0u8; 8]; + stream.read_exact(&mut ext_len)?; + payload_len = u64::from_be_bytes(ext_len); + } + + // Mask key (client→server MUST be masked) + let mut mask = [0u8; 4]; + if masked { + stream.read_exact(&mut mask)?; + } else { + let _ = self.send_close(1002, "Client frames must be masked"); + return Ok(None); + } + + // Control frame checks + if matches!(opcode, 0x8 | 0x9 | 0xA) { + if payload_len > 125 { + let _ = self.send_close(1002, "Control frame too large"); + return Ok(None); + } + if !fin { + let _ = self.send_close(1002, "Control frames must not be fragmented"); + return Ok(None); + } + } + + // Read payload + unmask + let mut payload = vec![0u8; payload_len as usize]; + if payload_len > 0 { + stream.read_exact(&mut payload)?; + for i in 0..payload.len() { + payload[i] ^= mask[i % 4]; + } + } + + match opcode { + 0x0 | 0x1 | 0x2 => { + // Continuation / Text / Binary + message_payload.extend(payload); + if fin { + break; // got full message + } else { + continue; // wait for more fragments + } + } + 0x8 => { + // Close + let (code, reason) = if payload.len() >= 2 { + let code = u16::from_be_bytes([payload[0], payload[1]]); + let reason = if payload.len() > 2 { + String::from_utf8_lossy(&payload[2..]).into_owned() + } else { + String::new() + }; + (code, reason) + } else { + (1000, String::new()) + }; + let _ = self.send_close(code, &reason); + return Ok(None); + } + 0x9 => { + // Ping → respond with Pong + let _ = self.send_pong(&payload); + continue; + } + 0xA => { + // Pong → ignore + continue; + } + _ => { + let _ = self.send_close(1002, "Unsupported opcode"); + return Ok(None); + } + } + } + + // Try parsing JSON into ClientMessage + let message = match String::from_utf8(message_payload.clone()) { + Ok(text) => match serde_json::from_str(&text) { + Ok(msg) => WsMessage::Message(msg), + Err(_) => WsMessage::String(text), + }, + Err(_) => WsMessage::Binary(message_payload), + }; + + Ok(Some(message)) + } } impl Clone for Client { @@ -82,102 +332,3 @@ impl Hash for Client { self.0.peer_addr().unwrap().hash(state); } } - -impl Client { - /// Read a full WebSocket message, handling fragmentation (FIN) - pub fn read(&self) -> crate::Result>> { - let mut stream = &self.0; - let mut message_payload = Vec::new(); - let mut final_frame = false; - - while !final_frame { - let mut header = [0u8; 2]; - if stream.read_exact(&mut header).is_err() { - return Ok(None); // connection closed - } - - let fin = header[0] & 0x80 != 0; - let opcode = header[0] & 0x0F; - let masked = header[1] & 0x80 != 0; - let mut payload_len = (header[1] & 0x7F) as u64; - - // Extended payload lengths - if payload_len == 126 { - let mut ext_len = [0u8; 2]; - stream.read_exact(&mut ext_len)?; - payload_len = u16::from_be_bytes(ext_len) as u64; - } else if payload_len == 127 { - let mut ext_len = [0u8; 8]; - stream.read_exact(&mut ext_len)?; - payload_len = u64::from_be_bytes(ext_len); - } - - // Mask key (client → server) - let mut mask = [0u8; 4]; - if masked { - stream.read_exact(&mut mask)?; - } - - // Read payload - let mut payload = vec![0u8; payload_len as usize]; - stream.read_exact(&mut payload)?; - - if masked { - for i in 0..payload.len() { - payload[i] ^= mask[i % 4]; - } - } - - match opcode { - 0x0 | 0x1 | 0x2 => { - // Continuation / Text / Binary - message_payload.extend(payload); - } - 0x8 => return Ok(None), // Close - 0x9 => continue, // Ping → ignore - 0xA => continue, // Pong → ignore - _ => return Err(anyhow!("Unsupported WebSocket opcode: {}", opcode).into()), - } - - final_frame = fin; - } - - // Try parsing JSON into ClientMessage - let message = match String::from_utf8(message_payload.clone()) { - Ok(text) => match serde_json::from_str(&text) { - Ok(msg) => WsMessage::Message(msg), - Err(_) => WsMessage::String(text), - }, - Err(_) => WsMessage::Binary(message_payload), - }; - - Ok(Some(message)) - } - - /// Send a JSON-serializable object as a WebSocket text frame - pub fn send(&self, m: T) -> crate::Result<()> { - let payload = serde_json::to_string(&m)?; - let payload_bytes = payload.as_bytes(); - - let mut stream = self.0.try_clone()?; - let mut header = Vec::new(); - header.push(0x81); // FIN=1, opcode=0x1 (text) - - let len = payload_bytes.len(); - if len < 126 { - header.push(len as u8); - } else if len <= 65535 { - header.push(126); - header.extend_from_slice(&(len as u16).to_be_bytes()); - } else { - header.push(127); - header.extend_from_slice(&(len as u64).to_be_bytes()); - } - - stream.write_all(&header)?; - stream.write_all(payload_bytes)?; - stream.flush()?; - - Ok(()) - } -} From 50bc3a85caacb7fa5678e0e19a65a69062390c58 Mon Sep 17 00:00:00 2001 From: Leo dev Date: Mon, 15 Sep 2025 17:16:30 +0200 Subject: [PATCH 09/11] Better modularity --- src/lib.rs | 66 ++--------------------------------------- src/requests/message.rs | 63 +++++++++++++++++++++++++++++++++++++++ src/requests/mod.rs | 49 ++++++++++++++++++++++++++++++ src/types.rs | 13 ++++---- 4 files changed, 120 insertions(+), 71 deletions(-) create mode 100644 src/requests/message.rs create mode 100644 src/requests/mod.rs diff --git a/src/lib.rs b/src/lib.rs index 734d2fd..12aea54 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,16 +6,13 @@ use std::{ }; pub mod macros; +pub mod requests; pub mod types; pub mod utils; pub use anyhow::Result; -use crate::{ - types::{ClientMessage, WsMessage}, - utils::client::Client, - utils::plugin::DynPlugin, -}; +use crate::{utils::client::Client, utils::plugin::DynPlugin}; pub use once_cell; #[derive(serde::Serialize, serde::Deserialize)] @@ -123,67 +120,10 @@ impl Server { continue 'outer; } } - } - match req { - Some(WsMessage::Message(req)) => match req { - ClientMessage::SendMessage { - channel_id, - contents, - } => { - Self::LOGGER.info(format!("SendMessage to {channel_id}: {contents}")); - let msg = self.wrap_err( - &client, - self.db.messages_db.insert( - &channel_id, - "idk", - &contents, - chrono::Utc::now().timestamp(), - ), - )?; - - for c in self.clients.lock().unwrap().iter() { - self.wrap_err( - &c, - c.send(types::ServerMessage::MessageCreate(msg.clone())), - )?; - } - } - - ClientMessage::EditMessage { - channel_id, - message_id, - new_contents, - } => { - Self::LOGGER.info(format!( - "EditMessage {message_id} in {channel_id}: {new_contents}" - )); - } - - ClientMessage::DeleteMessage { - channel_id, - message_id, - } => { - Self::LOGGER.info(format!("DeleteMessage {message_id} in {channel_id}")); - } - }, - - Some(WsMessage::Binary(b)) => { - Self::LOGGER.info(format!("Binary message: {b:?}")); - } - - Some(WsMessage::String(s)) => { - Self::LOGGER.info(format!("String message: {s}")); - } - - None => { - self.clients.lock().unwrap().remove(&client); - break; - } + self.call_request(r, &client)?; } } - - Ok(()) } /// When there is a error it removes the client diff --git a/src/requests/message.rs b/src/requests/message.rs new file mode 100644 index 0000000..b2ef946 --- /dev/null +++ b/src/requests/message.rs @@ -0,0 +1,63 @@ +use std::sync::Arc; + +use crate::{Server, types, utils::client::Client}; + +crate::logger!(LOGGER "Message Manager"); + +pub fn send( + server: &Arc, + client: &Client, + channel_id: &str, + contents: &str, +) -> crate::Result<()> { + LOGGER.info(format!("SendMessage to {channel_id}: {contents}")); + + if contents.is_empty() { + server.wrap_err( + &client, + client.send(types::data::ResponseError::InvalidRequest(format!( + "Invalid message: empty message" + ))), + )?; + return Ok(()); + } + + let msg = server.wrap_err( + &client, + server.db.messages_db.insert( + &channel_id, + "idk", + &contents, + chrono::Utc::now().timestamp(), + ), + )?; + + for c in server.clients.lock().unwrap().iter() { + server.wrap_err(&c, c.send(types::ServerMessage::MessageCreate(msg.clone())))?; + } + + Ok(()) +} + +pub fn edit( + _server: &Arc, + _client: &Client, + channel_id: &str, + message_id: &str, + new_contents: &str, +) -> crate::Result<()> { + LOGGER.info(format!( + "EditMessage {message_id} in {channel_id}: {new_contents}" + )); + Ok(()) +} + +pub fn delete( + _server: &Arc, + _client: &Client, + channel_id: &str, + message_id: &str, +) -> crate::Result<()> { + LOGGER.info(format!("DeleteMessage {message_id} in {channel_id}")); + Ok(()) +} diff --git a/src/requests/mod.rs b/src/requests/mod.rs new file mode 100644 index 0000000..3835048 --- /dev/null +++ b/src/requests/mod.rs @@ -0,0 +1,49 @@ +pub mod message; + +use std::sync::Arc; + +use crate::{ + Server, + types::{ClientMessage, WsMessage}, + utils::client::Client, +}; + +impl Server { + pub fn call_request( + self: &Arc, + req: &WsMessage, + client: &Client, + ) -> crate::Result<()> { + match req { + WsMessage::Message(req) => match req { + ClientMessage::SendMessage { + channel_id, + contents, + } => { + message::send(self, client, channel_id, contents)?; + } + + ClientMessage::EditMessage { + channel_id, + message_id, + new_contents, + } => message::edit(self, client, channel_id, message_id, new_contents)?, + + ClientMessage::DeleteMessage { + channel_id, + message_id, + } => message::delete(self, client, channel_id, message_id)?, + }, + + WsMessage::Binary(b) => { + Self::LOGGER.info(format!("Binary message: {b:?}")); + } + + WsMessage::String(s) => { + Self::LOGGER.info(format!("String message: {s}")); + } + } + + Ok(()) + } +} diff --git a/src/types.rs b/src/types.rs index 2fd8636..bddc96c 100644 --- a/src/types.rs +++ b/src/types.rs @@ -33,9 +33,6 @@ pub enum ServerMessage { user_id: String, }, - /// Error responses - Error(), - TempMessage { message: String, }, @@ -101,11 +98,11 @@ pub mod data { } #[derive(Debug, Clone, Serialize, Deserialize)] - #[serde(tag = "error", rename_all = "snake_case")] + #[serde(tag = "error", content = "message", rename_all = "snake_case")] pub enum ResponseError { - InvalidRequest { message: String }, - Unauthorized { message: String }, - NotFound { message: String }, - InternalError { message: String }, + InvalidRequest(String), + Unauthorized(String), + NotFound(String), + InternalError(String), } } From 0f3e0ac68ab256d8d75a82d95ac6790eb7ce4fa9 Mon Sep 17 00:00:00 2001 From: Leo dev Date: Mon, 15 Sep 2025 17:43:12 +0200 Subject: [PATCH 10/11] Improvements overall --- src/requests/message.rs | 18 +++++------------ src/requests/mod.rs | 10 ++++------ src/types.rs | 8 +++----- src/utils/database.rs | 44 ++++++++++++++++++----------------------- 4 files changed, 31 insertions(+), 49 deletions(-) diff --git a/src/requests/message.rs b/src/requests/message.rs index b2ef946..efbf998 100644 --- a/src/requests/message.rs +++ b/src/requests/message.rs @@ -24,7 +24,7 @@ pub fn send( let msg = server.wrap_err( &client, - server.db.messages_db.insert( + server.db.insert_message( &channel_id, "idk", &contents, @@ -42,22 +42,14 @@ pub fn send( pub fn edit( _server: &Arc, _client: &Client, - channel_id: &str, - message_id: &str, + message_id: usize, new_contents: &str, ) -> crate::Result<()> { - LOGGER.info(format!( - "EditMessage {message_id} in {channel_id}: {new_contents}" - )); + LOGGER.info(format!("EditMessage {message_id}: {new_contents}")); Ok(()) } -pub fn delete( - _server: &Arc, - _client: &Client, - channel_id: &str, - message_id: &str, -) -> crate::Result<()> { - LOGGER.info(format!("DeleteMessage {message_id} in {channel_id}")); +pub fn delete(_server: &Arc, _client: &Client, message_id: usize) -> crate::Result<()> { + LOGGER.info(format!("DeleteMessage {message_id}")); Ok(()) } diff --git a/src/requests/mod.rs b/src/requests/mod.rs index 3835048..36364c1 100644 --- a/src/requests/mod.rs +++ b/src/requests/mod.rs @@ -24,15 +24,13 @@ impl Server { } ClientMessage::EditMessage { - channel_id, message_id, new_contents, - } => message::edit(self, client, channel_id, message_id, new_contents)?, + } => message::edit(self, client, *message_id, new_contents)?, - ClientMessage::DeleteMessage { - channel_id, - message_id, - } => message::delete(self, client, channel_id, message_id)?, + ClientMessage::DeleteMessage { message_id } => { + message::delete(self, client, *message_id)? + } }, WsMessage::Binary(b) => { diff --git a/src/types.rs b/src/types.rs index bddc96c..86f3f65 100644 --- a/src/types.rs +++ b/src/types.rs @@ -12,15 +12,13 @@ pub enum ClientMessage { /// Edit a message (if allowed) EditMessage { - channel_id: String, - message_id: String, + message_id: usize, new_contents: String, }, /// Delete a message (if allowed) DeleteMessage { - channel_id: String, - message_id: String, + message_id: usize, }, } @@ -46,7 +44,7 @@ pub enum ServerMessage { /// A message was deleted MessageDelete { channel_id: String, - message_id: String, + message_id: usize, }, /// Presence updates diff --git a/src/utils/database.rs b/src/utils/database.rs index b76bdac..3e3ca67 100644 --- a/src/utils/database.rs +++ b/src/utils/database.rs @@ -1,42 +1,33 @@ use crate::{ServerConfig, types::data::Message}; use rusqlite::{Connection, Result, params}; -pub struct Database { - pub messages_db: MessagesDb, -} +pub struct Database(pub Connection); +// General use case impl Database { - pub fn new(config: &ServerConfig) -> Option { + pub fn new(_config: &ServerConfig) -> Option { let conn = Connection::open("main.db").ok()?; - let messages_db = MessagesDb(conn); - messages_db.init(config)?; - Some(Self { messages_db }) - } -} -unsafe impl Send for Database {} -unsafe impl Sync for Database {} - -pub struct MessagesDb(pub Connection); - -impl MessagesDb { - pub fn init(&self, _config: &ServerConfig) -> Option { - self.0 - .execute( - "CREATE TABLE IF NOT EXISTS chat ( + conn.execute( + "CREATE TABLE IF NOT EXISTS chat ( id INTEGER PRIMARY KEY AUTOINCREMENT, channel_id TEXT NOT NULL, user_id TEXT NOT NULL, contents TEXT NOT NULL, timestamp INTEGER NOT NULL )", - [], - ) - .ok() - } + [], + ) + .ok()?; + Some(Database(conn)) + } +} + +// For chat messages +impl Database { /// Insert a message into the DB - pub fn insert( + pub fn insert_message( &self, channel_id: &str, user_id: &str, @@ -61,7 +52,7 @@ impl MessagesDb { } /// Get a message by its ID - pub fn get_by_id(&self, message_id: usize) -> Result> { + pub fn get_message_by_id(&self, message_id: usize) -> Result> { let mut stmt = self.0.prepare( "SELECT id, channel_id, user_id, contents, timestamp FROM chat @@ -91,3 +82,6 @@ impl MessagesDb { Ok(None) } } + +unsafe impl Send for Database {} +unsafe impl Sync for Database {} From 2b7abc505ba0b885607707c8e1533a66940fb672 Mon Sep 17 00:00:00 2001 From: Leo dev Date: Mon, 15 Sep 2025 18:23:27 +0200 Subject: [PATCH 11/11] Handshake --- client.js | 2 ++ src/lib.rs | 35 +++++++++++++++++++++++++++++++++++ src/requests/message.rs | 2 +- src/types.rs | 38 ++++++++++++++++++++++++++------------ src/utils/client.rs | 16 ++++++++++++++-- 5 files changed, 78 insertions(+), 15 deletions(-) diff --git a/client.js b/client.js index e1b8742..f423cb9 100644 --- a/client.js +++ b/client.js @@ -9,6 +9,8 @@ const ws = new WebSocket('ws://localhost:7080'); ws.onopen = () => { console.log('WebSocket connection established'); + console.log('Initializing handshake'); + ws.send(JSON.stringify({ version: '0.0.1', auth_token: '' })) }; ws.onmessage = (event) => { diff --git a/src/lib.rs b/src/lib.rs index 12aea54..1f56eed 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -17,6 +17,8 @@ pub use once_cell; #[derive(serde::Serialize, serde::Deserialize)] pub struct ServerConfig { + server_name: String, + server_id: String, port: u16, channels: Vec, } @@ -34,6 +36,8 @@ impl Default for ServerConfig { fn default() -> Self { Self { port: 7080, + server_name: format!("Server Name"), + server_id: format!("offline-server"), channels: Vec::new(), } } @@ -108,6 +112,37 @@ impl Server { // Initialize client let client = Client::new(stream)?; + // Initialize handshake + self.wrap_err( + &client, + client.send(types::handshake::ServerDetails { + name: self.config.server_name.clone(), + id: self.config.server_id.clone(), + version: format!("0.0.1"), + }), + )?; + + match self.wrap_err(&client, client.read_t::())? { + Some(types::WsMessage::Message(_)) => { + // Do auth stuff + self.wrap_err( + &client, + client.send(types::ServerMessage::Authenticated { + user_id: format!(""), + }), + )?; + } + Some(_) => { + self.wrap_err( + &client, + client.send(types::ResponseError::InvalidHandshake(format!( + "Invalid handshake" + ))), + )?; + } + None => {} + } + // Insert to the set of all connected clients self.clients.lock().unwrap().insert(client.clone()); diff --git a/src/requests/message.rs b/src/requests/message.rs index efbf998..e69616e 100644 --- a/src/requests/message.rs +++ b/src/requests/message.rs @@ -15,7 +15,7 @@ pub fn send( if contents.is_empty() { server.wrap_err( &client, - client.send(types::data::ResponseError::InvalidRequest(format!( + client.send(types::ResponseError::InvalidRequest(format!( "Invalid message: empty message" ))), )?; diff --git a/src/types.rs b/src/types.rs index 86f3f65..1b0d26f 100644 --- a/src/types.rs +++ b/src/types.rs @@ -17,9 +17,7 @@ pub enum ClientMessage { }, /// Delete a message (if allowed) - DeleteMessage { - message_id: usize, - }, + DeleteMessage { message_id: usize }, } /// Messages sent *from the server* to the client @@ -27,9 +25,7 @@ pub enum ClientMessage { #[serde(tag = "type", content = "params", rename_all = "snake_case")] pub enum ServerMessage { /// Successful authentication - Authenticated { - user_id: String, - }, + Authenticated { user_id: String }, TempMessage { message: String, @@ -60,6 +56,16 @@ pub enum ServerMessage { }, } +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "error", content = "message", rename_all = "snake_case")] +pub enum ResponseError { + InvalidRequest(String), + InvalidHandshake(String), + Unauthorized(String), + NotFound(String), + InternalError(String), +} + /// WebSocket wrapper #[derive(Debug, Clone)] pub enum WsMessage Deserialize<'de>> { @@ -94,13 +100,21 @@ pub mod data { Text, Voice, } +} + +pub mod handshake { + use serde::{Deserialize, Serialize}; #[derive(Debug, Clone, Serialize, Deserialize)] - #[serde(tag = "error", content = "message", rename_all = "snake_case")] - pub enum ResponseError { - InvalidRequest(String), - Unauthorized(String), - NotFound(String), - InternalError(String), + pub struct ServerDetails { + pub version: String, + pub name: String, + pub id: String, + } + + #[derive(Debug, Clone, Serialize, Deserialize)] + pub struct ClientDetails { + pub version: String, + pub auth_token: String, } } diff --git a/src/utils/client.rs b/src/utils/client.rs index 4f1331f..9137eab 100644 --- a/src/utils/client.rs +++ b/src/utils/client.rs @@ -6,7 +6,7 @@ use std::{ }; use anyhow::anyhow; -use serde::Serialize; +use serde::{Deserialize, Serialize}; use crate::types::{ClientMessage, WsMessage}; @@ -196,7 +196,9 @@ impl Client { /// - Ok(Some(WsMessage)) on an application message (text/binary) /// - Ok(None) if the connection should be closed (close received / read EOF) /// - Err on protocol or IO errors. - pub fn read(&self) -> crate::Result>> { + pub fn read_t Deserialize<'de>>( + &self, + ) -> crate::Result>> { let mut stream = self.0.try_clone()?; let mut message_payload = Vec::new(); @@ -311,6 +313,16 @@ impl Client { Ok(Some(message)) } + + /// Read a full WebSocket message, handling fragmentation and control frames. + /// + /// Returns: + /// - Ok(Some(WsMessage)) on an application message (text/binary) + /// - Ok(None) if the connection should be closed (close received / read EOF) + /// - Err on protocol or IO errors. + pub fn read(&self) -> crate::Result>> { + self.read_t() + } } impl Clone for Client {