From acf1ea91a37e5f7c63e4da071980a4c10372e1cf Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 19 Aug 2026 21:16:29 +0200 Subject: [PATCH 1/7] Refactor --- src/{ => data}/config.rs | 0 src/data/mod.rs | 1 + src/main.rs | 2 +- src/server.rs | 2 +- 4 files changed, 3 insertions(+), 2 deletions(-) rename src/{ => data}/config.rs (100%) create mode 100644 src/data/mod.rs diff --git a/src/config.rs b/src/data/config.rs similarity index 100% rename from src/config.rs rename to src/data/config.rs diff --git a/src/data/mod.rs b/src/data/mod.rs new file mode 100644 index 0000000..ef68c36 --- /dev/null +++ b/src/data/mod.rs @@ -0,0 +1 @@ +pub mod config; diff --git a/src/main.rs b/src/main.rs index 25696ca..bacef3d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,4 +1,4 @@ -pub mod config; +pub mod data; pub mod protocol; pub mod server; pub mod signature; diff --git a/src/server.rs b/src/server.rs index 41365e0..a9a2fb8 100644 --- a/src/server.rs +++ b/src/server.rs @@ -11,7 +11,7 @@ use ed25519_dalek::{SigningKey, VerifyingKey}; use tokio::sync::Mutex; use crate::{ - config::Config, + data::config::Config, protocol::{ClientMethod, read_loop, send_socket}, types::ClientMeta, }; From 8539b66cd067eb2403b6ca97150e0ce047ce0895 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 19 Aug 2026 21:23:37 +0200 Subject: [PATCH 2/7] Message store --- Cargo.lock | 102 ++++++++++++++++++++++++++++++++++++ Cargo.toml | 1 + src/data/messages.rs | 112 ++++++++++++++++++++++++++++++++++++++++ src/data/mod.rs | 1 + src/protocol/message.rs | 0 src/server.rs | 5 +- 6 files changed, 220 insertions(+), 1 deletion(-) create mode 100644 src/data/messages.rs create mode 100644 src/protocol/message.rs diff --git a/Cargo.lock b/Cargo.lock index 12ee0c4..62c05f4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,6 +2,18 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "ahash" +version = "0.8.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" +dependencies = [ + "cfg-if", + "once_cell", + "version_check", + "zerocopy", +] + [[package]] name = "anyhow" version = "1.0.104" @@ -111,6 +123,16 @@ version = "1.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" +[[package]] +name = "cc" +version = "1.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "509591b7bcd67f4ef775afad7662703b4935daaa6ec0e5605cfb1090b32a2b6d" +dependencies = [ + "find-msvc-tools", + "shlex", +] + [[package]] name = "cfg-if" version = "1.0.4" @@ -229,18 +251,37 @@ dependencies = [ "bs58", "ed25519-dalek", "rand 0.8.7", + "rusqlite", "serde", "serde_json", "tokio", "tower-http", ] +[[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 = "fiat-crypto" version = "0.2.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d" +[[package]] +name = "find-msvc-tools" +version = "0.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d45db016d36b838f563236e9193d0ee6ce38f3f68b6c94e914b4929c96bbb890" + [[package]] name = "form_urlencoded" version = "1.2.2" @@ -323,6 +364,24 @@ dependencies = [ "wasip2", ] +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" +dependencies = [ + "ahash", +] + +[[package]] +name = "hashlink" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ba4ff7128dee98c7dc9794b6a411377e1404dba1c97deb8d1a55297bd25d8af" +dependencies = [ + "hashbrown", +] + [[package]] name = "http" version = "1.5.0" @@ -421,6 +480,17 @@ version = "0.2.189" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2" +[[package]] +name = "libsqlite3-sys" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c10584274047cb335c23d3e61bcef8e323adae7c5c8c760540f73610177fc3f" +dependencies = [ + "cc", + "pkg-config", + "vcpkg", +] + [[package]] name = "log" version = "0.4.33" @@ -494,6 +564,12 @@ dependencies = [ "spki", ] +[[package]] +name = "pkg-config" +version = "0.3.34" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f6b464fbc74e149a392436b17d523f769e057cb6877f6a5c4618bc6f11800548" + [[package]] name = "ppv-lite86" version = "0.2.21" @@ -586,6 +662,20 @@ dependencies = [ "getrandom 0.3.4", ] +[[package]] +name = "rusqlite" +version = "0.31.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b838eba278d213a8beaf485bd313fd580ca4505a00d5871caeb1457c55322cae" +dependencies = [ + "bitflags", + "fallible-iterator", + "fallible-streaming-iterator", + "hashlink", + "libsqlite3-sys", + "smallvec", +] + [[package]] name = "rustc_version" version = "0.4.1" @@ -695,6 +785,12 @@ dependencies = [ "digest", ] +[[package]] +name = "shlex" +version = "2.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8fadd59c855ef2080decdef8ff161eb6661b86933c9d82e5ba29dc602a55aba" + [[package]] name = "signature" version = "2.2.0" @@ -963,6 +1059,12 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[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" diff --git a/Cargo.toml b/Cargo.toml index ec9225b..4cd58c6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -13,3 +13,4 @@ ed25519-dalek = { version = "2", features = ["rand_core"] } rand = "0.8" bs58 = "0.5.1" tower-http = { version = "0.7.0", features = ["fs", "cors"] } +rusqlite = { version = "0.31", features = ["bundled"] } diff --git a/src/data/messages.rs b/src/data/messages.rs new file mode 100644 index 0000000..728f684 --- /dev/null +++ b/src/data/messages.rs @@ -0,0 +1,112 @@ +use anyhow::Result; +use rusqlite::{params, Connection}; +use std::collections::HashMap; +use std::path::PathBuf; +use std::sync::Mutex; + +pub struct MessageStore { + data_dir: PathBuf, + connections: Mutex>, +} + +#[derive(Debug, Clone)] +pub struct StoredMessage { + pub id: String, + pub author_pubkey: String, + pub content: String, + pub timestamp: i64, + pub signature: String, + pub channel_id: String, +} + +impl MessageStore { + pub fn new(data_dir: PathBuf) -> Result { + std::fs::create_dir_all(&data_dir)?; + Ok(Self { + data_dir, + connections: Mutex::new(HashMap::new()), + }) + } + + /// Opens (or reuses an already-open) connection for a channel, + /// creating the schema if this is the first time. + fn with_channel( + &self, + channel_id: &str, + f: impl FnOnce(&Connection) -> Result, + ) -> Result { + let mut conns = self.connections.lock().unwrap(); + + if !conns.contains_key(channel_id) { + let path = self.data_dir.join(format!("{channel_id}.db")); + let conn = Connection::open(&path)?; + + conn.execute( + "CREATE TABLE IF NOT EXISTS messages ( + id TEXT PRIMARY KEY, + author_pubkey TEXT NOT NULL, + content TEXT NOT NULL, + timestamp INTEGER NOT NULL, + signature TEXT NOT NULL, + channel_id TEXT NOT NULL + )", + [], + )?; + + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_messages_timestamp ON messages(timestamp)", + [], + )?; + + conns.insert(channel_id.to_string(), conn); + } + + let conn = conns.get(channel_id).unwrap(); + f(conn) + } + + pub fn insert_message(&self, channel_id: &str, msg: &StoredMessage) -> Result<()> { + self.with_channel(channel_id, |conn| { + conn.execute( + "INSERT INTO messages (id, author_pubkey, content, timestamp, signature, channel_id) + VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + params![ + msg.id, + msg.author_pubkey, + msg.content, + msg.timestamp, + msg.signature, + msg.channel_id, + ], + )?; + Ok(()) + }) + } + + /// Fetches the most recent `limit` messages, oldest-first (ready to render top-to-bottom). + pub fn get_recent_messages(&self, channel_id: &str, limit: u32) -> Result> { + self.with_channel(channel_id, |conn| { + let mut stmt = conn.prepare( + "SELECT id, author_pubkey, content, timestamp, signature, channel_id + FROM messages + ORDER BY timestamp DESC + LIMIT ?1", + )?; + + let rows = stmt.query_map(params![limit], |row| { + Ok(StoredMessage { + id: row.get(0)?, + author_pubkey: row.get(1)?, + content: row.get(2)?, + timestamp: row.get(3)?, + signature: row.get(4)?, + channel_id: row.get(5)?, + }) + })?; + + let mut messages: Vec = rows.collect::>()?; + messages.reverse(); // DESC query, then flip to oldest-first for display + Ok(messages) + }) + } +} diff --git a/src/data/mod.rs b/src/data/mod.rs index ef68c36..ca69f9c 100644 --- a/src/data/mod.rs +++ b/src/data/mod.rs @@ -1 +1,2 @@ pub mod config; +pub mod messages; diff --git a/src/protocol/message.rs b/src/protocol/message.rs new file mode 100644 index 0000000..e69de29 diff --git a/src/server.rs b/src/server.rs index a9a2fb8..02b1be6 100644 --- a/src/server.rs +++ b/src/server.rs @@ -1,5 +1,6 @@ use std::{ collections::HashMap, + path::PathBuf, sync::{Arc, atomic::AtomicU16}, }; @@ -11,7 +12,7 @@ use ed25519_dalek::{SigningKey, VerifyingKey}; use tokio::sync::Mutex; use crate::{ - data::config::Config, + data::{config::Config, messages::MessageStore}, protocol::{ClientMethod, read_loop, send_socket}, types::ClientMeta, }; @@ -27,6 +28,7 @@ pub struct Server { pub key: SigningKey, pub config: Config, pub clients: Mutex>, + pub message_store: MessageStore, } impl Server { @@ -35,6 +37,7 @@ impl Server { key: crate::signature::get().await?, config: Config::get().await?, clients: Mutex::new(HashMap::new()), + message_store: MessageStore::new(PathBuf::from("item"))?, })) } } From 1dd75f489b1c417bc6accb8697976b173f1ca742 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 19 Aug 2026 21:50:23 +0200 Subject: [PATCH 3/7] Send message --- Cargo.lock | 99 ++++++++++++++++++++++++++++++++++++++++- Cargo.toml | 1 + src/data/messages.rs | 45 ++++++++++--------- src/protocol/message.rs | 47 +++++++++++++++++++ src/protocol/mod.rs | 19 +++++++- src/server.rs | 2 +- 6 files changed, 189 insertions(+), 24 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 62c05f4..25589f5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -117,6 +117,12 @@ dependencies = [ "tinyvec", ] +[[package]] +name = "bumpalo" +version = "3.20.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" + [[package]] name = "bytes" version = "1.12.1" @@ -256,6 +262,7 @@ dependencies = [ "serde_json", "tokio", "tower-http", + "uuid", ] [[package]] @@ -360,10 +367,21 @@ checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" dependencies = [ "cfg-if", "libc", - "r-efi", + "r-efi 5.3.0", "wasip2", ] +[[package]] +name = "getrandom" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099" +dependencies = [ + "cfg-if", + "libc", + "r-efi 6.0.0", +] + [[package]] name = "hashbrown" version = "0.14.5" @@ -474,6 +492,17 @@ version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" +[[package]] +name = "js-sys" +version = "0.3.104" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0e0c1080212aad755ea003d18543e8768dd432c48819efd73a7bf1e39b7a5a3a" +dependencies = [ + "cfg-if", + "futures-util", + "wasm-bindgen", +] + [[package]] name = "libc" version = "0.2.189" @@ -603,6 +632,12 @@ version = "5.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" +[[package]] +name = "r-efi" +version = "6.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" + [[package]] name = "rand" version = "0.8.7" @@ -685,6 +720,12 @@ dependencies = [ "semver", ] +[[package]] +name = "rustversion" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cf54715a573b99ac80df0bc206da022bcd442c974952c7b9720069370852e21f" + [[package]] name = "ryu" version = "1.0.23" @@ -1059,6 +1100,17 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[package]] +name = "uuid" +version = "1.24.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2cefc03fd367c0c6d4305de1b312cf00248c4114f4a0418ce6a6af769e3b0bd9" +dependencies = [ + "getrandom 0.4.3", + "js-sys", + "wasm-bindgen", +] + [[package]] name = "vcpkg" version = "0.2.15" @@ -1086,6 +1138,51 @@ dependencies = [ "wit-bindgen", ] +[[package]] +name = "wasm-bindgen" +version = "0.2.127" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1b70935747edd64d89de3efa29d73789b806c15798f8e7dca4d8ac356b50ce70" +dependencies = [ + "cfg-if", + "once_cell", + "rustversion", + "wasm-bindgen-macro", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-macro" +version = "0.2.127" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "77775f8f3f7217702089053b94958f8f54061a3f663417df76e19cbdcca29bc1" +dependencies = [ + "quote", + "wasm-bindgen-macro-support", +] + +[[package]] +name = "wasm-bindgen-macro-support" +version = "0.2.127" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e11d33f857dc2fb11b8bc75aee111aa9cbeb12cd9f25efd3d4c2a3dd4e235284" +dependencies = [ + "bumpalo", + "proc-macro2", + "quote", + "syn 2.0.119", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-shared" +version = "0.2.127" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7ef64dbcc55df09c7e5a46182d181c2cfa3e925f3da937ea764728b4bbb9dcbf" +dependencies = [ + "unicode-ident", +] + [[package]] name = "windows-link" version = "0.2.1" diff --git a/Cargo.toml b/Cargo.toml index 4cd58c6..219512b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,3 +14,4 @@ rand = "0.8" bs58 = "0.5.1" tower-http = { version = "0.7.0", features = ["fs", "cors"] } rusqlite = { version = "0.31", features = ["bundled"] } +uuid = { version = "1.24.1", features = ["v4"] } diff --git a/src/data/messages.rs b/src/data/messages.rs index 728f684..d2ecf22 100644 --- a/src/data/messages.rs +++ b/src/data/messages.rs @@ -1,5 +1,6 @@ use anyhow::Result; -use rusqlite::{params, Connection}; +use rusqlite::{Connection, params}; +use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::path::PathBuf; use std::sync::Mutex; @@ -9,14 +10,19 @@ pub struct MessageStore { connections: Mutex>, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MessageData { + pub content: String, + pub timestamp: u64, + pub signature: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] pub struct StoredMessage { pub id: String, - pub author_pubkey: String, - pub content: String, - pub timestamp: i64, - pub signature: String, - pub channel_id: String, + pub author: String, + #[serde(flatten)] + pub data: MessageData, } impl MessageStore { @@ -48,7 +54,6 @@ impl MessageStore { content TEXT NOT NULL, timestamp INTEGER NOT NULL, signature TEXT NOT NULL, - channel_id TEXT NOT NULL )", [], )?; @@ -68,15 +73,14 @@ impl MessageStore { pub fn insert_message(&self, channel_id: &str, msg: &StoredMessage) -> Result<()> { self.with_channel(channel_id, |conn| { conn.execute( - "INSERT INTO messages (id, author_pubkey, content, timestamp, signature, channel_id) + "INSERT INTO messages (id, author_pubkey, content, timestamp, signature) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", params![ msg.id, - msg.author_pubkey, - msg.content, - msg.timestamp, - msg.signature, - msg.channel_id, + msg.author, + msg.data.content, + msg.data.timestamp, + msg.data.signature, ], )?; Ok(()) @@ -87,7 +91,7 @@ impl MessageStore { pub fn get_recent_messages(&self, channel_id: &str, limit: u32) -> Result> { self.with_channel(channel_id, |conn| { let mut stmt = conn.prepare( - "SELECT id, author_pubkey, content, timestamp, signature, channel_id + "SELECT id, author_pubkey, content, timestamp, signature FROM messages ORDER BY timestamp DESC LIMIT ?1", @@ -96,11 +100,12 @@ impl MessageStore { let rows = stmt.query_map(params![limit], |row| { Ok(StoredMessage { id: row.get(0)?, - author_pubkey: row.get(1)?, - content: row.get(2)?, - timestamp: row.get(3)?, - signature: row.get(4)?, - channel_id: row.get(5)?, + author: row.get(1)?, + data: MessageData { + content: row.get(2)?, + timestamp: row.get(3)?, + signature: row.get(4)?, + }, }) })?; diff --git a/src/protocol/message.rs b/src/protocol/message.rs index e69de29..de7d09b 100644 --- a/src/protocol/message.rs +++ b/src/protocol/message.rs @@ -0,0 +1,47 @@ +use crate::data::messages::{MessageData, StoredMessage}; +use crate::server::Server; +use axum::extract::ws::WebSocket; +use ed25519_dalek::{Signature, Verifier, VerifyingKey}; +use std::sync::Arc; +use std::time::{SystemTime, UNIX_EPOCH}; +use tokio::sync::Mutex; + +pub async fn send_message( + server: &Arc, + verifying_key: VerifyingKey, + socket: &Arc>, + message: MessageData, + channel_id: String, +) -> anyhow::Result<()> { + let server_timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as u64; + + if server_timestamp.saturating_sub(message.timestamp) > 2000 { + anyhow::bail!( + "Message timestamp out of range ({server_timestamp} - {}) < 2000", + message.timestamp + ); + } + + let server_pubkey_string = crate::signature::to_string(&server.key.verifying_key()); + let signed_string = format!( + "{}@{}@{}", + message.timestamp, server_pubkey_string, message.content + ); + + let signature = crate::signature::from_string_sig(&message.signature) + .map_err(|_| anyhow::anyhow!("Invalid signature encoding"))?; + + verifying_key + .verify(signed_string.as_bytes(), &signature) + .map_err(|_| anyhow::anyhow!("Signature verification failed"))?; + + let stored = StoredMessage { + id: uuid::Uuid::new_v4().to_string(), + author: crate::signature::to_string(&verifying_key), + data: message, + }; + + server.message_store.insert_message(&channel_id, &stored)?; + + Ok(()) +} diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index a8ee7c5..d3eead3 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -1,12 +1,14 @@ use std::{borrow::Cow, sync::Arc}; use axum::extract::ws::{Message, Utf8Bytes, WebSocket}; +use ed25519_dalek::VerifyingKey; use serde::{Deserialize, Serialize}; use tokio::sync::Mutex; -use crate::types::ClientMeta; +use crate::{server::Server, types::ClientMeta}; pub mod initialize; +pub mod message; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(tag = "method")] @@ -35,6 +37,11 @@ pub enum ServerMethod { hostname: String, }, + SendMessage { + channel_id: String, + data: crate::data::messages::MessageData, + }, + Meta(ClientMeta), Error { @@ -42,7 +49,11 @@ pub enum ServerMethod { }, } -pub async fn read_loop(socket: &Arc>) -> anyhow::Result<()> { +pub async fn read_loop( + server: &Arc, + verifying_key: VerifyingKey, + socket: &Arc>, +) -> anyhow::Result<()> { while let Some(message) = read_socket(&mut *socket.lock().await).await? { match message { ServerMethod::Initialize { .. } => { @@ -61,6 +72,10 @@ pub async fn read_loop(socket: &Arc>) -> anyhow::Result<()> { ServerMethod::Error { error } => { eprintln!("Client error: {error}"); } + + ServerMethod::SendMessage { channel_id, data } => { + message::send_message(server, verifying_key, socket, data, channel_id).await?; + } } } diff --git a/src/server.rs b/src/server.rs index 02b1be6..3406455 100644 --- a/src/server.rs +++ b/src/server.rs @@ -69,7 +69,7 @@ impl Server { clients.connections.insert(conid, client.clone()); - if let Err(e) = read_loop(&client).await { + if let Err(e) = read_loop(&s, public_key, &client).await { eprintln!("Failed to handle client: {e}"); } else { println!("Client connection closed") From 4fe471243cf511ae728aae44f6f63398b555c85b Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 20 Aug 2026 11:24:39 +0200 Subject: [PATCH 4/7] Fix messages --- src/data/messages.rs | 4 ++-- src/protocol/message.rs | 4 ++-- src/server.rs | 2 +- src/types.rs | 7 ++++++- 4 files changed, 11 insertions(+), 6 deletions(-) diff --git a/src/data/messages.rs b/src/data/messages.rs index d2ecf22..6c613da 100644 --- a/src/data/messages.rs +++ b/src/data/messages.rs @@ -53,7 +53,7 @@ impl MessageStore { author_pubkey TEXT NOT NULL, content TEXT NOT NULL, timestamp INTEGER NOT NULL, - signature TEXT NOT NULL, + signature TEXT NOT NULL )", [], )?; @@ -74,7 +74,7 @@ impl MessageStore { self.with_channel(channel_id, |conn| { conn.execute( "INSERT INTO messages (id, author_pubkey, content, timestamp, signature) - VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + VALUES (?1, ?2, ?3, ?4, ?5)", params![ msg.id, msg.author, diff --git a/src/protocol/message.rs b/src/protocol/message.rs index de7d09b..b5a8458 100644 --- a/src/protocol/message.rs +++ b/src/protocol/message.rs @@ -1,7 +1,7 @@ use crate::data::messages::{MessageData, StoredMessage}; use crate::server::Server; use axum::extract::ws::WebSocket; -use ed25519_dalek::{Signature, Verifier, VerifyingKey}; +use ed25519_dalek::{Verifier, VerifyingKey}; use std::sync::Arc; use std::time::{SystemTime, UNIX_EPOCH}; use tokio::sync::Mutex; @@ -9,7 +9,7 @@ use tokio::sync::Mutex; pub async fn send_message( server: &Arc, verifying_key: VerifyingKey, - socket: &Arc>, + _socket: &Arc>, message: MessageData, channel_id: String, ) -> anyhow::Result<()> { diff --git a/src/server.rs b/src/server.rs index 3406455..6497b57 100644 --- a/src/server.rs +++ b/src/server.rs @@ -37,7 +37,7 @@ impl Server { key: crate::signature::get().await?, config: Config::get().await?, clients: Mutex::new(HashMap::new()), - message_store: MessageStore::new(PathBuf::from("item"))?, + message_store: MessageStore::new(PathBuf::from("messages"))?, })) } } diff --git a/src/types.rs b/src/types.rs index 98fde83..9a0ef5a 100644 --- a/src/types.rs +++ b/src/types.rs @@ -1,9 +1,14 @@ use serde::{Deserialize, Serialize}; #[derive(Debug, Clone, Serialize, Deserialize)] -pub struct ClientMeta {} +#[serde(rename_all = "camelCase")] +pub struct ClientMeta { + pub display_name: String, + pub avatar: Option, +} #[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] pub struct ServerMeta { pub name: String, pub description: String, From 0bce209dd139e8ac3e0f6315c4012872cc3fd74f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 20 Aug 2026 15:08:55 +0200 Subject: [PATCH 5/7] Get Messages --- .gitignore | 1 + src/data/messages.rs | 15 +++++++++++---- src/protocol/message.rs | 26 ++++++++++++++++++++++++++ src/protocol/mod.rs | 25 ++++++++++++++++++++++--- 4 files changed, 60 insertions(+), 7 deletions(-) diff --git a/.gitignore b/.gitignore index 57c39e9..4fc816a 100644 --- a/.gitignore +++ b/.gitignore @@ -27,3 +27,4 @@ target private.key config.json +messages diff --git a/src/data/messages.rs b/src/data/messages.rs index 6c613da..8434097 100644 --- a/src/data/messages.rs +++ b/src/data/messages.rs @@ -87,17 +87,24 @@ impl MessageStore { }) } - /// Fetches the most recent `limit` messages, oldest-first (ready to render top-to-bottom). - pub fn get_recent_messages(&self, channel_id: &str, limit: u32) -> Result> { + /// Fetches the most recent `limit` messages by the `offset`. + pub fn get_recent_messages( + &self, + channel_id: &str, + limit: u32, + chunk: u32, + ) -> Result> { self.with_channel(channel_id, |conn| { + let offset = chunk * limit; + let mut stmt = conn.prepare( "SELECT id, author_pubkey, content, timestamp, signature FROM messages ORDER BY timestamp DESC - LIMIT ?1", + LIMIT ?1 OFFSET ?2", )?; - let rows = stmt.query_map(params![limit], |row| { + let rows = stmt.query_map(params![limit, offset], |row| { Ok(StoredMessage { id: row.get(0)?, author: row.get(1)?, diff --git a/src/protocol/message.rs b/src/protocol/message.rs index b5a8458..0f3b76f 100644 --- a/src/protocol/message.rs +++ b/src/protocol/message.rs @@ -1,7 +1,9 @@ use crate::data::messages::{MessageData, StoredMessage}; +use crate::protocol::{ClientMethod, send_socket}; use crate::server::Server; use axum::extract::ws::WebSocket; use ed25519_dalek::{Verifier, VerifyingKey}; +use std::collections::HashMap; use std::sync::Arc; use std::time::{SystemTime, UNIX_EPOCH}; use tokio::sync::Mutex; @@ -45,3 +47,27 @@ pub async fn send_message( Ok(()) } + +pub async fn get_messages( + server: &Arc, + _verifying_key: VerifyingKey, + socket: &Arc>, + channel_id: String, + chunk: u32, +) -> anyhow::Result<()> { + const CHUNK_SIZE: u32 = 16; + + let messages = server + .message_store + .get_recent_messages(&channel_id, CHUNK_SIZE, chunk)?; + + send_socket( + &mut *socket.lock().await, + &ClientMethod::Messages { + messages: HashMap::from([(channel_id, messages)]), + }, + ) + .await?; + + Ok(()) +} diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index d3eead3..d15c575 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -1,11 +1,11 @@ -use std::{borrow::Cow, sync::Arc}; +use std::{borrow::Cow, collections::HashMap, sync::Arc}; use axum::extract::ws::{Message, Utf8Bytes, WebSocket}; use ed25519_dalek::VerifyingKey; use serde::{Deserialize, Serialize}; use tokio::sync::Mutex; -use crate::{server::Server, types::ClientMeta}; +use crate::{data::messages::StoredMessage, server::Server, types::ClientMeta}; pub mod initialize; pub mod message; @@ -21,6 +21,10 @@ pub enum ClientMethod { hostname: String, }, + Messages { + messages: HashMap>, + }, + Error { error: Cow<'static, str>, }, @@ -42,6 +46,11 @@ pub enum ServerMethod { data: crate::data::messages::MessageData, }, + GetMessages { + channel_id: String, + chunk: u32, + }, + Meta(ClientMeta), Error { @@ -54,7 +63,11 @@ pub async fn read_loop( verifying_key: VerifyingKey, socket: &Arc>, ) -> anyhow::Result<()> { - while let Some(message) = read_socket(&mut *socket.lock().await).await? { + let mut socket_lock = socket.lock().await; + + while let Some(message) = read_socket(&mut *socket_lock).await? { + drop(socket_lock); + match message { ServerMethod::Initialize { .. } => { send_socket( @@ -76,7 +89,13 @@ pub async fn read_loop( ServerMethod::SendMessage { channel_id, data } => { message::send_message(server, verifying_key, socket, data, channel_id).await?; } + + ServerMethod::GetMessages { channel_id, chunk } => { + message::get_messages(server, verifying_key, socket, channel_id, chunk).await?; + } } + + socket_lock = socket.lock().await; } Ok(()) From 943dc967033a1452a6930cbee301160bb7338900 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 20 Aug 2026 16:39:56 +0200 Subject: [PATCH 6/7] Message broadcast --- src/protocol/message.rs | 6 +++++ src/server.rs | 57 +++++++++++++++++++++++++++++++---------- 2 files changed, 49 insertions(+), 14 deletions(-) diff --git a/src/protocol/message.rs b/src/protocol/message.rs index 0f3b76f..b19542c 100644 --- a/src/protocol/message.rs +++ b/src/protocol/message.rs @@ -45,6 +45,12 @@ pub async fn send_message( server.message_store.insert_message(&channel_id, &stored)?; + server + .broadcast(&ClientMethod::Messages { + messages: HashMap::from([(channel_id, vec![stored])]), + }) + .await?; + Ok(()) } diff --git a/src/server.rs b/src/server.rs index 6497b57..a48b07c 100644 --- a/src/server.rs +++ b/src/server.rs @@ -9,7 +9,7 @@ use axum::{ response::Response, }; use ed25519_dalek::{SigningKey, VerifyingKey}; -use tokio::sync::Mutex; +use tokio::{sync::Mutex, task::JoinSet}; use crate::{ data::{config::Config, messages::MessageStore}, @@ -21,13 +21,13 @@ pub struct UserConnections { pub meta: ClientMeta, pub counter: AtomicU16, pub public_key: VerifyingKey, - pub connections: HashMap>>, + pub connections: Mutex>>>, } pub struct Server { pub key: SigningKey, pub config: Config, - pub clients: Mutex>, + pub clients: Mutex>>, pub message_store: MessageStore, } @@ -53,21 +53,27 @@ impl Server { let mut clients_meta = s.clients.lock().await; - let clients = - clients_meta - .entry(public_key) - .or_insert_with(|| UserConnections { + let clients = clients_meta + .entry(public_key) + .or_insert_with(|| { + Arc::new(UserConnections { meta, public_key: public_key, counter: AtomicU16::new(0), - connections: HashMap::new(), - }); + connections: Mutex::new(HashMap::new()), + }) + }) + .clone(); let conid = clients .counter .fetch_add(1, std::sync::atomic::Ordering::Relaxed); - clients.connections.insert(conid, client.clone()); + clients + .connections + .lock() + .await + .insert(conid, client.clone()); if let Err(e) = read_loop(&s, public_key, &client).await { eprintln!("Failed to handle client: {e}"); @@ -75,9 +81,11 @@ impl Server { println!("Client connection closed") } - clients.connections.remove(&conid); + let mut connections = clients.connections.lock().await; - if clients.connections.len() == 0 { + connections.remove(&conid); + + if connections.len() == 0 { clients_meta.remove(&public_key); } } @@ -88,11 +96,32 @@ impl Server { } }) } + + pub async fn broadcast(self: &Arc, message: &ClientMethod) -> anyhow::Result<()> { + let mut set = JoinSet::new(); + + for (_, client) in self.clients.lock().await.iter() { + let msg = message.clone(); + let client = client.clone(); + + set.spawn(async move { client.send(&msg).await }); + } + + // Await all spawned tasks to finish + while let Some(res) = set.join_next().await { + // handle task panic or errors if necessary + if let Ok(Err(e)) = res { + eprintln!("Failed to send to a client: {:?}", e); + } + } + + Ok(()) + } } impl UserConnections { pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> { - for (_, conn) in &self.connections { + for (_, conn) in self.connections.lock().await.iter() { send_socket(&mut *conn.lock().await, message).await?; } @@ -100,7 +129,7 @@ impl UserConnections { } pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result { - if let Some(conn) = self.connections.get(&id) { + if let Some(conn) = self.connections.lock().await.get(&id) { send_socket(&mut *conn.lock().await, message).await?; Ok(true) From 7c2cc0afe9effa8d162439694cd729c2c4794e05 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 20 Aug 2026 19:20:35 +0200 Subject: [PATCH 7/7] fix --- src/server.rs | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/server.rs b/src/server.rs index a48b07c..9440f7e 100644 --- a/src/server.rs +++ b/src/server.rs @@ -75,12 +75,16 @@ impl Server { .await .insert(conid, client.clone()); + drop(clients_meta); + if let Err(e) = read_loop(&s, public_key, &client).await { eprintln!("Failed to handle client: {e}"); } else { println!("Client connection closed") } + let mut clients_meta = s.clients.lock().await; + let mut connections = clients.connections.lock().await; connections.remove(&conid);