Merge pull request #8 from recurse-chat/refactor/server-and-logging

Server refactor, session/voice cleanup, init voice listing, and structured logging
This commit is contained in:
2026-08-29 07:43:02 -04:00
committed by GitHub
18 changed files with 1029 additions and 528 deletions
Generated
+227 -2
View File
@@ -24,6 +24,65 @@ dependencies = [
"zerocopy", "zerocopy",
] ]
[[package]]
name = "aho-corasick"
version = "1.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c982642fa9e8606056828ee9a8505737230110bb1099153c79efe865c59d12ba"
dependencies = [
"memchr",
]
[[package]]
name = "anstream"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "824a212faf96e9acacdbd09febd34438f8f711fb84e09a8916013cd7815ca28d"
dependencies = [
"anstyle",
"anstyle-parse",
"anstyle-query",
"anstyle-wincon",
"colorchoice",
"is_terminal_polyfill",
"utf8parse",
]
[[package]]
name = "anstyle"
version = "1.0.14"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "940b3a0ca603d1eade50a4846a2afffd5ef57a9feac2c0e2ec2e14f9ead76000"
[[package]]
name = "anstyle-parse"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "52ce7f38b242319f7cabaa6813055467063ecdc9d355bbb4ce0c68908cd8130e"
dependencies = [
"utf8parse",
]
[[package]]
name = "anstyle-query"
version = "1.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc"
dependencies = [
"windows-sys",
]
[[package]]
name = "anstyle-wincon"
version = "3.0.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d"
dependencies = [
"anstyle",
"once_cell_polyfill",
"windows-sys",
]
[[package]] [[package]]
name = "anyhow" name = "anyhow"
version = "1.0.104" version = "1.0.104"
@@ -103,6 +162,12 @@ version = "1.8.3"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06" checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06"
[[package]]
name = "bitflags"
version = "1.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a"
[[package]] [[package]]
name = "bitflags" name = "bitflags"
version = "2.13.1" version = "2.13.1"
@@ -204,6 +269,12 @@ version = "0.5.4"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a" checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a"
[[package]]
name = "colorchoice"
version = "1.0.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570"
[[package]] [[package]]
name = "const-oid" name = "const-oid"
version = "0.9.6" version = "0.9.6"
@@ -312,6 +383,37 @@ version = "2.11.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06" checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06"
[[package]]
name = "defmt"
version = "1.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e2953bfe4f93bbd20cc71198842756f77d161884c99ebbabc41d80231ded88d1"
dependencies = [
"bitflags 1.3.2",
"defmt-macros",
]
[[package]]
name = "defmt-macros"
version = "1.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bad9c72e7ca2137e0dc3813245a0d282fd6daad32fd800af018306a9169b5fe8"
dependencies = [
"defmt-parser",
"proc-macro2",
"quote",
"syn 2.0.119",
]
[[package]]
name = "defmt-parser"
version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "10d60334b3b2e7c9d91ef8150abfb6fa4c1c39ebbcf4a81c2e346aad939fee3e"
dependencies = [
"thiserror",
]
[[package]] [[package]]
name = "der" name = "der"
version = "0.7.10" version = "0.7.10"
@@ -378,7 +480,9 @@ dependencies = [
"chacha20poly1305", "chacha20poly1305",
"curve25519-dalek 5.0.0", "curve25519-dalek 5.0.0",
"ed25519-dalek", "ed25519-dalek",
"env_logger",
"futures-util", "futures-util",
"log",
"rand 0.8.7", "rand 0.8.7",
"rusqlite", "rusqlite",
"serde", "serde",
@@ -390,6 +494,29 @@ dependencies = [
"x25519-dalek", "x25519-dalek",
] ]
[[package]]
name = "env_filter"
version = "2.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "900d271a03799a1ee8d1ca9b19893b48ca674a9284fefcfb85f05e74ed314217"
dependencies = [
"log",
"regex",
]
[[package]]
name = "env_logger"
version = "0.11.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "de671bd27a75a797dc9ae289ba1e77276e75e2026408aab65185384e2d5cd3f6"
dependencies = [
"anstream",
"anstyle",
"env_filter",
"jiff",
"log",
]
[[package]] [[package]]
name = "fallible-iterator" name = "fallible-iterator"
version = "0.3.0" version = "0.3.0"
@@ -648,12 +775,54 @@ dependencies = [
"hybrid-array", "hybrid-array",
] ]
[[package]]
name = "is_terminal_polyfill"
version = "1.70.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695"
[[package]] [[package]]
name = "itoa" name = "itoa"
version = "1.0.18" version = "1.0.18"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682"
[[package]]
name = "jiff"
version = "0.2.35"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "668b7183bd07af9a4885f5c35b0cc5c83c4607a913c16b7e17291832910d2dcc"
dependencies = [
"defmt",
"jiff-core",
"jiff-static",
"log",
"portable-atomic",
"portable-atomic-util",
"serde_core",
]
[[package]]
name = "jiff-core"
version = "0.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7feca88439efe53da3754500c1851dedf3cb36c524dd5cf8225cc0794de95d09"
dependencies = [
"defmt",
]
[[package]]
name = "jiff-static"
version = "0.2.35"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3a69dcb3a21cfb32ce1cd056169337ca284af0766dd766e7878819b251a49204"
dependencies = [
"jiff-core",
"proc-macro2",
"quote",
"syn 2.0.119",
]
[[package]] [[package]]
name = "js-sys" name = "js-sys"
version = "0.3.104" version = "0.3.104"
@@ -733,6 +902,12 @@ version = "1.21.4"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50"
[[package]]
name = "once_cell_polyfill"
version = "1.70.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe"
[[package]] [[package]]
name = "percent-encoding" name = "percent-encoding"
version = "2.3.2" version = "2.3.2"
@@ -771,6 +946,21 @@ dependencies = [
"universal-hash", "universal-hash",
] ]
[[package]]
name = "portable-atomic"
version = "1.15.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "05c8b63e8d9609db387f0324918f81d68fe27748f084ef092fb35954d0539a85"
[[package]]
name = "portable-atomic-util"
version = "0.2.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a106d1259c23fac8e543272398ae0e3c0b8d33c88ed73d0cc71b0f1d902618"
dependencies = [
"portable-atomic",
]
[[package]] [[package]]
name = "ppv-lite86" name = "ppv-lite86"
version = "0.2.21" version = "0.2.21"
@@ -875,13 +1065,42 @@ version = "0.10.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69"
[[package]]
name = "regex"
version = "1.13.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f020237b6c8eed93db2e2cb53c00c60a8e1bc73da7d073199a1180401450218d"
dependencies = [
"aho-corasick",
"memchr",
"regex-automata",
"regex-syntax",
]
[[package]]
name = "regex-automata"
version = "0.4.18"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ad8553b9b26413251cbf30e620595c7a41b3887f03da04579c0e6b0d6a06b4b2"
dependencies = [
"aho-corasick",
"memchr",
"regex-syntax",
]
[[package]]
name = "regex-syntax"
version = "0.8.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4"
[[package]] [[package]]
name = "rusqlite" name = "rusqlite"
version = "0.31.0" version = "0.31.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b838eba278d213a8beaf485bd313fd580ca4505a00d5871caeb1457c55322cae" checksum = "b838eba278d213a8beaf485bd313fd580ca4505a00d5871caeb1457c55322cae"
dependencies = [ dependencies = [
"bitflags", "bitflags 2.13.1",
"fallible-iterator", "fallible-iterator",
"fallible-streaming-iterator", "fallible-streaming-iterator",
"hashlink", "hashlink",
@@ -1204,7 +1423,7 @@ version = "0.7.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b11f75e912b0c2be01b63d8cf8057b8c3f97cf34abb3d431a3a4c8675498e233" checksum = "b11f75e912b0c2be01b63d8cf8057b8c3f97cf34abb3d431a3a4c8675498e233"
dependencies = [ dependencies = [
"bitflags", "bitflags 2.13.1",
"bytes", "bytes",
"futures-core", "futures-core",
"futures-util", "futures-util",
@@ -1299,6 +1518,12 @@ dependencies = [
"ctutils", "ctutils",
] ]
[[package]]
name = "utf8parse"
version = "0.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821"
[[package]] [[package]]
name = "uuid" name = "uuid"
version = "1.24.1" version = "1.24.1"
+2
View File
@@ -10,6 +10,8 @@ serde = { version = "1.0.229", features = ["serde_derive"] }
serde_json = "1.0.151" serde_json = "1.0.151"
tokio = { version = "1.53.1", features = ["rt", "rt-multi-thread", "macros", "sync", "fs"] } tokio = { version = "1.53.1", features = ["rt", "rt-multi-thread", "macros", "sync", "fs"] }
ed25519-dalek = { version = "2", features = ["rand_core"] } ed25519-dalek = { version = "2", features = ["rand_core"] }
env_logger = "0.11"
log = "0.4"
rand = "0.8" rand = "0.8"
bs58 = "0.5.1" bs58 = "0.5.1"
tower-http = { version = "0.7.0", features = ["fs", "cors"] } tower-http = { version = "0.7.0", features = ["fs", "cors"] }
+10 -2
View File
@@ -20,8 +20,10 @@ pub async fn get() -> anyhow::Result<SigningKey> {
if !private_key_path.exists() { if !private_key_path.exists() {
let key = SigningKey::generate(&mut OsRng); let key = SigningKey::generate(&mut OsRng);
tokio::fs::write(private_key_path, &key.to_bytes()).await?; tokio::fs::write(private_key_path, &key.to_bytes()).await?;
log::info!("Generated new server signing key at private.key");
Ok(key) Ok(key)
} else { } else {
log::debug!("Loaded existing server signing key from private.key");
Ok(SigningKey::from_bytes( Ok(SigningKey::from_bytes(
&tokio::fs::read(private_key_path) &tokio::fs::read(private_key_path)
.await? .await?
@@ -140,9 +142,11 @@ pub async fn crypto_handshake(
server: &Arc<Server>, server: &Arc<Server>,
mut socket: WebSocket, mut socket: WebSocket,
) -> anyhow::Result<EnclaveWebSocket> { ) -> anyhow::Result<EnclaveWebSocket> {
log::info!("Handshake: sending server x25519 public key");
socket socket
.send(axum::extract::ws::Message::Binary( .send(axum::extract::ws::Message::Binary(
server.x_keypair.0.to_bytes().to_vec().into(), server.identity.x25519.public.to_bytes().to_vec().into(),
)) ))
.await?; .await?;
@@ -159,9 +163,13 @@ pub async fn crypto_handshake(
"Failed to get proper length of client x key" "Failed to get proper length of client x key"
))?); ))?);
let shared_secret = server.x_keypair.1.diffie_hellman(&client_pubkey); log::debug!("Handshake: received client x25519 public key");
let shared_secret = server.identity.x25519.secret.diffie_hellman(&client_pubkey);
let cipher = Arc::new(Mutex::new(SessionCipher::new(&shared_secret)?)); let cipher = Arc::new(Mutex::new(SessionCipher::new(&shared_secret)?));
log::info!("Handshake: session cipher established with client");
Ok(EnclaveWebSocket::new(socket, cipher)) Ok(EnclaveWebSocket::new(socket, cipher))
} }
+9 -3
View File
@@ -58,11 +58,17 @@ impl Config {
if !config_path.exists() { if !config_path.exists() {
let config = Config::new(); let config = Config::new();
tokio::fs::write(config_path, &serde_json::to_string_pretty(&config)?).await?; tokio::fs::write(config_path, &serde_json::to_string_pretty(&config)?).await?;
log::info!("No config.json found, wrote default config");
Ok(config) Ok(config)
} else { } else {
Ok(serde_json::from_str( let config: Config =
&tokio::fs::read_to_string(config_path).await?, serde_json::from_str(&tokio::fs::read_to_string(config_path).await?)?;
)?) log::info!(
"Loaded config: {} (port {})",
config.meta.name,
config.port
);
Ok(config)
} }
} }
} }
+11 -2
View File
@@ -3,7 +3,6 @@ pub mod data;
pub mod protocol; pub mod protocol;
pub mod server; pub mod server;
pub mod types; pub mod types;
pub mod vc_server;
pub mod ws; pub mod ws;
use std::{ use std::{
@@ -26,9 +25,15 @@ use crate::server::Server;
#[tokio::main] #[tokio::main]
async fn main() -> anyhow::Result<()> { async fn main() -> anyhow::Result<()> {
env_logger::Builder::from_env(env_logger::Env::default().default_filter_or("info"))
.format_timestamp_millis()
.init();
log::info!("Starting enclave-server v{}", env!("CARGO_PKG_VERSION"));
let server = Server::new().await?; let server = Server::new().await?;
let udp_server = tokio::spawn(server.clone().start_udp_server()); let udp_server = tokio::spawn(server.voice.clone().run(server.config.port));
let cors = CorsLayer::new() let cors = CorsLayer::new()
.allow_origin(Any) .allow_origin(Any)
@@ -48,8 +53,12 @@ async fn main() -> anyhow::Result<()> {
)) ))
.await?; .await?;
log::info!("HTTP/WS server listening on 0.0.0.0:{}", server.config.port);
axum::serve(listener, app).await?; axum::serve(listener, app).await?;
log::info!("Shutting down UDP voice server");
udp_server.abort(); udp_server.abort();
Ok(()) Ok(())
+30 -6
View File
@@ -5,12 +5,14 @@ use std::{
use ed25519_dalek::{Signer, VerifyingKey}; use ed25519_dalek::{Signer, VerifyingKey};
use crate::{server::Server, ws::EnclaveWebSocket}; use crate::{
server::Server,
types::ClientMeta,
ws::EnclaveWebSocket,
};
use super::*; use super::*;
use crate::server::UserConnections;
impl UserConnections {
pub async fn initialize( pub async fn initialize(
server: &Arc<Server>, server: &Arc<Server>,
socket: &EnclaveWebSocket, socket: &EnclaveWebSocket,
@@ -23,6 +25,8 @@ impl UserConnections {
hostname, hostname,
}) = socket.read().await? }) = socket.read().await?
else { else {
log::warn!("Client sent the wrong method during initialization");
socket socket
.send(&ClientMethod::Error { .send(&ClientMethod::Error {
error: Cow::Borrowed("Initialization required"), error: Cow::Borrowed("Initialization required"),
@@ -37,6 +41,8 @@ impl UserConnections {
let server_timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as u64; let server_timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_millis() as u64;
if server_timestamp.saturating_sub(timestamp) > 2000 { if server_timestamp.saturating_sub(timestamp) > 2000 {
log::warn!("Client timestamp wasn't correct (server {server_timestamp}, client {timestamp})");
socket socket
.send(&ClientMethod::Error { .send(&ClientMethod::Error {
error: Cow::Borrowed( error: Cow::Borrowed(
@@ -49,6 +55,8 @@ impl UserConnections {
} }
if !server.config.hostnames.contains(&hostname) { if !server.config.hostnames.contains(&hostname) {
log::warn!("Client sent invalid hostname: {hostname}");
socket.send( socket.send(
&ClientMethod::Error { &ClientMethod::Error {
error: Cow::Owned(format!("Invalid Hostname, to avoid man-in-the-middle attacks, please use the correct hostname(s): {}", server.config.hostnames.clone().into_iter().collect::<Vec<_>>().join(", "))), error: Cow::Owned(format!("Invalid Hostname, to avoid man-in-the-middle attacks, please use the correct hostname(s): {}", server.config.hostnames.clone().into_iter().collect::<Vec<_>>().join(", "))),
@@ -60,6 +68,8 @@ impl UserConnections {
} }
let Ok(public_key) = crate::crypto::from_string(&public_key_string) else { let Ok(public_key) = crate::crypto::from_string(&public_key_string) else {
log::warn!("Client sent an invalid public key: {public_key_string}");
socket socket
.send(&ClientMethod::Error { .send(&ClientMethod::Error {
error: Cow::Borrowed("Invalid public key"), error: Cow::Borrowed("Invalid public key"),
@@ -76,6 +86,8 @@ impl UserConnections {
) )
.is_err() .is_err()
{ {
log::warn!("Client signature verification failed for {public_key_string}");
socket socket
.send(&ClientMethod::Error { .send(&ClientMethod::Error {
error: Cow::Borrowed("Invalid signature"), error: Cow::Borrowed("Invalid signature"),
@@ -85,11 +97,13 @@ impl UserConnections {
return Err(anyhow::anyhow!("Invalid signature")); return Err(anyhow::anyhow!("Invalid signature"));
} }
log::info!("Client authenticated: {public_key_string} (hostname: {hostname})");
{ {
socket socket
.send(&ClientMethod::Initialized { .send(&ClientMethod::Initialized {
public_key: crate::crypto::to_string(&server.key.verifying_key()), public_key: crate::crypto::to_string(&server.identity.key.verifying_key()),
signature: crate::crypto::to_string_sig(&server.key.sign( signature: crate::crypto::to_string_sig(&server.identity.key.sign(
format!("{server_timestamp}@{hostname}@{public_key_string}").as_bytes(), format!("{server_timestamp}@{hostname}@{public_key_string}").as_bytes(),
)), )),
@@ -100,6 +114,8 @@ impl UserConnections {
} }
let Some(ServerMethod::Meta(meta)) = socket.read().await? else { let Some(ServerMethod::Meta(meta)) = socket.read().await? else {
log::warn!("Client {public_key_string} didn't send meta during initialization");
socket socket
.send(&ClientMethod::Error { .send(&ClientMethod::Error {
error: Cow::Borrowed("Expected meta"), error: Cow::Borrowed("Expected meta"),
@@ -111,6 +127,14 @@ impl UserConnections {
)); ));
}; };
for pin in server.voice.pins.lock().await.values() {
socket
.send(&ClientMethod::UserJoinedVoice {
channel_id: pin.channel_id.clone(),
pubkey: crate::crypto::to_string(&pin.pubkey),
})
.await?;
}
Ok((public_key, meta)) Ok((public_key, meta))
} }
}
+28 -8
View File
@@ -22,7 +22,7 @@ pub async fn send_message(
); );
} }
let server_pubkey_string = crate::crypto::to_string(&server.key.verifying_key()); let server_pubkey_string = crate::crypto::to_string(&server.identity.key.verifying_key());
let signed_string = format!( let signed_string = format!(
"{}@{}@{}", "{}@{}@{}",
message.timestamp, server_pubkey_string, message.content message.timestamp, server_pubkey_string, message.content
@@ -42,9 +42,16 @@ pub async fn send_message(
data: message, data: message,
}; };
server.message_store.insert_message(&channel_id, &stored)?; server.store.messages.insert_message(&channel_id, &stored)?;
log::info!(
"Message sent by {} in channel {channel_id} (id {})",
stored.author,
stored.id
);
server server
.sessions
.broadcast(&ClientMethod::Messages { .broadcast(&ClientMethod::Messages {
messages: HashMap::from([(channel_id, vec![stored])]), messages: HashMap::from([(channel_id, vec![stored])]),
}) })
@@ -63,9 +70,12 @@ pub async fn get_messages(
const CHUNK_SIZE: u32 = 16; const CHUNK_SIZE: u32 = 16;
let messages = server let messages = server
.message_store .store
.messages
.get_recent_messages(&channel_id, CHUNK_SIZE, chunk)?; .get_recent_messages(&channel_id, CHUNK_SIZE, chunk)?;
log::debug!("Serving {} messages for {channel_id} (chunk {chunk})", messages.len());
socket socket
.send(&ClientMethod::Messages { .send(&ClientMethod::Messages {
messages: HashMap::from([(channel_id, messages)]), messages: HashMap::from([(channel_id, messages)]),
@@ -84,7 +94,8 @@ pub async fn edit_message(
new_signature: String, new_signature: String,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let existing = server let existing = server
.message_store .store
.messages
.get_message(&channel_id, &message_id)? .get_message(&channel_id, &message_id)?
.ok_or_else(|| anyhow::anyhow!("Message not found"))?; .ok_or_else(|| anyhow::anyhow!("Message not found"))?;
@@ -93,7 +104,7 @@ pub async fn edit_message(
anyhow::bail!("Not authorized to edit this message"); anyhow::bail!("Not authorized to edit this message");
} }
let server_pubkey_string = crate::crypto::to_string(&server.key.verifying_key()); let server_pubkey_string = crate::crypto::to_string(&server.identity.key.verifying_key());
let signed_string = format!( let signed_string = format!(
"{}@{}@{}", "{}@{}@{}",
existing.data.timestamp, server_pubkey_string, new_content existing.data.timestamp, server_pubkey_string, new_content
@@ -107,9 +118,12 @@ pub async fn edit_message(
.map_err(|_| anyhow::anyhow!("Signature verification failed"))?; .map_err(|_| anyhow::anyhow!("Signature verification failed"))?;
server server
.message_store .store
.messages
.update_message(&channel_id, &message_id, &new_content, &new_signature)?; .update_message(&channel_id, &message_id, &new_content, &new_signature)?;
log::info!("Message {message_id} edited in channel {channel_id}");
let updated = StoredMessage { let updated = StoredMessage {
id: message_id, id: message_id,
author: author_pubkey, author: author_pubkey,
@@ -122,6 +136,7 @@ pub async fn edit_message(
}; };
server server
.sessions
.broadcast(&ClientMethod::MessageEdited { .broadcast(&ClientMethod::MessageEdited {
channel_id: channel_id.clone(), channel_id: channel_id.clone(),
message: updated, message: updated,
@@ -138,7 +153,8 @@ pub async fn delete_message(
channel_id: String, channel_id: String,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let existing = server let existing = server
.message_store .store
.messages
.get_message(&channel_id, &message_id)? .get_message(&channel_id, &message_id)?
.ok_or_else(|| anyhow::anyhow!("Message not found"))?; .ok_or_else(|| anyhow::anyhow!("Message not found"))?;
@@ -148,10 +164,14 @@ pub async fn delete_message(
} }
server server
.message_store .store
.messages
.delete_message(&channel_id, &message_id)?; .delete_message(&channel_id, &message_id)?;
log::info!("Message {message_id} deleted in channel {channel_id}");
server server
.sessions
.broadcast(&ClientMethod::MessageDeleted { .broadcast(&ClientMethod::MessageDeleted {
channel_id: channel_id.clone(), channel_id: channel_id.clone(),
message_id, message_id,
+5 -1
View File
@@ -118,9 +118,13 @@ pub async fn read_loop(
verifying_key: VerifyingKey, verifying_key: VerifyingKey,
socket: &Arc<crate::ws::EnclaveWebSocket>, socket: &Arc<crate::ws::EnclaveWebSocket>,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let pubkey_string = crate::crypto::to_string(&verifying_key);
while let Some(message) = socket.read().await? { while let Some(message) = socket.read().await? {
match message { match message {
ServerMethod::Initialize { .. } => { ServerMethod::Initialize { .. } => {
log::warn!("Client {pubkey_string} sent Initialize after already initializing");
socket socket
.send(&ClientMethod::Error { .send(&ClientMethod::Error {
error: Cow::Borrowed("Already initialized"), error: Cow::Borrowed("Already initialized"),
@@ -132,7 +136,7 @@ pub async fn read_loop(
ServerMethod::Meta(meta) => {} ServerMethod::Meta(meta) => {}
ServerMethod::Error { error } => { ServerMethod::Error { error } => {
eprintln!("Client error: {error}"); log::warn!("Client error (client {pubkey_string}): {error}");
} }
ServerMethod::SendMessage { channel_id, data } => { ServerMethod::SendMessage { channel_id, data } => {
+3 -1
View File
@@ -10,7 +10,9 @@ pub async fn get_users(
socket: &Arc<crate::ws::EnclaveWebSocket>, socket: &Arc<crate::ws::EnclaveWebSocket>,
pubkeys: Vec<String>, pubkeys: Vec<String>,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
let users = server.user_store.get_users(&pubkeys).await?; let users = server.store.users.get_users(&pubkeys).await?;
log::debug!("Serving {} user infos", users.len());
socket.send(&ClientMethod::Users { users }).await?; socket.send(&ClientMethod::Users { users }).await?;
+18 -23
View File
@@ -10,14 +10,18 @@ pub async fn join(
socket: &Arc<crate::ws::EnclaveWebSocket>, socket: &Arc<crate::ws::EnclaveWebSocket>,
channel_id: String, channel_id: String,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
{ let user = server
let pin = rand::random::<u64>() % (1 << 53); .sessions
.get(&verifying_key)
server
.voice_pins
.lock()
.await .await
.insert(pin, (verifying_key, channel_id.clone())); .ok_or_else(|| anyhow::anyhow!("Not connected"))?;
let pin = server.voice.join(verifying_key, user, &channel_id).await;
log::info!(
"User {} joining voice channel {channel_id}",
crate::crypto::to_string(&verifying_key)
);
socket socket
.send(&ClientMethod::JoinVoice { .send(&ClientMethod::JoinVoice {
@@ -25,9 +29,9 @@ pub async fn join(
pin, pin,
}) })
.await?; .await?;
}
server server
.sessions
.broadcast(&ClientMethod::UserJoinedVoice { .broadcast(&ClientMethod::UserJoinedVoice {
channel_id, channel_id,
pubkey: crate::crypto::to_string(&verifying_key), pubkey: crate::crypto::to_string(&verifying_key),
@@ -38,25 +42,16 @@ pub async fn join(
} }
pub async fn leave(server: &Arc<Server>, verifying_key: VerifyingKey) -> anyhow::Result<()> { pub async fn leave(server: &Arc<Server>, verifying_key: VerifyingKey) -> anyhow::Result<()> {
let Some(channel_id) = ({ let Some(channel_id) = server.voice.remove(verifying_key).await else {
let clients = server.clients.lock().await; log::debug!(
"User {} requested LeaveVoice but wasn't in any channel",
let Some(client) = clients.get(&verifying_key) else { crate::crypto::to_string(&verifying_key)
return Ok(()); );
};
client.voice.lock().await.take().map(|v| v.channel_id)
}) else {
return Ok(()); return Ok(());
}; };
server server
.voice_pins .sessions
.lock()
.await
.retain(|_, v| v.0 != verifying_key);
server
.broadcast(&ClientMethod::UserLeftVoice { .broadcast(&ClientMethod::UserLeftVoice {
channel_id, channel_id,
pubkey: crate::crypto::to_string(&verifying_key), pubkey: crate::crypto::to_string(&verifying_key),
-191
View File
@@ -1,191 +0,0 @@
use std::{
collections::HashMap,
net::SocketAddr,
path::PathBuf,
sync::{Arc, atomic::AtomicU16},
};
use axum::{
extract::{WebSocketUpgrade, ws::WebSocket},
response::Response,
};
use ed25519_dalek::{SigningKey, VerifyingKey};
use tokio::{
net::UdpSocket,
sync::{Mutex, OnceCell},
task::JoinSet,
time::Instant,
};
use crate::{
crypto::SessionCipher,
data::{config::Config, messages::MessageStore, users::UserMetaStore},
protocol::{ClientMethod, read_loop},
types::ClientMeta,
};
use x25519_dalek::{PublicKey as X25519Public, StaticSecret as X25519Secret};
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<HashMap<u16, Arc<crate::ws::EnclaveWebSocket>>>,
pub cihper: Arc<Mutex<SessionCipher>>,
pub voice: Mutex<Option<VoiceConnection>>,
}
pub struct Server {
pub key: SigningKey,
pub x_keypair: (X25519Public, X25519Secret),
pub config: Config,
pub clients: Mutex<HashMap<VerifyingKey, Arc<UserConnections>>>,
pub voice_pins: Mutex<HashMap<u64, (VerifyingKey, String)>>,
pub message_store: MessageStore,
pub user_store: UserMetaStore,
pub voice_socket: OnceCell<UdpSocket>,
}
impl Server {
pub async fn new() -> anyhow::Result<Arc<Self>> {
let key = crate::crypto::get().await?;
Ok(Arc::new(Self {
x_keypair: (
crate::crypto::ed25519_verifying_key_to_x25519(&key.verifying_key())
.ok_or(anyhow::anyhow!("Failed to convert ed pubkey to x"))?,
crate::crypto::ed25519_signing_key_to_x25519(&key),
),
key,
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"))?,
voice_socket: OnceCell::new(),
}))
}
}
impl Server {
pub async fn ws_handler(self: &Arc<Self>, ws: WebSocketUpgrade) -> Response {
let s = self.clone();
ws.on_upgrade(move |socket: WebSocket| async move {
let mut client = match crate::crypto::crypto_handshake(&s, socket).await {
Ok(client) => client,
Err(err) => {
eprintln!("Failed to initialize crypto: {err}");
return;
}
};
match UserConnections::initialize(&s, &client).await {
Ok((public_key, meta)) => {
if let Err(e) = s
.user_store
.upsert_user(&crate::crypto::to_string(&public_key), &meta)
.await
{
eprintln!("Failed to upsert client: {e}");
}
let mut clients_meta = s.clients.lock().await;
let clients = clients_meta
.entry(public_key)
.or_insert_with(|| {
Arc::new(UserConnections {
meta,
public_key: public_key,
counter: AtomicU16::new(0),
connections: Mutex::new(HashMap::new()),
voice: Mutex::new(None),
cihper: client.cipher.clone(),
})
})
.clone();
client.cipher = clients.cihper.clone();
let client = Arc::new(client);
let conid = clients
.counter
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
clients
.connections
.lock()
.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}");
}
let mut clients_meta = s.clients.lock().await;
let mut connections = clients.connections.lock().await;
connections.remove(&conid);
if connections.len() == 0 {
clients_meta.remove(&public_key);
}
}
Err(e) => {
eprintln!("Failed to initialize client: {e}")
}
}
})
}
pub async fn broadcast(self: &Arc<Self>, 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.lock().await.iter() {
conn.send(message).await?;
}
Ok(())
}
pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result<bool> {
if let Some(conn) = self.connections.lock().await.get(&id) {
conn.send(message).await?;
Ok(true)
} else {
Ok(false)
}
}
}
+26
View File
@@ -0,0 +1,26 @@
use ed25519_dalek::SigningKey;
use x25519_dalek::{PublicKey as X25519Public, StaticSecret as X25519Secret};
pub struct ServerIdentity {
pub key: SigningKey,
pub x25519: X25519KeyPair,
}
pub struct X25519KeyPair {
pub public: X25519Public,
pub secret: X25519Secret,
}
impl ServerIdentity {
pub async fn load() -> anyhow::Result<Self> {
let key = crate::crypto::get().await?;
let x25519 = X25519KeyPair {
public: crate::crypto::ed25519_verifying_key_to_x25519(&key.verifying_key())
.ok_or(anyhow::anyhow!("Failed to convert ed pubkey to x"))?,
secret: crate::crypto::ed25519_signing_key_to_x25519(&key),
};
Ok(Self { key, x25519 })
}
}
+102
View File
@@ -0,0 +1,102 @@
pub mod identity;
pub mod session;
pub mod store;
pub mod vc_server;
use std::sync::Arc;
use axum::{
extract::{WebSocketUpgrade, ws::WebSocket},
response::Response,
};
use crate::{
data::config::Config,
protocol::{read_loop, ClientMethod},
};
pub use identity::{ServerIdentity, X25519KeyPair};
pub use session::{SessionRegistry, UserConnections};
pub use store::DataStore;
pub use vc_server::{VoiceConnection, VoicePin, VoiceServer};
pub struct Server {
pub identity: ServerIdentity,
pub config: Config,
pub sessions: SessionRegistry,
pub voice: Arc<VoiceServer>,
pub store: DataStore,
}
impl Server {
pub async fn new() -> anyhow::Result<Arc<Self>> {
Ok(Arc::new(Self {
identity: ServerIdentity::load().await?,
config: Config::get().await?,
sessions: SessionRegistry::new(),
voice: Arc::new(VoiceServer::new()),
store: DataStore::new()?,
}))
}
pub async fn ws_handler(self: &Arc<Self>, ws: WebSocketUpgrade) -> Response {
let s = self.clone();
ws.on_upgrade(move |socket: WebSocket| async move {
let client = match crate::crypto::crypto_handshake(&s, socket).await {
Ok(client) => client,
Err(err) => {
log::error!("Failed to initialize crypto: {err}");
return;
}
};
match crate::protocol::initialize::initialize(&s, &client).await {
Ok((public_key, meta)) => {
if let Err(e) = s
.store
.users
.upsert_user(&crate::crypto::to_string(&public_key), &meta)
.await
{
log::error!("Failed to upsert client: {e}");
}
let pubkey_string = crate::crypto::to_string(&public_key);
log::info!("Client connected: {pubkey_string}");
let (client, conid) = s.sessions.register(public_key, meta, client).await;
if let Err(e) = read_loop(&s, public_key, &client).await {
log::warn!("Client read loop errored: {e}");
}
if s.sessions.deregister(public_key, conid).await {
log::debug!("Deregistered last connection for {pubkey_string}");
if let Some(channel_id) = s.voice.remove(public_key).await {
log::info!(
"User left voice after disconnect: {pubkey_string} ({channel_id})"
);
s.sessions
.broadcast(&ClientMethod::UserLeftVoice {
channel_id,
pubkey: pubkey_string.clone(),
})
.await
.ok();
}
}
log::info!("Client disconnected: {pubkey_string}");
}
Err(e) => {
log::warn!("Failed to initialize client: {e}")
}
}
})
}
}
+143
View File
@@ -0,0 +1,143 @@
use std::{
collections::HashMap,
sync::{Arc, atomic::AtomicU16},
};
use ed25519_dalek::VerifyingKey;
use tokio::sync::Mutex;
use crate::{
crypto::SessionCipher,
protocol::ClientMethod,
types::ClientMeta,
ws::EnclaveWebSocket,
};
pub struct UserConnections {
pub meta: ClientMeta,
pub counter: AtomicU16,
pub public_key: VerifyingKey,
pub connections: Mutex<HashMap<u16, Arc<EnclaveWebSocket>>>,
pub cipher: Arc<Mutex<SessionCipher>>,
}
pub struct SessionRegistry {
pub clients: Mutex<HashMap<VerifyingKey, Arc<UserConnections>>>,
}
impl Default for SessionRegistry {
fn default() -> Self {
Self::new()
}
}
impl SessionRegistry {
pub fn new() -> Self {
Self {
clients: Mutex::new(HashMap::new()),
}
}
pub async fn get(&self, public_key: &VerifyingKey) -> Option<Arc<UserConnections>> {
self.clients.lock().await.get(public_key).cloned()
}
/// Registers a new websocket connection for a user, returning the
/// connection and its id. The websocket is given the shared cipher
/// of the user's existing connections so voice traffic stays keyed
/// consistently across devices.
pub async fn register(
&self,
public_key: VerifyingKey,
meta: ClientMeta,
mut client: EnclaveWebSocket,
) -> (Arc<EnclaveWebSocket>, u16) {
let mut clients = self.clients.lock().await;
let user = clients
.entry(public_key)
.or_insert_with(|| {
Arc::new(UserConnections {
meta,
public_key,
counter: AtomicU16::new(0),
connections: Mutex::new(HashMap::new()),
cipher: client.cipher.clone(),
})
})
.clone();
client.cipher = user.cipher.clone();
let client = Arc::new(client);
let conid = user.counter.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
user.connections.lock().await.insert(conid, client.clone());
log::debug!(
"Registered connection {conid} for user {}",
crate::crypto::to_string(&public_key)
);
(client, conid)
}
/// Removes a connection from a user's session. Returns `true` when
/// that was the user's last connection and they have been dropped
/// from the registry entirely.
pub async fn deregister(&self, public_key: VerifyingKey, conid: u16) -> bool {
let mut clients = self.clients.lock().await;
let Some(user) = clients.get(&public_key).cloned() else {
return false;
};
let mut connections = user.connections.lock().await;
connections.remove(&conid);
log::debug!(
"Deregistered connection {conid} for user {}",
crate::crypto::to_string(&public_key)
);
if connections.is_empty() {
clients.remove(&public_key);
true
} else {
false
}
}
pub async fn broadcast(&self, message: &ClientMethod) -> anyhow::Result<()> {
let users: Vec<_> = self.clients.lock().await.values().cloned().collect();
for user in users {
if let Err(e) = user.send(message).await {
log::warn!("Failed to send to a client: {e:?}");
}
}
Ok(())
}
}
impl UserConnections {
pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> {
for (_, conn) in self.connections.lock().await.iter() {
conn.send(message).await?;
}
Ok(())
}
pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result<bool> {
if let Some(conn) = self.connections.lock().await.get(&id) {
conn.send(message).await?;
Ok(true)
} else {
Ok(false)
}
}
}
+21
View File
@@ -0,0 +1,21 @@
use std::path::PathBuf;
use crate::data::{messages::MessageStore, users::UserMetaStore};
pub struct DataStore {
pub messages: MessageStore,
pub users: UserMetaStore,
}
impl DataStore {
pub fn new() -> anyhow::Result<Self> {
let store = Self {
messages: MessageStore::new(PathBuf::from("messages"))?,
users: UserMetaStore::new(PathBuf::from("users.db"))?,
};
log::debug!("Opened data store (messages/, users.db)");
Ok(store)
}
}
+268
View File
@@ -0,0 +1,268 @@
use std::{
collections::HashMap,
net::{IpAddr, Ipv4Addr, SocketAddr},
sync::Arc,
};
use anyhow::Context;
use ed25519_dalek::VerifyingKey;
use tokio::{
net::UdpSocket,
sync::{Mutex, OnceCell},
time::Instant,
};
use crate::{
crypto::SessionCipher,
protocol::ClientMethod,
server::UserConnections,
};
pub struct VoicePin {
pub pubkey: VerifyingKey,
pub channel_id: String,
}
pub struct VoiceConnection {
pub user: Arc<UserConnections>,
pub channel_id: String,
pub addr: SocketAddr,
pub last_speaking_sent: Instant,
}
pub struct VoiceServer {
pub pins: Mutex<HashMap<u64, VoicePin>>,
pub socket: OnceCell<UdpSocket>,
pub participants: Mutex<HashMap<VerifyingKey, VoiceConnection>>,
}
impl Default for VoiceServer {
fn default() -> Self {
Self::new()
}
}
impl VoiceServer {
pub fn new() -> Self {
Self {
pins: Mutex::new(HashMap::new()),
socket: OnceCell::new(),
participants: Mutex::new(HashMap::new()),
}
}
/// Adds a user to a voice channel and allocates a one-time pin that
/// their next UDP packet must carry to bind their address.
pub async fn join(
&self,
public_key: VerifyingKey,
user: Arc<UserConnections>,
channel_id: &str,
) -> u64 {
let pin = rand::random::<u64>() % (1 << 53);
self.participants.lock().await.insert(
public_key,
VoiceConnection {
user,
channel_id: channel_id.to_string(),
addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0),
last_speaking_sent: Instant::now(),
},
);
self.pins.lock().await.insert(
pin,
VoicePin {
pubkey: public_key,
channel_id: channel_id.to_string(),
},
);
log::debug!(
"User {} joined voice channel {channel_id} (pin allocated)",
crate::crypto::to_string(&public_key)
);
pin
}
/// Removes a user from every voice channel/pin they hold, returning
/// the channel they were in, if any.
pub async fn remove(&self, public_key: VerifyingKey) -> Option<String> {
let channel_id = self
.participants
.lock()
.await
.remove(&public_key)?
.channel_id;
self.pins.lock().await.retain(|_, pin| pin.pubkey != public_key);
log::info!(
"User {} left voice channel {channel_id}",
crate::crypto::to_string(&public_key)
);
Some(channel_id)
}
/// Runs the UDP loop that terminates voice audio: binds addresses on
/// the first (pin-bearing) packet and relays encrypted audio between
/// the participants of each channel.
pub async fn run(self: Arc<Self>, port: u16) -> anyhow::Result<()> {
let socket = UdpSocket::bind(("0.0.0.0", port)).await?;
self.socket
.set(socket)
.map_err(|_| anyhow::anyhow!("UDP server already started"))?;
let mut buf = [0u8; 4096];
log::info!("UDP voice server listening on port {port}");
loop {
let (len, addr) = self.get_voice_socket()?.recv_from(&mut buf).await?;
if len < 8 {
log::warn!("Dropping voice packet too short for pin from {addr}");
continue;
}
let pin_bytes: [u8; 8] = buf[..8].try_into().unwrap();
let pin = u64::from_be_bytes(pin_bytes);
// is_first_packet distinguishes the plaintext pin-bootstrap packet
// from subsequent encrypted audio packets.
let (sender_pubkey, channel_id, payload, is_first_packet) = {
let mut pins = self.pins.lock().await;
if let Some(pin) = pins.remove(&pin) {
log::debug!("Voice address {addr} bound via pin for user {}", crate::crypto::to_string(&pin.pubkey));
(pin.pubkey, pin.channel_id, &buf[8..len], true)
} else {
drop(pins);
match self.find_sender(&addr).await {
Some((pubkey, channel_id)) => (pubkey, channel_id, &buf[..len], false),
None => {
log::warn!("Voice packet from unknown address {addr}");
continue;
}
}
}
};
if let Some(participant) = self.participants.lock().await.get_mut(&sender_pubkey) {
participant.addr = addr;
}
// The pin-bearing bootstrap packet carries no payload to decrypt —
// it's purely "here's my pin, bind my address."
if is_first_packet {
continue;
}
let Some(cipher) = ({
let participants = self.participants.lock().await;
participants.get(&sender_pubkey).map(|p| p.user.cipher.clone())
}) else {
log::warn!("Voice packet from user not in any channel, dropping");
continue;
};
let decrypted_payload = match cipher.lock().await.decrypt(payload) {
Ok(pt) => pt,
Err(e) => {
log::warn!("Dropping voice packet: decryption failed: {e}");
continue;
}
};
let s = self.clone();
tokio::spawn(async move {
if let Err(e) = s
.relay_voice(&sender_pubkey, &channel_id, &decrypted_payload)
.await
{
log::warn!("Voice relay error: {e}");
}
});
}
}
/// Looks up which known voice participant a UDP address belongs to,
/// for packets arriving after the initial pin-bearing packet.
async fn find_sender(&self, addr: &SocketAddr) -> Option<(VerifyingKey, String)> {
let participants = self.participants.lock().await;
for (pubkey, participant) in participants.iter() {
if *addr == participant.addr {
return Some((*pubkey, participant.channel_id.clone()));
}
}
None
}
/// Sends `payload` to every voice participant currently in `channel_id`.
async fn relay_voice(
&self,
sender: &VerifyingKey,
channel_id: &str,
payload: &[u8],
) -> anyhow::Result<()> {
let mut participants = self.participants.lock().await;
for (_pubkey, participant) in participants.iter_mut() {
if channel_id != participant.channel_id {
continue;
}
let now = Instant::now();
if now.duration_since(participant.last_speaking_sent).as_millis() >= 600 {
participant
.user
.send(&ClientMethod::Speaking {
pubkey: crate::crypto::to_string(sender),
})
.await?;
participant.last_speaking_sent = now;
}
if *sender == participant.user.public_key {
continue;
}
let _ = self
.udp_send_to(&participant.user.cipher, &participant.addr, payload)
.await;
}
Ok(())
}
pub fn get_voice_socket(&self) -> anyhow::Result<&UdpSocket> {
self.socket
.get()
.context("Failed to get voice socket")
}
pub async fn udp_send_to(
&self,
cipher: &Arc<Mutex<SessionCipher>>,
addr: &SocketAddr,
payload: &[u8],
) -> anyhow::Result<()> {
let socket = self.get_voice_socket()?;
socket
.send_to(&cipher.lock().await.encrypt(payload)?, addr)
.await?;
Ok(())
}
}
-179
View File
@@ -1,179 +0,0 @@
use std::{net::SocketAddr, sync::Arc};
use anyhow::Context;
use ed25519_dalek::VerifyingKey;
use tokio::{net::UdpSocket, sync::Mutex};
use crate::{crypto::SessionCipher, protocol::ClientMethod, server::Server};
use tokio::time::Instant;
impl Server {
pub async fn start_udp_server(self: Arc<Self>) -> 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];
eprintln!("[vc] UDP server listening on port {}", self.config.port);
loop {
let (len, addr) = self.get_voice_socket()?.recv_from(&mut buf).await?;
if len < 8 {
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);
// is_first_packet distinguishes the plaintext pin-bootstrap packet
// from subsequent encrypted audio packets.
let (sender_pubkey, channel_id, payload, is_first_packet) = {
let mut pins = self.voice_pins.lock().await;
if let Some((pubkey, channel_id)) = pins.remove(&pin) {
(pubkey, channel_id, &buf[8..len], true)
} else {
drop(pins);
match self.find_voice_sender(&addr).await {
Some((pubkey, channel_id)) => (pubkey, channel_id, &buf[..len], false),
None => {
continue;
}
}
}
};
let clients = self.clients.lock().await;
let Some(user) = clients.get(&sender_pubkey).cloned() else {
eprintln!("[vc] sender pubkey not found in clients, dropping");
continue;
};
drop(clients);
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);
// The pin-bearing bootstrap packet carries no payload to decrypt —
// it's purely "here's my pin, bind my address." Everything after
// this first packet is the real, encrypted audio stream.
if is_first_packet {
continue;
}
let decrypted_payload = match user.cihper.lock().await.decrypt(payload) {
Ok(pt) => pt,
Err(e) => {
eprintln!("[vc] dropping packet: decryption failed: {e}");
continue;
}
};
let s = self.clone();
tokio::spawn(async move {
if let Err(e) = s
.relay_voice(&sender_pubkey, &channel_id, &decrypted_payload)
.await
{
eprintln!("{e}");
}
});
}
}
/// 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(voice) = &*user.voice.lock().await {
if *addr == voice.addr {
return Some((*pubkey, voice.channel_id.clone()));
}
}
}
None
}
/// Sends `payload` to every voice participant currently in `channel_id`.
async fn relay_voice(
&self,
sender: &VerifyingKey,
channel_id: &str,
payload: &[u8],
) -> anyhow::Result<()> {
let clients = self.clients.lock().await;
for (_pubkey, user) in clients.iter() {
let Some(voice) = &mut *user.voice.lock().await else {
continue;
};
if channel_id != voice.channel_id {
continue;
}
let now = Instant::now();
if now.duration_since(voice.last_speaking_sent).as_millis() >= 600 {
for conn in user.connections.lock().await.values() {
conn.send(&ClientMethod::Speaking {
pubkey: crate::crypto::to_string(sender),
})
.await?;
}
voice.last_speaking_sent = now;
}
if *sender == user.public_key {
continue;
}
let _ = self.udp_send_to(&user.cihper, &voice.addr, payload).await;
}
Ok(())
}
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,
cipher: &Arc<Mutex<SessionCipher>>,
addr: &SocketAddr,
payload: &[u8],
) -> anyhow::Result<()> {
let socket = self.get_voice_socket()?;
socket
.send_to(&cipher.lock().await.encrypt(payload)?, addr)
.await?;
Ok(())
}
}
+22 -6
View File
@@ -31,10 +31,16 @@ impl EnclaveWebSocket {
pub async fn read(&self) -> anyhow::Result<Option<ServerMethod>> { pub async fn read(&self) -> anyhow::Result<Option<ServerMethod>> {
match self.rx.lock().await.next().await.transpose()? { match self.rx.lock().await.next().await.transpose()? {
Some(Message::Text(text)) => match serde_json::from_str(&text.to_string()) { Some(Message::Text(text)) => {
Ok(msg) => Ok(Some(msg)), let text = text.to_string();
match serde_json::from_str::<ServerMethod>(&text) {
Ok(msg) => {
log::debug!("Received message: {msg:?}");
Ok(Some(msg))
}
Err(e) => { Err(e) => {
log::warn!("Failed to parse client message: {e}");
self.send(&ClientMethod::Error { self.send(&ClientMethod::Error {
error: Cow::Owned(format!("Unable to parse message: {e}")), error: Cow::Owned(format!("Unable to parse message: {e}")),
}) })
@@ -42,15 +48,20 @@ impl EnclaveWebSocket {
Ok(None) Ok(None)
} }
}, }
}
Some(Message::Binary(encrypted)) => { Some(Message::Binary(encrypted)) => {
let text = String::from_utf8(self.cipher.lock().await.decrypt(&encrypted)?)?; let text = String::from_utf8(self.cipher.lock().await.decrypt(&encrypted)?)?;
match serde_json::from_str(&text.to_string()) { match serde_json::from_str::<ServerMethod>(&text) {
Ok(msg) => Ok(Some(msg)), Ok(msg) => {
log::debug!("Received message: {msg:?}");
Ok(Some(msg))
}
Err(e) => { Err(e) => {
log::warn!("Failed to parse client message: {e}");
self.send(&ClientMethod::Error { self.send(&ClientMethod::Error {
error: Cow::Owned(format!("Unable to parse message: {e}")), error: Cow::Owned(format!("Unable to parse message: {e}")),
}) })
@@ -67,7 +78,10 @@ impl EnclaveWebSocket {
Ok(None) Ok(None)
} }
Some(_) => Ok(None), Some(other) => {
log::debug!("Ignoring websocket message: {other:?}");
Ok(None)
}
None => Ok(None), None => Ok(None),
} }
@@ -76,6 +90,8 @@ impl EnclaveWebSocket {
pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> { pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> {
let text = serde_json::to_string(message)?; let text = serde_json::to_string(message)?;
log::debug!("Sending message: {message:?}");
let encrypted = self.cipher.lock().await.encrypt(text.as_bytes())?; let encrypted = self.cipher.lock().await.encrypt(text.as_bytes())?;
self.tx self.tx