From 756018020e174558becaae7bf88b263947ca2863 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 26 Aug 2026 01:45:39 +0200 Subject: [PATCH 01/11] Voice chat channel --- src/data/config.rs | 34 ++++++++++++++++++++++++---------- src/types.rs | 3 ++- 2 files changed, 26 insertions(+), 11 deletions(-) diff --git a/src/data/config.rs b/src/data/config.rs index ea4141b..5a8693e 100644 --- a/src/data/config.rs +++ b/src/data/config.rs @@ -19,18 +19,32 @@ impl Config { name: "New Server".to_string(), description: String::new(), - channels: vec![Channel { - id: "text-channels".to_string(), - name: "Text Channels".to_string(), + channels: vec![ + Channel { + id: "text-channels".to_string(), + name: "Text Channels".to_string(), - data: ChannelKind::Category { - channels: vec![Channel { - id: "general".to_string(), - name: "General".to_string(), - data: ChannelKind::Text, - }], + data: ChannelKind::Category { + channels: vec![Channel { + id: "general".to_string(), + name: "General".to_string(), + data: ChannelKind::Text, + }], + }, }, - }], + Channel { + id: "voice-channels".to_string(), + name: "Voice Channels".to_string(), + + data: ChannelKind::Category { + channels: vec![Channel { + id: "vc-1".to_string(), + name: "VC 1".to_string(), + data: ChannelKind::Voice { max_users: 255 }, + }], + }, + }, + ], }, port: 3415, diff --git a/src/types.rs b/src/types.rs index 9a0ef5a..f35f431 100644 --- a/src/types.rs +++ b/src/types.rs @@ -19,8 +19,9 @@ pub struct ServerMeta { #[serde(tag = "kind")] #[serde(rename_all = "camelCase")] pub enum ChannelKind { - Text, Category { channels: Vec }, + Voice { max_users: u8 }, + Text, } #[derive(Debug, Clone, Serialize, Deserialize)] From f94ceea2e3dff6dfcfb017fc511a6f48450c7aa8 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 26 Aug 2026 11:52:39 +0200 Subject: [PATCH 02/11] Refactored crypto --- src/{signature.rs => crypto.rs} | 0 src/main.rs | 2 +- src/protocol/initialize.rs | 8 ++++---- src/protocol/message.rs | 14 +++++++------- src/server.rs | 4 ++-- 5 files changed, 14 insertions(+), 14 deletions(-) rename src/{signature.rs => crypto.rs} (100%) diff --git a/src/signature.rs b/src/crypto.rs similarity index 100% rename from src/signature.rs rename to src/crypto.rs diff --git a/src/main.rs b/src/main.rs index bacef3d..1dc0627 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,7 +1,7 @@ +pub mod crypto; pub mod data; pub mod protocol; pub mod server; -pub mod signature; pub mod types; use std::{ diff --git a/src/protocol/initialize.rs b/src/protocol/initialize.rs index 650cc48..7b0140e 100644 --- a/src/protocol/initialize.rs +++ b/src/protocol/initialize.rs @@ -66,7 +66,7 @@ impl UserConnections { return Err(anyhow::anyhow!("Client's hostname wasn't correct")); } - let Ok(public_key) = crate::signature::from_string(&public_key_string) else { + let Ok(public_key) = crate::crypto::from_string(&public_key_string) else { send_socket( &mut socket, &ClientMethod::Error { @@ -81,7 +81,7 @@ impl UserConnections { if public_key .verify_strict( format!("{timestamp}@{hostname}").as_bytes(), - &crate::signature::from_string_sig(&signature)?, + &crate::crypto::from_string_sig(&signature)?, ) .is_err() { @@ -100,8 +100,8 @@ impl UserConnections { send_socket( &mut socket, &ClientMethod::Initialized { - public_key: crate::signature::to_string(&server.key.verifying_key()), - signature: crate::signature::to_string_sig(&server.key.sign( + public_key: crate::crypto::to_string(&server.key.verifying_key()), + signature: crate::crypto::to_string_sig(&server.key.sign( format!("{server_timestamp}@{hostname}@{public_key_string}").as_bytes(), )), diff --git a/src/protocol/message.rs b/src/protocol/message.rs index 3c4b2c8..ec96e43 100644 --- a/src/protocol/message.rs +++ b/src/protocol/message.rs @@ -24,13 +24,13 @@ pub async fn send_message( ); } - let server_pubkey_string = crate::signature::to_string(&server.key.verifying_key()); + let server_pubkey_string = crate::crypto::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) + let signature = crate::crypto::from_string_sig(&message.signature) .map_err(|_| anyhow::anyhow!("Invalid signature encoding"))?; verifying_key @@ -39,7 +39,7 @@ pub async fn send_message( let stored = StoredMessage { id: uuid::Uuid::new_v4().to_string(), - author: crate::signature::to_string(&verifying_key), + author: crate::crypto::to_string(&verifying_key), is_edited: false, data: message, }; @@ -92,18 +92,18 @@ pub async fn edit_message( .get_message(&channel_id, &message_id)? .ok_or_else(|| anyhow::anyhow!("Message not found"))?; - let author_pubkey = crate::signature::to_string(&verifying_key); + let author_pubkey = crate::crypto::to_string(&verifying_key); if existing.author != author_pubkey { anyhow::bail!("Not authorized to edit this message"); } - let server_pubkey_string = crate::signature::to_string(&server.key.verifying_key()); + let server_pubkey_string = crate::crypto::to_string(&server.key.verifying_key()); let signed_string = format!( "{}@{}@{}", existing.data.timestamp, server_pubkey_string, new_content ); - let signature = crate::signature::from_string_sig(&new_signature) + let signature = crate::crypto::from_string_sig(&new_signature) .map_err(|_| anyhow::anyhow!("Invalid signature encoding"))?; verifying_key @@ -146,7 +146,7 @@ pub async fn delete_message( .get_message(&channel_id, &message_id)? .ok_or_else(|| anyhow::anyhow!("Message not found"))?; - let author_pubkey = crate::signature::to_string(&verifying_key); + let author_pubkey = crate::crypto::to_string(&verifying_key); if existing.author != author_pubkey { anyhow::bail!("Not authorized to delete this message"); } diff --git a/src/server.rs b/src/server.rs index c6ca1f3..2f20ec4 100644 --- a/src/server.rs +++ b/src/server.rs @@ -35,7 +35,7 @@ pub struct Server { impl Server { pub async fn new() -> anyhow::Result> { Ok(Arc::new(Self { - key: crate::signature::get().await?, + key: crate::crypto::get().await?, config: Config::get().await?, clients: Mutex::new(HashMap::new()), message_store: MessageStore::new(PathBuf::from("messages"))?, @@ -53,7 +53,7 @@ impl Server { Ok((client, public_key, meta)) => { if let Err(e) = s .user_store - .upsert_user(&crate::signature::to_string(&public_key), &meta) + .upsert_user(&crate::crypto::to_string(&public_key), &meta) .await { eprintln!("Failed to upsert client: {e}"); From 5486dfce6884f12b4f16ed011876407e9462bdb7 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 26 Aug 2026 12:37:11 +0200 Subject: [PATCH 03/11] JoinVoice events --- src/protocol/mod.rs | 25 +++++++++++++++++++++++++ src/server.rs | 2 ++ 2 files changed, 27 insertions(+) diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index d526c29..a6e2504 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -43,6 +43,11 @@ pub enum ClientMethod { channel_id: String, message_id: String, }, + + JoinVoice { + channel_id: String, + pin: u64, + }, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -87,6 +92,10 @@ pub enum ServerMethod { message_id: String, channel_id: String, }, + + JoinVoice { + channel_id: String, + }, } pub async fn read_loop( @@ -152,6 +161,22 @@ pub async fn read_loop( ServerMethod::GetUsers { pubkeys } => { user::get_users(server, verifying_key, socket, pubkeys).await?; } + + ServerMethod::JoinVoice { channel_id } => { + let pin = rand::random(); + + server + .voice_pins + .lock() + .await + .insert(pin, (verifying_key, channel_id.clone())); + + send_socket( + &mut *socket.lock().await, + &ClientMethod::JoinVoice { channel_id, pin }, + ) + .await?; + } } socket_lock = socket.lock().await; diff --git a/src/server.rs b/src/server.rs index 2f20ec4..2dce567 100644 --- a/src/server.rs +++ b/src/server.rs @@ -28,6 +28,7 @@ pub struct Server { pub key: SigningKey, pub config: Config, pub clients: Mutex>>, + pub voice_pins: Mutex>, pub message_store: MessageStore, pub user_store: UserMetaStore, } @@ -38,6 +39,7 @@ impl Server { key: crate::crypto::get().await?, config: Config::get().await?, clients: Mutex::new(HashMap::new()), + voice_pins: Mutex::new(HashMap::new()), message_store: MessageStore::new(PathBuf::from("messages"))?, user_store: UserMetaStore::new(PathBuf::from("users.db"))?, })) From ccbe3c10f04d16ef22ed9fdf1a0e0a1c9b47ccdf Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 26 Aug 2026 13:19:34 +0200 Subject: [PATCH 04/11] Testing voice udp server --- src/main.rs | 5 ++ src/server.rs | 11 ++++- src/vc_server.rs | 118 +++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 133 insertions(+), 1 deletion(-) create mode 100644 src/vc_server.rs diff --git a/src/main.rs b/src/main.rs index 1dc0627..68a08b9 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3,6 +3,7 @@ pub mod data; pub mod protocol; pub mod server; pub mod types; +pub mod vc_server; use std::{ net::{IpAddr, Ipv4Addr, SocketAddr}, @@ -26,6 +27,8 @@ use crate::server::Server; async fn main() -> anyhow::Result<()> { let server = Server::new().await?; + let udp_server = tokio::spawn(server.clone().start_udp_server()); + let cors = CorsLayer::new() .allow_origin(Any) .allow_methods(Any) @@ -46,6 +49,8 @@ async fn main() -> anyhow::Result<()> { axum::serve(listener, app).await?; + udp_server.abort(); + Ok(()) } diff --git a/src/server.rs b/src/server.rs index 2dce567..0584278 100644 --- a/src/server.rs +++ b/src/server.rs @@ -1,5 +1,6 @@ use std::{ collections::HashMap, + net::SocketAddr, path::PathBuf, sync::{Arc, atomic::AtomicU16}, }; @@ -9,7 +10,11 @@ use axum::{ response::Response, }; use ed25519_dalek::{SigningKey, VerifyingKey}; -use tokio::{sync::Mutex, task::JoinSet}; +use tokio::{ + net::UdpSocket, + sync::{Mutex, OnceCell}, + task::JoinSet, +}; use crate::{ data::{config::Config, messages::MessageStore, users::UserMetaStore}, @@ -22,6 +27,7 @@ pub struct UserConnections { pub counter: AtomicU16, pub public_key: VerifyingKey, pub connections: Mutex>>>, + pub voice: Mutex>, } pub struct Server { @@ -31,6 +37,7 @@ pub struct Server { pub voice_pins: Mutex>, pub message_store: MessageStore, pub user_store: UserMetaStore, + pub voice_socket: OnceCell, } impl Server { @@ -42,6 +49,7 @@ impl Server { voice_pins: Mutex::new(HashMap::new()), message_store: MessageStore::new(PathBuf::from("messages"))?, user_store: UserMetaStore::new(PathBuf::from("users.db"))?, + voice_socket: OnceCell::new(), })) } } @@ -73,6 +81,7 @@ impl Server { public_key: public_key, counter: AtomicU16::new(0), connections: Mutex::new(HashMap::new()), + voice: Mutex::new(None), }) }) .clone(); diff --git a/src/vc_server.rs b/src/vc_server.rs new file mode 100644 index 0000000..b8ad59c --- /dev/null +++ b/src/vc_server.rs @@ -0,0 +1,118 @@ +use std::{net::SocketAddr, sync::Arc}; + +use anyhow::Context; +use ed25519_dalek::VerifyingKey; +use tokio::net::UdpSocket; + +use crate::server::Server; + +impl Server { + pub async fn start_udp_server(self: Arc) -> anyhow::Result<()> { + let socket = UdpSocket::bind(("0.0.0.0", self.config.port)).await?; + + self.voice_socket + .set(socket) + .map_err(|_| anyhow::anyhow!("UDP server already started"))?; + + let mut buf = [0u8; 4096]; + + println!("UDP server started"); + + loop { + let (len, addr) = self.get_voice_socket()?.recv_from(&mut buf).await?; + + if len < 8 { + continue; // too short to even contain a pincode, drop silently + } + + let pin_bytes: [u8; 8] = buf[..8].try_into().unwrap(); + let pin = u64::from_be_bytes(pin_bytes); + let payload = &buf[8..len]; + + // Resolve the sender's identity + channel for this pin. + // First packet for a pin consumes it (single-use) and binds the address. + let sender = { + let mut pins = self.voice_pins.lock().await; + + if let Some((pubkey, channel_id)) = pins.remove(&pin) { + Some((pubkey, channel_id)) + } else { + None + } + }; + + let (sender_pubkey, channel_id) = match sender { + Some(v) => v, + + // Not a first-time pin — check if this addr is already a known + // voice participant, so we know who's speaking and where to relay. + None => match self.find_voice_sender(&addr).await { + Some(v) => v, + None => continue, // unknown pin, unknown addr — drop + }, + }; + + let clients = self.clients.lock().await; + let Some(user) = clients.get(&sender_pubkey).cloned() else { + continue; // pin referenced a user that's since disconnected + }; + drop(clients); + + // Record/refresh this user's known voice address + channel. + *user.voice.lock().await = Some((addr, channel_id.clone())); + + self.relay_voice(&sender_pubkey, &channel_id, payload).await; + } + } + + /// Looks up which known voice participant a UDP address belongs to, + /// for packets arriving after the initial pin-bearing packet. + async fn find_voice_sender(&self, addr: &SocketAddr) -> Option<(VerifyingKey, String)> { + let clients = self.clients.lock().await; + + for (pubkey, user) in clients.iter() { + if let Some((user_addr, channel)) = &*user.voice.lock().await { + if *user_addr == *addr { + return Some((*pubkey, channel.clone())); + } + } + } + + None + } + + /// Sends `payload` to every other voice participant currently in `channel_id`. + async fn relay_voice(&self, sender: &VerifyingKey, channel_id: &str, payload: &[u8]) { + let clients = self.clients.lock().await; + + for (pubkey, user) in clients.iter() { + // if pubkey == sender { + // continue; + // } + + let Some((addr, channel)) = &*user.voice.lock().await else { + continue; + }; + + if channel_id != channel { + continue; + } + + let _ = self.udp_send_to(addr, payload).await; + } + } + + pub fn get_voice_socket(&self) -> anyhow::Result<&UdpSocket> { + self.voice_socket + .get() + .context("Failed to get voice socket") + } + + pub async fn udp_send_to(&self, addr: &SocketAddr, payload: &[u8]) -> anyhow::Result<()> { + let socket = self.get_voice_socket()?; + + socket.send_to(payload, addr).await?; + + Ok(()) + } +} From 4b7ef4e15bae6894f8f14e8df3f4e01bdf3bf3c9 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 26 Aug 2026 14:31:14 +0200 Subject: [PATCH 05/11] Fixing UDP protocol --- src/vc_server.rs | 47 +++++++++++++++++++++-------------------------- 1 file changed, 21 insertions(+), 26 deletions(-) diff --git a/src/vc_server.rs b/src/vc_server.rs index b8ad59c..e62a9a9 100644 --- a/src/vc_server.rs +++ b/src/vc_server.rs @@ -25,32 +25,27 @@ impl Server { continue; // too short to even contain a pincode, drop silently } - let pin_bytes: [u8; 8] = buf[..8].try_into().unwrap(); - let pin = u64::from_be_bytes(pin_bytes); - let payload = &buf[8..len]; + let pin_bytes: [u8; 8] = buf[..8].try_into().unwrap(); + let pin = u64::from_be_bytes(pin_bytes); - // Resolve the sender's identity + channel for this pin. - // First packet for a pin consumes it (single-use) and binds the address. - let sender = { - let mut pins = self.voice_pins.lock().await; + // First packet carries the pin for authentication; subsequent packets + // are raw audio identified by UDP address alone. + let (sender_pubkey, channel_id, payload) = { + let mut pins = self.voice_pins.lock().await; - if let Some((pubkey, channel_id)) = pins.remove(&pin) { - Some((pubkey, channel_id)) - } else { - None + if let Some((pubkey, channel_id)) = pins.remove(&pin) { + // First packet: pin consumed, strip 8-byte prefix + (pubkey, channel_id, &buf[8..len]) + } else { + drop(pins); + + // Not a first-time pin — match by address + match self.find_voice_sender(&addr).await { + Some((pubkey, channel_id)) => (pubkey, channel_id, &buf[..]), + None => continue, } - }; - - let (sender_pubkey, channel_id) = match sender { - Some(v) => v, - - // Not a first-time pin — check if this addr is already a known - // voice participant, so we know who's speaking and where to relay. - None => match self.find_voice_sender(&addr).await { - Some(v) => v, - None => continue, // unknown pin, unknown addr — drop - }, - }; + } + }; let clients = self.clients.lock().await; let Some(user) = clients.get(&sender_pubkey).cloned() else { @@ -81,11 +76,11 @@ impl Server { None } - /// Sends `payload` to every other voice participant currently in `channel_id`. - async fn relay_voice(&self, sender: &VerifyingKey, channel_id: &str, payload: &[u8]) { + /// Sends `payload` to every voice participant currently in `channel_id`. + async fn relay_voice(&self, _sender: &VerifyingKey, channel_id: &str, payload: &[u8]) { let clients = self.clients.lock().await; - for (pubkey, user) in clients.iter() { + for (_pubkey, user) in clients.iter() { // if pubkey == sender { // continue; // } From 2bee3c1d32bc3735acd5c8b0e9807d3efb6ab92a Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 26 Aug 2026 14:55:53 +0200 Subject: [PATCH 06/11] Fix auth --- src/protocol/mod.rs | 2 +- src/vc_server.rs | 59 ++++++++++++++++++++++++++------------------- 2 files changed, 35 insertions(+), 26 deletions(-) diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index a6e2504..784573b 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -163,7 +163,7 @@ pub async fn read_loop( } ServerMethod::JoinVoice { channel_id } => { - let pin = rand::random(); + let pin = rand::random::() % (1 << 53); server .voice_pins diff --git a/src/vc_server.rs b/src/vc_server.rs index e62a9a9..fc04d8b 100644 --- a/src/vc_server.rs +++ b/src/vc_server.rs @@ -16,44 +16,50 @@ impl Server { let mut buf = [0u8; 4096]; - println!("UDP server started"); + eprintln!("[vc] UDP server listening on port {}", self.config.port); loop { let (len, addr) = self.get_voice_socket()?.recv_from(&mut buf).await?; + eprintln!("[vc] UDP packet received: {len} bytes from {addr}"); + if len < 8 { - continue; // too short to even contain a pincode, drop silently + eprintln!("[vc] dropping packet too short for pin"); + continue; } - let pin_bytes: [u8; 8] = buf[..8].try_into().unwrap(); - let pin = u64::from_be_bytes(pin_bytes); + let pin_bytes: [u8; 8] = buf[..8].try_into().unwrap(); + let pin = u64::from_be_bytes(pin_bytes); - // First packet carries the pin for authentication; subsequent packets - // are raw audio identified by UDP address alone. - let (sender_pubkey, channel_id, payload) = { - let mut pins = self.voice_pins.lock().await; + let (sender_pubkey, channel_id, payload) = { + let mut pins = self.voice_pins.lock().await; - if let Some((pubkey, channel_id)) = pins.remove(&pin) { - // First packet: pin consumed, strip 8-byte prefix - (pubkey, channel_id, &buf[8..len]) - } else { - drop(pins); + if let Some((pubkey, channel_id)) = pins.remove(&pin) { + eprintln!("[vc] pin {pin} authenticated for channel {channel_id}"); + (pubkey, channel_id, &buf[8..len]) + } else { + drop(pins); - // Not a first-time pin — match by address - match self.find_voice_sender(&addr).await { - Some((pubkey, channel_id)) => (pubkey, channel_id, &buf[..]), - None => continue, + match self.find_voice_sender(&addr).await { + Some((pubkey, channel_id)) => { + eprintln!("[vc] known sender in channel {channel_id}"); + (pubkey, channel_id, &buf[..]) + } + None => { + eprintln!("[vc] unknown pin {pin} and unknown addr {addr}, dropping"); + continue; + } + } } - } - }; + }; let clients = self.clients.lock().await; let Some(user) = clients.get(&sender_pubkey).cloned() else { - continue; // pin referenced a user that's since disconnected + eprintln!("[vc] sender pubkey not found in clients, dropping"); + continue; }; drop(clients); - // Record/refresh this user's known voice address + channel. *user.voice.lock().await = Some((addr, channel_id.clone())); self.relay_voice(&sender_pubkey, &channel_id, payload).await; @@ -80,11 +86,8 @@ impl Server { async fn relay_voice(&self, _sender: &VerifyingKey, channel_id: &str, payload: &[u8]) { let clients = self.clients.lock().await; + let mut sent = 0; for (_pubkey, user) in clients.iter() { - // if pubkey == sender { - // continue; - // } - let Some((addr, channel)) = &*user.voice.lock().await else { continue; }; @@ -94,7 +97,13 @@ impl Server { } let _ = self.udp_send_to(addr, payload).await; + sent += 1; } + + eprintln!( + "[vc] relayed {} bytes to {sent} users in {channel_id}", + payload.len() + ); } pub fn get_voice_socket(&self) -> anyhow::Result<&UdpSocket> { From 8887c96f0519bdc9ec52aad81777abed9c16c0df Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 27 Aug 2026 01:47:53 +0200 Subject: [PATCH 07/11] User speaking --- src/protocol/mod.rs | 4 +++ src/server.rs | 9 +++++- src/vc_server.rs | 67 +++++++++++++++++++++++++++++---------------- 3 files changed, 55 insertions(+), 25 deletions(-) diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index 784573b..50688af 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -48,6 +48,10 @@ pub enum ClientMethod { channel_id: String, pin: u64, }, + + Speaking { + pubkey: String, + }, } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/src/server.rs b/src/server.rs index 0584278..c5bbc74 100644 --- a/src/server.rs +++ b/src/server.rs @@ -14,6 +14,7 @@ use tokio::{ net::UdpSocket, sync::{Mutex, OnceCell}, task::JoinSet, + time::Instant, }; use crate::{ @@ -22,12 +23,18 @@ use crate::{ types::ClientMeta, }; +pub struct VoiceConnection { + pub addr: SocketAddr, + pub channel_id: String, + pub last_speaking_sent: Instant, +} + pub struct UserConnections { pub meta: ClientMeta, pub counter: AtomicU16, pub public_key: VerifyingKey, pub connections: Mutex>>>, - pub voice: Mutex>, + pub voice: Mutex>, } pub struct Server { diff --git a/src/vc_server.rs b/src/vc_server.rs index fc04d8b..25ed6c5 100644 --- a/src/vc_server.rs +++ b/src/vc_server.rs @@ -4,7 +4,12 @@ use anyhow::Context; use ed25519_dalek::VerifyingKey; use tokio::net::UdpSocket; -use crate::server::Server; +use crate::{ + protocol::{ClientMethod, send_socket}, + server::Server, +}; + +use tokio::time::Instant; impl Server { pub async fn start_udp_server(self: Arc) -> anyhow::Result<()> { @@ -21,8 +26,6 @@ impl Server { loop { let (len, addr) = self.get_voice_socket()?.recv_from(&mut buf).await?; - eprintln!("[vc] UDP packet received: {len} bytes from {addr}"); - if len < 8 { eprintln!("[vc] dropping packet too short for pin"); continue; @@ -35,18 +38,13 @@ impl Server { let mut pins = self.voice_pins.lock().await; if let Some((pubkey, channel_id)) = pins.remove(&pin) { - eprintln!("[vc] pin {pin} authenticated for channel {channel_id}"); (pubkey, channel_id, &buf[8..len]) } else { drop(pins); match self.find_voice_sender(&addr).await { - Some((pubkey, channel_id)) => { - eprintln!("[vc] known sender in channel {channel_id}"); - (pubkey, channel_id, &buf[..]) - } + Some((pubkey, channel_id)) => (pubkey, channel_id, &buf[..]), None => { - eprintln!("[vc] unknown pin {pin} and unknown addr {addr}, dropping"); continue; } } @@ -60,9 +58,14 @@ impl Server { }; drop(clients); - *user.voice.lock().await = Some((addr, channel_id.clone())); + *user.voice.lock().await = Some(crate::server::VoiceConnection { + addr, + channel_id: channel_id.clone(), + last_speaking_sent: Instant::now(), + }); - self.relay_voice(&sender_pubkey, &channel_id, payload).await; + self.relay_voice(&sender_pubkey, &channel_id, payload) + .await?; } } @@ -72,9 +75,9 @@ impl Server { let clients = self.clients.lock().await; for (pubkey, user) in clients.iter() { - if let Some((user_addr, channel)) = &*user.voice.lock().await { - if *user_addr == *addr { - return Some((*pubkey, channel.clone())); + if let Some(voice) = &*user.voice.lock().await { + if *addr == voice.addr { + return Some((*pubkey, voice.channel_id.clone())); } } } @@ -83,27 +86,43 @@ impl Server { } /// Sends `payload` to every voice participant currently in `channel_id`. - async fn relay_voice(&self, _sender: &VerifyingKey, channel_id: &str, payload: &[u8]) { + async fn relay_voice( + &self, + sender: &VerifyingKey, + channel_id: &str, + payload: &[u8], + ) -> anyhow::Result<()> { let clients = self.clients.lock().await; - let mut sent = 0; for (_pubkey, user) in clients.iter() { - let Some((addr, channel)) = &*user.voice.lock().await else { + let Some(voice) = &mut *user.voice.lock().await else { continue; }; - if channel_id != channel { + if channel_id != voice.channel_id { continue; } - let _ = self.udp_send_to(addr, payload).await; - sent += 1; + let now = Instant::now(); + + if now.duration_since(voice.last_speaking_sent).as_millis() >= 1000 { + for conn in user.connections.lock().await.values() { + let _ = send_socket( + &mut *conn.lock().await, + &ClientMethod::Speaking { + pubkey: crate::crypto::to_string(sender), + }, + ) + .await; + } + } + + voice.last_speaking_sent = now; + + let _ = self.udp_send_to(&voice.addr, payload).await; } - eprintln!( - "[vc] relayed {} bytes to {sent} users in {channel_id}", - payload.len() - ); + Ok(()) } pub fn get_voice_socket(&self) -> anyhow::Result<&UdpSocket> { From 02a444d6ede4c406835fcc9087ca35357011015c Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 27 Aug 2026 02:22:56 +0200 Subject: [PATCH 08/11] Fix idkx --- src/protocol/mod.rs | 39 ++++++++++++++++++++++++++++----------- src/vc_server.rs | 6 +++--- 2 files changed, 31 insertions(+), 14 deletions(-) diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index 50688af..7fbd942 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -49,6 +49,11 @@ pub enum ClientMethod { pin: u64, }, + UserJoinedVoice { + channel_id: String, + pubkey: String, + }, + Speaking { pubkey: String, }, @@ -167,19 +172,31 @@ pub async fn read_loop( } ServerMethod::JoinVoice { channel_id } => { - let pin = rand::random::() % (1 << 53); + { + let pin = rand::random::() % (1 << 53); + + server + .voice_pins + .lock() + .await + .insert(pin, (verifying_key, channel_id.clone())); + + send_socket( + &mut *socket.lock().await, + &ClientMethod::JoinVoice { + channel_id: channel_id.clone(), + pin, + }, + ) + .await?; + } server - .voice_pins - .lock() - .await - .insert(pin, (verifying_key, channel_id.clone())); - - send_socket( - &mut *socket.lock().await, - &ClientMethod::JoinVoice { channel_id, pin }, - ) - .await?; + .broadcast(&ClientMethod::UserJoinedVoice { + channel_id, + pubkey: crate::crypto::to_string(&verifying_key), + }) + .await?; } } diff --git a/src/vc_server.rs b/src/vc_server.rs index 25ed6c5..289c289 100644 --- a/src/vc_server.rs +++ b/src/vc_server.rs @@ -105,7 +105,7 @@ impl Server { let now = Instant::now(); - if now.duration_since(voice.last_speaking_sent).as_millis() >= 1000 { + if now.duration_since(voice.last_speaking_sent).as_millis() >= 600 { for conn in user.connections.lock().await.values() { let _ = send_socket( &mut *conn.lock().await, @@ -115,9 +115,9 @@ impl Server { ) .await; } - } - voice.last_speaking_sent = now; + voice.last_speaking_sent = now; + } let _ = self.udp_send_to(&voice.addr, payload).await; } From a4d675bda911275cc4b0f67560ee73f051c4ea2f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 27 Aug 2026 16:18:19 +0200 Subject: [PATCH 09/11] Better websocket --- Cargo.lock | 13 +++++++ Cargo.toml | 1 + src/main.rs | 1 + src/protocol/initialize.rs | 74 ++++++++++++++++---------------------- src/protocol/message.rs | 18 ++++------ src/protocol/mod.rs | 72 ++++++------------------------------- src/protocol/user.rs | 11 ++---- src/server.rs | 13 ++++--- src/vc_server.rs | 15 +++----- src/ws.rs | 65 +++++++++++++++++++++++++++++++++ 10 files changed, 141 insertions(+), 142 deletions(-) create mode 100644 src/ws.rs diff --git a/Cargo.lock b/Cargo.lock index 25589f5..28940fa 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -256,6 +256,7 @@ dependencies = [ "axum", "bs58", "ed25519-dalek", + "futures-util", "rand 0.8.7", "rusqlite", "serde", @@ -313,6 +314,17 @@ version = "0.3.34" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "92d699e522242e69e3003b94ecc1f960f3a5e015aa7c5d7486e65ad01dd94f5e" +[[package]] +name = "futures-macro" +version = "0.3.34" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9fb9654ba8355388abeb8dcb4fc62f511300867002afc858860463bdd9fe0c44" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.3", +] + [[package]] name = "futures-sink" version = "0.3.34" @@ -332,6 +344,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0d50a92467f8ba5dd6e3ee5d4bd04d73ab2e4e1c44474a0674821dfce14b79bc" dependencies = [ "futures-core", + "futures-macro", "futures-sink", "futures-task", "pin-project-lite", diff --git a/Cargo.toml b/Cargo.toml index 219512b..2d0a27a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -15,3 +15,4 @@ 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"] } +futures-util = "0.3.34" diff --git a/src/main.rs b/src/main.rs index 68a08b9..a4524bf 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,6 +4,7 @@ pub mod protocol; pub mod server; pub mod types; pub mod vc_server; +pub mod ws; use std::{ net::{IpAddr, Ipv4Addr, SocketAddr}, diff --git a/src/protocol/initialize.rs b/src/protocol/initialize.rs index 7b0140e..f5e4778 100644 --- a/src/protocol/initialize.rs +++ b/src/protocol/initialize.rs @@ -3,10 +3,9 @@ use std::{ time::{SystemTime, UNIX_EPOCH}, }; -use axum::extract::ws::WebSocket; use ed25519_dalek::{Signer, VerifyingKey}; -use crate::server::Server; +use crate::{server::Server, ws::EnclaveWebSocket}; use super::*; use crate::server::UserConnections; @@ -14,23 +13,21 @@ use crate::server::UserConnections; impl UserConnections { pub async fn initialize( server: &Arc, - mut socket: WebSocket, - ) -> anyhow::Result<(WebSocket, VerifyingKey, ClientMeta)> { + socket: Arc, + ) -> anyhow::Result<(Arc, VerifyingKey, ClientMeta)> { let Some(ServerMethod::Initialize { public_key: public_key_string, signature, timestamp, hostname, - }) = read_socket(&mut socket).await? + }) = socket.read().await? else { - send_socket( - &mut socket, - &ClientMethod::Error { + socket + .send(&ClientMethod::Error { error: Cow::Borrowed("Initialization required"), - }, - ) - .await?; + }) + .await?; return Err(anyhow::anyhow!( "Failed to initialize: Client sent the wrong method" @@ -40,23 +37,20 @@ impl UserConnections { let server_timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as u64; if server_timestamp.saturating_sub(timestamp) > 2000 { - send_socket( - &mut socket, - &ClientMethod::Error { + socket + .send(&ClientMethod::Error { error: Cow::Borrowed( "Timestamp doesn't match, make sure it's in secs and is (<= 2secs)", ), - }, - ) - .await?; + }) + .await?; return Err(anyhow::anyhow!("Client tampstamp wasn't correct")); } if hostname != server.config.public_hostname || !server.config.hostnames.contains(&hostname) { - send_socket( - &mut socket, + socket.send( &ClientMethod::Error { error: Cow::Owned(format!("Invalid Hostname, to avoid man-in-the-middle attacks, please use the correct hostname: {}", server.config.public_hostname)), }, @@ -67,13 +61,11 @@ impl UserConnections { } let Ok(public_key) = crate::crypto::from_string(&public_key_string) else { - send_socket( - &mut socket, - &ClientMethod::Error { + socket + .send(&ClientMethod::Error { error: Cow::Borrowed("Invalid public key"), - }, - ) - .await?; + }) + .await?; return Err(anyhow::anyhow!("Invalid public key")); }; @@ -85,21 +77,18 @@ impl UserConnections { ) .is_err() { - send_socket( - &mut socket, - &ClientMethod::Error { + socket + .send(&ClientMethod::Error { error: Cow::Borrowed("Invalid signature"), - }, - ) - .await?; + }) + .await?; return Err(anyhow::anyhow!("Invalid signature")); } { - send_socket( - &mut socket, - &ClientMethod::Initialized { + socket + .send(&ClientMethod::Initialized { public_key: crate::crypto::to_string(&server.key.verifying_key()), signature: crate::crypto::to_string_sig(&server.key.sign( format!("{server_timestamp}@{hostname}@{public_key_string}").as_bytes(), @@ -107,19 +96,16 @@ impl UserConnections { timestamp: server_timestamp, hostname, - }, - ) - .await?; + }) + .await?; } - let Some(ServerMethod::Meta(meta)) = read_socket(&mut socket).await? else { - send_socket( - &mut socket, - &ClientMethod::Error { + let Some(ServerMethod::Meta(meta)) = socket.read().await? else { + socket + .send(&ClientMethod::Error { error: Cow::Borrowed("Expected meta"), - }, - ) - .await?; + }) + .await?; return Err(anyhow::anyhow!( "Expected meta, client called another method" diff --git a/src/protocol/message.rs b/src/protocol/message.rs index ec96e43..91d7a1e 100644 --- a/src/protocol/message.rs +++ b/src/protocol/message.rs @@ -1,17 +1,15 @@ use crate::data::messages::{MessageData, StoredMessage}; -use crate::protocol::{ClientMethod, send_socket}; +use crate::protocol::ClientMethod; 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; pub async fn send_message( server: &Arc, verifying_key: VerifyingKey, - _socket: &Arc>, + _socket: &Arc, message: MessageData, channel_id: String, ) -> anyhow::Result<()> { @@ -58,7 +56,7 @@ pub async fn send_message( pub async fn get_messages( server: &Arc, _verifying_key: VerifyingKey, - socket: &Arc>, + socket: &Arc, channel_id: String, chunk: u32, ) -> anyhow::Result<()> { @@ -68,13 +66,11 @@ pub async fn get_messages( .message_store .get_recent_messages(&channel_id, CHUNK_SIZE, chunk)?; - send_socket( - &mut *socket.lock().await, - &ClientMethod::Messages { + socket + .send(&ClientMethod::Messages { messages: HashMap::from([(channel_id, messages)]), - }, - ) - .await?; + }) + .await?; Ok(()) } diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index 7fbd942..f8e1814 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -1,9 +1,7 @@ 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::{data::messages::StoredMessage, server::Server, types::ClientMeta}; @@ -110,22 +108,16 @@ pub enum ServerMethod { pub async fn read_loop( server: &Arc, verifying_key: VerifyingKey, - socket: &Arc>, + socket: &Arc, ) -> anyhow::Result<()> { - let mut socket_lock = socket.lock().await; - - while let Some(message) = read_socket(&mut *socket_lock).await? { - drop(socket_lock); - + while let Some(message) = socket.read().await? { match message { ServerMethod::Initialize { .. } => { - send_socket( - &mut *socket.lock().await, - &ClientMethod::Error { + socket + .send(&ClientMethod::Error { error: Cow::Borrowed("Already initialized"), - }, - ) - .await?; + }) + .await?; } #[allow(unused_variables)] @@ -181,14 +173,12 @@ pub async fn read_loop( .await .insert(pin, (verifying_key, channel_id.clone())); - send_socket( - &mut *socket.lock().await, - &ClientMethod::JoinVoice { + socket + .send(&ClientMethod::JoinVoice { channel_id: channel_id.clone(), pin, - }, - ) - .await?; + }) + .await?; } server @@ -199,49 +189,7 @@ pub async fn read_loop( .await?; } } - - socket_lock = socket.lock().await; } Ok(()) } - -pub async fn read_socket(socket: &mut WebSocket) -> anyhow::Result> { - match socket.recv().await.transpose()? { - Some(Message::Text(text)) => match serde_json::from_str(&text.to_string()) { - Ok(msg) => Ok(Some(msg)), - - Err(e) => { - send_socket( - socket, - &ClientMethod::Error { - error: Cow::Owned(format!("Unable to parse message: {e}")), - }, - ) - .await?; - - Ok(None) - } - }, - - Some(Message::Ping(v)) => { - socket.send(Message::Pong(v)).await?; - - Ok(None) - } - - Some(_) => Ok(None), - - None => Ok(None), - } -} - -pub async fn send_socket(socket: &mut WebSocket, message: &ClientMethod) -> anyhow::Result<()> { - socket - .send(Message::Text(Utf8Bytes::from(serde_json::to_string( - message, - )?))) - .await?; - - Ok(()) -} diff --git a/src/protocol/user.rs b/src/protocol/user.rs index 8bc0ac1..c56e76f 100644 --- a/src/protocol/user.rs +++ b/src/protocol/user.rs @@ -1,23 +1,18 @@ use std::sync::Arc; -use axum::extract::ws::WebSocket; use ed25519_dalek::VerifyingKey; -use tokio::sync::Mutex; -use crate::{ - protocol::{ClientMethod, send_socket}, - server::Server, -}; +use crate::{protocol::ClientMethod, server::Server}; pub async fn get_users( server: &Arc, _verifying_key: VerifyingKey, - socket: &Arc>, + socket: &Arc, pubkeys: Vec, ) -> anyhow::Result<()> { let users = server.user_store.get_users(&pubkeys).await?; - send_socket(&mut *socket.lock().await, &ClientMethod::Users { users }).await?; + socket.send(&ClientMethod::Users { users }).await?; Ok(()) } diff --git a/src/server.rs b/src/server.rs index c5bbc74..39895f4 100644 --- a/src/server.rs +++ b/src/server.rs @@ -19,8 +19,9 @@ use tokio::{ use crate::{ data::{config::Config, messages::MessageStore, users::UserMetaStore}, - protocol::{ClientMethod, read_loop, send_socket}, + protocol::{ClientMethod, read_loop}, types::ClientMeta, + ws::EnclaveWebSocket, }; pub struct VoiceConnection { @@ -33,7 +34,7 @@ pub struct UserConnections { pub meta: ClientMeta, pub counter: AtomicU16, pub public_key: VerifyingKey, - pub connections: Mutex>>>, + pub connections: Mutex>>, pub voice: Mutex>, } @@ -66,7 +67,7 @@ impl Server { let s = self.clone(); ws.on_upgrade(move |socket: WebSocket| async move { - match UserConnections::initialize(&s, socket).await { + match UserConnections::initialize(&s, Arc::new(EnclaveWebSocket::new(socket))).await { Ok((client, public_key, meta)) => { if let Err(e) = s .user_store @@ -76,8 +77,6 @@ impl Server { eprintln!("Failed to upsert client: {e}"); } - let client = Arc::new(Mutex::new(client)); - let mut clients_meta = s.clients.lock().await; let clients = clients_meta @@ -154,7 +153,7 @@ impl Server { impl UserConnections { pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> { for (_, conn) in self.connections.lock().await.iter() { - send_socket(&mut *conn.lock().await, message).await?; + conn.send(message).await?; } Ok(()) @@ -162,7 +161,7 @@ impl UserConnections { pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result { if let Some(conn) = self.connections.lock().await.get(&id) { - send_socket(&mut *conn.lock().await, message).await?; + conn.send(message).await?; Ok(true) } else { diff --git a/src/vc_server.rs b/src/vc_server.rs index 289c289..cabb20c 100644 --- a/src/vc_server.rs +++ b/src/vc_server.rs @@ -4,10 +4,7 @@ use anyhow::Context; use ed25519_dalek::VerifyingKey; use tokio::net::UdpSocket; -use crate::{ - protocol::{ClientMethod, send_socket}, - server::Server, -}; +use crate::{protocol::ClientMethod, server::Server}; use tokio::time::Instant; @@ -107,13 +104,11 @@ impl Server { if now.duration_since(voice.last_speaking_sent).as_millis() >= 600 { for conn in user.connections.lock().await.values() { - let _ = send_socket( - &mut *conn.lock().await, - &ClientMethod::Speaking { + let _ = conn + .send(&ClientMethod::Speaking { pubkey: crate::crypto::to_string(sender), - }, - ) - .await; + }) + .await; } voice.last_speaking_sent = now; diff --git a/src/ws.rs b/src/ws.rs new file mode 100644 index 0000000..de5e692 --- /dev/null +++ b/src/ws.rs @@ -0,0 +1,65 @@ +use std::borrow::Cow; + +use axum::extract::ws::{Message, Utf8Bytes, WebSocket}; +use futures_util::{ + SinkExt, StreamExt, + stream::{SplitSink, SplitStream}, +}; +use tokio::sync::Mutex; + +use crate::protocol::{ClientMethod, ServerMethod}; + +pub struct EnclaveWebSocket { + tx: Mutex>, + rx: Mutex>, +} + +impl EnclaveWebSocket { + pub fn new(ws: WebSocket) -> Self { + let (tx, rx) = ws.split(); + + Self { + tx: Mutex::new(tx), + rx: Mutex::new(rx), + } + } + + pub async fn read(&self) -> anyhow::Result> { + match self.rx.lock().await.next().await.transpose()? { + Some(Message::Text(text)) => match serde_json::from_str(&text.to_string()) { + Ok(msg) => Ok(Some(msg)), + + Err(e) => { + self.send(&ClientMethod::Error { + error: Cow::Owned(format!("Unable to parse message: {e}")), + }) + .await?; + + Ok(None) + } + }, + + Some(Message::Ping(v)) => { + self.tx.lock().await.send(Message::Pong(v)).await?; + + Ok(None) + } + + Some(_) => Ok(None), + + None => Ok(None), + } + } + + pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> { + self.tx + .lock() + .await + .send(Message::Text(Utf8Bytes::from(serde_json::to_string( + message, + )?))) + .await?; + + Ok(()) + } +} From 7cc7a131c3ce125bf23132defb1cf87e6e76bd8a Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 27 Aug 2026 17:12:03 +0200 Subject: [PATCH 10/11] Speaking indicator --- src/vc_server.rs | 27 +++++++++++++++++---------- 1 file changed, 17 insertions(+), 10 deletions(-) diff --git a/src/vc_server.rs b/src/vc_server.rs index cabb20c..5a0a946 100644 --- a/src/vc_server.rs +++ b/src/vc_server.rs @@ -55,11 +55,19 @@ impl Server { }; drop(clients); - *user.voice.lock().await = Some(crate::server::VoiceConnection { - addr, - channel_id: channel_id.clone(), - last_speaking_sent: Instant::now(), - }); + let mut voice = user.voice.lock().await; + + if let Some(voice) = &mut *voice { + voice.addr = addr; + } else { + *voice = Some(crate::server::VoiceConnection { + addr, + channel_id: channel_id.clone(), + last_speaking_sent: Instant::now(), + }); + } + + drop(voice); self.relay_voice(&sender_pubkey, &channel_id, payload) .await?; @@ -104,11 +112,10 @@ impl Server { if now.duration_since(voice.last_speaking_sent).as_millis() >= 600 { for conn in user.connections.lock().await.values() { - let _ = conn - .send(&ClientMethod::Speaking { - pubkey: crate::crypto::to_string(sender), - }) - .await; + conn.send(&ClientMethod::Speaking { + pubkey: crate::crypto::to_string(sender), + }) + .await?; } voice.last_speaking_sent = now; From f8e510bdb6096fcc987256a2c2cd5944a0671ce8 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 27 Aug 2026 20:32:30 +0200 Subject: [PATCH 11/11] More voice stuff --- src/protocol/mod.rs | 34 ++++++++-------------- src/protocol/voice.rs | 67 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 79 insertions(+), 22 deletions(-) create mode 100644 src/protocol/voice.rs diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index f8e1814..9cb3a1b 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -8,6 +8,7 @@ use crate::{data::messages::StoredMessage, server::Server, types::ClientMeta}; pub mod initialize; pub mod message; pub mod user; +pub mod voice; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(tag = "method")] @@ -52,6 +53,11 @@ pub enum ClientMethod { pubkey: String, }, + UserLeftVoice { + channel_id: String, + pubkey: String, + }, + Speaking { pubkey: String, }, @@ -103,6 +109,8 @@ pub enum ServerMethod { JoinVoice { channel_id: String, }, + + LeaveVoice, } pub async fn read_loop( @@ -164,29 +172,11 @@ pub async fn read_loop( } ServerMethod::JoinVoice { channel_id } => { - { - let pin = rand::random::() % (1 << 53); + voice::join(server, verifying_key, socket, channel_id).await?; + } - server - .voice_pins - .lock() - .await - .insert(pin, (verifying_key, channel_id.clone())); - - socket - .send(&ClientMethod::JoinVoice { - channel_id: channel_id.clone(), - pin, - }) - .await?; - } - - server - .broadcast(&ClientMethod::UserJoinedVoice { - channel_id, - pubkey: crate::crypto::to_string(&verifying_key), - }) - .await?; + ServerMethod::LeaveVoice => { + voice::leave(server, verifying_key).await?; } } } diff --git a/src/protocol/voice.rs b/src/protocol/voice.rs new file mode 100644 index 0000000..563225b --- /dev/null +++ b/src/protocol/voice.rs @@ -0,0 +1,67 @@ +use std::sync::Arc; + +use ed25519_dalek::VerifyingKey; + +use crate::{protocol::ClientMethod, server::Server}; + +pub async fn join( + server: &Arc, + verifying_key: VerifyingKey, + socket: &Arc, + channel_id: String, +) -> anyhow::Result<()> { + { + let pin = rand::random::() % (1 << 53); + + server + .voice_pins + .lock() + .await + .insert(pin, (verifying_key, channel_id.clone())); + + socket + .send(&ClientMethod::JoinVoice { + channel_id: channel_id.clone(), + pin, + }) + .await?; + } + + server + .broadcast(&ClientMethod::UserJoinedVoice { + channel_id, + pubkey: crate::crypto::to_string(&verifying_key), + }) + .await?; + + Ok(()) +} + +pub async fn leave(server: &Arc, verifying_key: VerifyingKey) -> anyhow::Result<()> { + let Some(channel_id) = ({ + let clients = server.clients.lock().await; + + let Some(client) = clients.get(&verifying_key) else { + return Ok(()); + }; + + client.voice.lock().await.take().map(|v| v.channel_id) + }) else { + return Ok(()); + }; + + server + .voice_pins + .lock() + .await + .retain(|_, v| v.0 != verifying_key); + + server + .broadcast(&ClientMethod::UserLeftVoice { + channel_id, + pubkey: crate::crypto::to_string(&verifying_key), + }) + .await?; + + Ok(()) +}