Generated
+13
@@ -256,6 +256,7 @@ dependencies = [
|
|||||||
"axum",
|
"axum",
|
||||||
"bs58",
|
"bs58",
|
||||||
"ed25519-dalek",
|
"ed25519-dalek",
|
||||||
|
"futures-util",
|
||||||
"rand 0.8.7",
|
"rand 0.8.7",
|
||||||
"rusqlite",
|
"rusqlite",
|
||||||
"serde",
|
"serde",
|
||||||
@@ -313,6 +314,17 @@ version = "0.3.34"
|
|||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "92d699e522242e69e3003b94ecc1f960f3a5e015aa7c5d7486e65ad01dd94f5e"
|
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]]
|
[[package]]
|
||||||
name = "futures-sink"
|
name = "futures-sink"
|
||||||
version = "0.3.34"
|
version = "0.3.34"
|
||||||
@@ -332,6 +344,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
|||||||
checksum = "0d50a92467f8ba5dd6e3ee5d4bd04d73ab2e4e1c44474a0674821dfce14b79bc"
|
checksum = "0d50a92467f8ba5dd6e3ee5d4bd04d73ab2e4e1c44474a0674821dfce14b79bc"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"futures-core",
|
"futures-core",
|
||||||
|
"futures-macro",
|
||||||
"futures-sink",
|
"futures-sink",
|
||||||
"futures-task",
|
"futures-task",
|
||||||
"pin-project-lite",
|
"pin-project-lite",
|
||||||
|
|||||||
@@ -15,3 +15,4 @@ bs58 = "0.5.1"
|
|||||||
tower-http = { version = "0.7.0", features = ["fs", "cors"] }
|
tower-http = { version = "0.7.0", features = ["fs", "cors"] }
|
||||||
rusqlite = { version = "0.31", features = ["bundled"] }
|
rusqlite = { version = "0.31", features = ["bundled"] }
|
||||||
uuid = { version = "1.24.1", features = ["v4"] }
|
uuid = { version = "1.24.1", features = ["v4"] }
|
||||||
|
futures-util = "0.3.34"
|
||||||
|
|||||||
+24
-10
@@ -19,18 +19,32 @@ impl Config {
|
|||||||
name: "New Server".to_string(),
|
name: "New Server".to_string(),
|
||||||
description: String::new(),
|
description: String::new(),
|
||||||
|
|
||||||
channels: vec![Channel {
|
channels: vec![
|
||||||
id: "text-channels".to_string(),
|
Channel {
|
||||||
name: "Text Channels".to_string(),
|
id: "text-channels".to_string(),
|
||||||
|
name: "Text Channels".to_string(),
|
||||||
|
|
||||||
data: ChannelKind::Category {
|
data: ChannelKind::Category {
|
||||||
channels: vec![Channel {
|
channels: vec![Channel {
|
||||||
id: "general".to_string(),
|
id: "general".to_string(),
|
||||||
name: "General".to_string(),
|
name: "General".to_string(),
|
||||||
data: ChannelKind::Text,
|
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,
|
port: 3415,
|
||||||
|
|||||||
+7
-1
@@ -1,8 +1,10 @@
|
|||||||
|
pub mod crypto;
|
||||||
pub mod data;
|
pub mod data;
|
||||||
pub mod protocol;
|
pub mod protocol;
|
||||||
pub mod server;
|
pub mod server;
|
||||||
pub mod signature;
|
|
||||||
pub mod types;
|
pub mod types;
|
||||||
|
pub mod vc_server;
|
||||||
|
pub mod ws;
|
||||||
|
|
||||||
use std::{
|
use std::{
|
||||||
net::{IpAddr, Ipv4Addr, SocketAddr},
|
net::{IpAddr, Ipv4Addr, SocketAddr},
|
||||||
@@ -26,6 +28,8 @@ use crate::server::Server;
|
|||||||
async fn main() -> anyhow::Result<()> {
|
async fn main() -> anyhow::Result<()> {
|
||||||
let server = Server::new().await?;
|
let server = Server::new().await?;
|
||||||
|
|
||||||
|
let udp_server = tokio::spawn(server.clone().start_udp_server());
|
||||||
|
|
||||||
let cors = CorsLayer::new()
|
let cors = CorsLayer::new()
|
||||||
.allow_origin(Any)
|
.allow_origin(Any)
|
||||||
.allow_methods(Any)
|
.allow_methods(Any)
|
||||||
@@ -46,6 +50,8 @@ async fn main() -> anyhow::Result<()> {
|
|||||||
|
|
||||||
axum::serve(listener, app).await?;
|
axum::serve(listener, app).await?;
|
||||||
|
|
||||||
|
udp_server.abort();
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+34
-48
@@ -3,10 +3,9 @@ use std::{
|
|||||||
time::{SystemTime, UNIX_EPOCH},
|
time::{SystemTime, UNIX_EPOCH},
|
||||||
};
|
};
|
||||||
|
|
||||||
use axum::extract::ws::WebSocket;
|
|
||||||
use ed25519_dalek::{Signer, VerifyingKey};
|
use ed25519_dalek::{Signer, VerifyingKey};
|
||||||
|
|
||||||
use crate::server::Server;
|
use crate::{server::Server, ws::EnclaveWebSocket};
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
use crate::server::UserConnections;
|
use crate::server::UserConnections;
|
||||||
@@ -14,23 +13,21 @@ use crate::server::UserConnections;
|
|||||||
impl UserConnections {
|
impl UserConnections {
|
||||||
pub async fn initialize(
|
pub async fn initialize(
|
||||||
server: &Arc<Server>,
|
server: &Arc<Server>,
|
||||||
mut socket: WebSocket,
|
socket: Arc<EnclaveWebSocket>,
|
||||||
) -> anyhow::Result<(WebSocket, VerifyingKey, ClientMeta)> {
|
) -> anyhow::Result<(Arc<EnclaveWebSocket>, VerifyingKey, ClientMeta)> {
|
||||||
let Some(ServerMethod::Initialize {
|
let Some(ServerMethod::Initialize {
|
||||||
public_key: public_key_string,
|
public_key: public_key_string,
|
||||||
signature,
|
signature,
|
||||||
|
|
||||||
timestamp,
|
timestamp,
|
||||||
hostname,
|
hostname,
|
||||||
}) = read_socket(&mut socket).await?
|
}) = socket.read().await?
|
||||||
else {
|
else {
|
||||||
send_socket(
|
socket
|
||||||
&mut socket,
|
.send(&ClientMethod::Error {
|
||||||
&ClientMethod::Error {
|
|
||||||
error: Cow::Borrowed("Initialization required"),
|
error: Cow::Borrowed("Initialization required"),
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
|
|
||||||
return Err(anyhow::anyhow!(
|
return Err(anyhow::anyhow!(
|
||||||
"Failed to initialize: Client sent the wrong method"
|
"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;
|
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 {
|
||||||
send_socket(
|
socket
|
||||||
&mut socket,
|
.send(&ClientMethod::Error {
|
||||||
&ClientMethod::Error {
|
|
||||||
error: Cow::Borrowed(
|
error: Cow::Borrowed(
|
||||||
"Timestamp doesn't match, make sure it's in secs and is (<= 2secs)",
|
"Timestamp doesn't match, make sure it's in secs and is (<= 2secs)",
|
||||||
),
|
),
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
|
|
||||||
return Err(anyhow::anyhow!("Client tampstamp wasn't correct"));
|
return Err(anyhow::anyhow!("Client tampstamp wasn't correct"));
|
||||||
}
|
}
|
||||||
|
|
||||||
if hostname != server.config.public_hostname || !server.config.hostnames.contains(&hostname)
|
if hostname != server.config.public_hostname || !server.config.hostnames.contains(&hostname)
|
||||||
{
|
{
|
||||||
send_socket(
|
socket.send(
|
||||||
&mut socket,
|
|
||||||
&ClientMethod::Error {
|
&ClientMethod::Error {
|
||||||
error: Cow::Owned(format!("Invalid Hostname, to avoid man-in-the-middle attacks, please use the correct hostname: {}", server.config.public_hostname)),
|
error: Cow::Owned(format!("Invalid Hostname, to avoid man-in-the-middle attacks, please use the correct hostname: {}", server.config.public_hostname)),
|
||||||
},
|
},
|
||||||
@@ -66,14 +60,12 @@ impl UserConnections {
|
|||||||
return Err(anyhow::anyhow!("Client's hostname wasn't correct"));
|
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(
|
socket
|
||||||
&mut socket,
|
.send(&ClientMethod::Error {
|
||||||
&ClientMethod::Error {
|
|
||||||
error: Cow::Borrowed("Invalid public key"),
|
error: Cow::Borrowed("Invalid public key"),
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
|
|
||||||
return Err(anyhow::anyhow!("Invalid public key"));
|
return Err(anyhow::anyhow!("Invalid public key"));
|
||||||
};
|
};
|
||||||
@@ -81,45 +73,39 @@ impl UserConnections {
|
|||||||
if public_key
|
if public_key
|
||||||
.verify_strict(
|
.verify_strict(
|
||||||
format!("{timestamp}@{hostname}").as_bytes(),
|
format!("{timestamp}@{hostname}").as_bytes(),
|
||||||
&crate::signature::from_string_sig(&signature)?,
|
&crate::crypto::from_string_sig(&signature)?,
|
||||||
)
|
)
|
||||||
.is_err()
|
.is_err()
|
||||||
{
|
{
|
||||||
send_socket(
|
socket
|
||||||
&mut socket,
|
.send(&ClientMethod::Error {
|
||||||
&ClientMethod::Error {
|
|
||||||
error: Cow::Borrowed("Invalid signature"),
|
error: Cow::Borrowed("Invalid signature"),
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
|
|
||||||
return Err(anyhow::anyhow!("Invalid signature"));
|
return Err(anyhow::anyhow!("Invalid signature"));
|
||||||
}
|
}
|
||||||
|
|
||||||
{
|
{
|
||||||
send_socket(
|
socket
|
||||||
&mut socket,
|
.send(&ClientMethod::Initialized {
|
||||||
&ClientMethod::Initialized {
|
public_key: crate::crypto::to_string(&server.key.verifying_key()),
|
||||||
public_key: crate::signature::to_string(&server.key.verifying_key()),
|
signature: crate::crypto::to_string_sig(&server.key.sign(
|
||||||
signature: crate::signature::to_string_sig(&server.key.sign(
|
|
||||||
format!("{server_timestamp}@{hostname}@{public_key_string}").as_bytes(),
|
format!("{server_timestamp}@{hostname}@{public_key_string}").as_bytes(),
|
||||||
)),
|
)),
|
||||||
|
|
||||||
timestamp: server_timestamp,
|
timestamp: server_timestamp,
|
||||||
hostname,
|
hostname,
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
let Some(ServerMethod::Meta(meta)) = read_socket(&mut socket).await? else {
|
let Some(ServerMethod::Meta(meta)) = socket.read().await? else {
|
||||||
send_socket(
|
socket
|
||||||
&mut socket,
|
.send(&ClientMethod::Error {
|
||||||
&ClientMethod::Error {
|
|
||||||
error: Cow::Borrowed("Expected meta"),
|
error: Cow::Borrowed("Expected meta"),
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
|
|
||||||
return Err(anyhow::anyhow!(
|
return Err(anyhow::anyhow!(
|
||||||
"Expected meta, client called another method"
|
"Expected meta, client called another method"
|
||||||
|
|||||||
+14
-18
@@ -1,17 +1,15 @@
|
|||||||
use crate::data::messages::{MessageData, StoredMessage};
|
use crate::data::messages::{MessageData, StoredMessage};
|
||||||
use crate::protocol::{ClientMethod, send_socket};
|
use crate::protocol::ClientMethod;
|
||||||
use crate::server::Server;
|
use crate::server::Server;
|
||||||
use axum::extract::ws::WebSocket;
|
|
||||||
use ed25519_dalek::{Verifier, VerifyingKey};
|
use ed25519_dalek::{Verifier, VerifyingKey};
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::{SystemTime, UNIX_EPOCH};
|
use std::time::{SystemTime, UNIX_EPOCH};
|
||||||
use tokio::sync::Mutex;
|
|
||||||
|
|
||||||
pub async fn send_message(
|
pub async fn send_message(
|
||||||
server: &Arc<Server>,
|
server: &Arc<Server>,
|
||||||
verifying_key: VerifyingKey,
|
verifying_key: VerifyingKey,
|
||||||
_socket: &Arc<Mutex<WebSocket>>,
|
_socket: &Arc<crate::ws::EnclaveWebSocket>,
|
||||||
message: MessageData,
|
message: MessageData,
|
||||||
channel_id: String,
|
channel_id: String,
|
||||||
) -> anyhow::Result<()> {
|
) -> anyhow::Result<()> {
|
||||||
@@ -24,13 +22,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!(
|
let signed_string = format!(
|
||||||
"{}@{}@{}",
|
"{}@{}@{}",
|
||||||
message.timestamp, server_pubkey_string, message.content
|
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"))?;
|
.map_err(|_| anyhow::anyhow!("Invalid signature encoding"))?;
|
||||||
|
|
||||||
verifying_key
|
verifying_key
|
||||||
@@ -39,7 +37,7 @@ pub async fn send_message(
|
|||||||
|
|
||||||
let stored = StoredMessage {
|
let stored = StoredMessage {
|
||||||
id: uuid::Uuid::new_v4().to_string(),
|
id: uuid::Uuid::new_v4().to_string(),
|
||||||
author: crate::signature::to_string(&verifying_key),
|
author: crate::crypto::to_string(&verifying_key),
|
||||||
is_edited: false,
|
is_edited: false,
|
||||||
data: message,
|
data: message,
|
||||||
};
|
};
|
||||||
@@ -58,7 +56,7 @@ pub async fn send_message(
|
|||||||
pub async fn get_messages(
|
pub async fn get_messages(
|
||||||
server: &Arc<Server>,
|
server: &Arc<Server>,
|
||||||
_verifying_key: VerifyingKey,
|
_verifying_key: VerifyingKey,
|
||||||
socket: &Arc<Mutex<WebSocket>>,
|
socket: &Arc<crate::ws::EnclaveWebSocket>,
|
||||||
channel_id: String,
|
channel_id: String,
|
||||||
chunk: u32,
|
chunk: u32,
|
||||||
) -> anyhow::Result<()> {
|
) -> anyhow::Result<()> {
|
||||||
@@ -68,13 +66,11 @@ pub async fn get_messages(
|
|||||||
.message_store
|
.message_store
|
||||||
.get_recent_messages(&channel_id, CHUNK_SIZE, chunk)?;
|
.get_recent_messages(&channel_id, CHUNK_SIZE, chunk)?;
|
||||||
|
|
||||||
send_socket(
|
socket
|
||||||
&mut *socket.lock().await,
|
.send(&ClientMethod::Messages {
|
||||||
&ClientMethod::Messages {
|
|
||||||
messages: HashMap::from([(channel_id, messages)]),
|
messages: HashMap::from([(channel_id, messages)]),
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
@@ -92,18 +88,18 @@ pub async fn edit_message(
|
|||||||
.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"))?;
|
||||||
|
|
||||||
let author_pubkey = crate::signature::to_string(&verifying_key);
|
let author_pubkey = crate::crypto::to_string(&verifying_key);
|
||||||
if existing.author != author_pubkey {
|
if existing.author != author_pubkey {
|
||||||
anyhow::bail!("Not authorized to edit this message");
|
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!(
|
let signed_string = format!(
|
||||||
"{}@{}@{}",
|
"{}@{}@{}",
|
||||||
existing.data.timestamp, server_pubkey_string, new_content
|
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"))?;
|
.map_err(|_| anyhow::anyhow!("Invalid signature encoding"))?;
|
||||||
|
|
||||||
verifying_key
|
verifying_key
|
||||||
@@ -146,7 +142,7 @@ pub async fn delete_message(
|
|||||||
.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"))?;
|
||||||
|
|
||||||
let author_pubkey = crate::signature::to_string(&verifying_key);
|
let author_pubkey = crate::crypto::to_string(&verifying_key);
|
||||||
if existing.author != author_pubkey {
|
if existing.author != author_pubkey {
|
||||||
anyhow::bail!("Not authorized to delete this message");
|
anyhow::bail!("Not authorized to delete this message");
|
||||||
}
|
}
|
||||||
|
|||||||
+37
-53
@@ -1,15 +1,14 @@
|
|||||||
use std::{borrow::Cow, collections::HashMap, sync::Arc};
|
use std::{borrow::Cow, collections::HashMap, sync::Arc};
|
||||||
|
|
||||||
use axum::extract::ws::{Message, Utf8Bytes, WebSocket};
|
|
||||||
use ed25519_dalek::VerifyingKey;
|
use ed25519_dalek::VerifyingKey;
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use tokio::sync::Mutex;
|
|
||||||
|
|
||||||
use crate::{data::messages::StoredMessage, server::Server, types::ClientMeta};
|
use crate::{data::messages::StoredMessage, server::Server, types::ClientMeta};
|
||||||
|
|
||||||
pub mod initialize;
|
pub mod initialize;
|
||||||
pub mod message;
|
pub mod message;
|
||||||
pub mod user;
|
pub mod user;
|
||||||
|
pub mod voice;
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
#[serde(tag = "method")]
|
#[serde(tag = "method")]
|
||||||
@@ -43,6 +42,25 @@ pub enum ClientMethod {
|
|||||||
channel_id: String,
|
channel_id: String,
|
||||||
message_id: String,
|
message_id: String,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
JoinVoice {
|
||||||
|
channel_id: String,
|
||||||
|
pin: u64,
|
||||||
|
},
|
||||||
|
|
||||||
|
UserJoinedVoice {
|
||||||
|
channel_id: String,
|
||||||
|
pubkey: String,
|
||||||
|
},
|
||||||
|
|
||||||
|
UserLeftVoice {
|
||||||
|
channel_id: String,
|
||||||
|
pubkey: String,
|
||||||
|
},
|
||||||
|
|
||||||
|
Speaking {
|
||||||
|
pubkey: String,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
@@ -87,27 +105,27 @@ pub enum ServerMethod {
|
|||||||
message_id: String,
|
message_id: String,
|
||||||
channel_id: String,
|
channel_id: String,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
JoinVoice {
|
||||||
|
channel_id: String,
|
||||||
|
},
|
||||||
|
|
||||||
|
LeaveVoice,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn read_loop(
|
pub async fn read_loop(
|
||||||
server: &Arc<Server>,
|
server: &Arc<Server>,
|
||||||
verifying_key: VerifyingKey,
|
verifying_key: VerifyingKey,
|
||||||
socket: &Arc<Mutex<WebSocket>>,
|
socket: &Arc<crate::ws::EnclaveWebSocket>,
|
||||||
) -> anyhow::Result<()> {
|
) -> anyhow::Result<()> {
|
||||||
let mut socket_lock = socket.lock().await;
|
while let Some(message) = socket.read().await? {
|
||||||
|
|
||||||
while let Some(message) = read_socket(&mut *socket_lock).await? {
|
|
||||||
drop(socket_lock);
|
|
||||||
|
|
||||||
match message {
|
match message {
|
||||||
ServerMethod::Initialize { .. } => {
|
ServerMethod::Initialize { .. } => {
|
||||||
send_socket(
|
socket
|
||||||
&mut *socket.lock().await,
|
.send(&ClientMethod::Error {
|
||||||
&ClientMethod::Error {
|
|
||||||
error: Cow::Borrowed("Already initialized"),
|
error: Cow::Borrowed("Already initialized"),
|
||||||
},
|
})
|
||||||
)
|
.await?;
|
||||||
.await?;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(unused_variables)]
|
#[allow(unused_variables)]
|
||||||
@@ -152,50 +170,16 @@ pub async fn read_loop(
|
|||||||
ServerMethod::GetUsers { pubkeys } => {
|
ServerMethod::GetUsers { pubkeys } => {
|
||||||
user::get_users(server, verifying_key, socket, pubkeys).await?;
|
user::get_users(server, verifying_key, socket, pubkeys).await?;
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
socket_lock = socket.lock().await;
|
ServerMethod::JoinVoice { channel_id } => {
|
||||||
}
|
voice::join(server, verifying_key, socket, channel_id).await?;
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn read_socket(socket: &mut WebSocket) -> anyhow::Result<Option<ServerMethod>> {
|
|
||||||
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)) => {
|
ServerMethod::LeaveVoice => {
|
||||||
socket.send(Message::Pong(v)).await?;
|
voice::leave(server, verifying_key).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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,23 +1,18 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
use axum::extract::ws::WebSocket;
|
|
||||||
use ed25519_dalek::VerifyingKey;
|
use ed25519_dalek::VerifyingKey;
|
||||||
use tokio::sync::Mutex;
|
|
||||||
|
|
||||||
use crate::{
|
use crate::{protocol::ClientMethod, server::Server};
|
||||||
protocol::{ClientMethod, send_socket},
|
|
||||||
server::Server,
|
|
||||||
};
|
|
||||||
|
|
||||||
pub async fn get_users(
|
pub async fn get_users(
|
||||||
server: &Arc<Server>,
|
server: &Arc<Server>,
|
||||||
_verifying_key: VerifyingKey,
|
_verifying_key: VerifyingKey,
|
||||||
socket: &Arc<Mutex<WebSocket>>,
|
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.user_store.get_users(&pubkeys).await?;
|
||||||
|
|
||||||
send_socket(&mut *socket.lock().await, &ClientMethod::Users { users }).await?;
|
socket.send(&ClientMethod::Users { users }).await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,67 @@
|
|||||||
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
use ed25519_dalek::VerifyingKey;
|
||||||
|
|
||||||
|
use crate::{protocol::ClientMethod, server::Server};
|
||||||
|
|
||||||
|
pub async fn join(
|
||||||
|
server: &Arc<Server>,
|
||||||
|
verifying_key: VerifyingKey,
|
||||||
|
socket: &Arc<crate::ws::EnclaveWebSocket>,
|
||||||
|
channel_id: String,
|
||||||
|
) -> anyhow::Result<()> {
|
||||||
|
{
|
||||||
|
let pin = rand::random::<u64>() % (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<Server>, 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(())
|
||||||
|
}
|
||||||
+27
-10
@@ -1,5 +1,6 @@
|
|||||||
use std::{
|
use std::{
|
||||||
collections::HashMap,
|
collections::HashMap,
|
||||||
|
net::SocketAddr,
|
||||||
path::PathBuf,
|
path::PathBuf,
|
||||||
sync::{Arc, atomic::AtomicU16},
|
sync::{Arc, atomic::AtomicU16},
|
||||||
};
|
};
|
||||||
@@ -9,37 +10,54 @@ use axum::{
|
|||||||
response::Response,
|
response::Response,
|
||||||
};
|
};
|
||||||
use ed25519_dalek::{SigningKey, VerifyingKey};
|
use ed25519_dalek::{SigningKey, VerifyingKey};
|
||||||
use tokio::{sync::Mutex, task::JoinSet};
|
use tokio::{
|
||||||
|
net::UdpSocket,
|
||||||
|
sync::{Mutex, OnceCell},
|
||||||
|
task::JoinSet,
|
||||||
|
time::Instant,
|
||||||
|
};
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
data::{config::Config, messages::MessageStore, users::UserMetaStore},
|
data::{config::Config, messages::MessageStore, users::UserMetaStore},
|
||||||
protocol::{ClientMethod, read_loop, send_socket},
|
protocol::{ClientMethod, read_loop},
|
||||||
types::ClientMeta,
|
types::ClientMeta,
|
||||||
|
ws::EnclaveWebSocket,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
pub struct VoiceConnection {
|
||||||
|
pub addr: SocketAddr,
|
||||||
|
pub channel_id: String,
|
||||||
|
pub last_speaking_sent: Instant,
|
||||||
|
}
|
||||||
|
|
||||||
pub struct UserConnections {
|
pub struct UserConnections {
|
||||||
pub meta: ClientMeta,
|
pub meta: ClientMeta,
|
||||||
pub counter: AtomicU16,
|
pub counter: AtomicU16,
|
||||||
pub public_key: VerifyingKey,
|
pub public_key: VerifyingKey,
|
||||||
pub connections: Mutex<HashMap<u16, Arc<Mutex<WebSocket>>>>,
|
pub connections: Mutex<HashMap<u16, Arc<crate::ws::EnclaveWebSocket>>>,
|
||||||
|
pub voice: Mutex<Option<VoiceConnection>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct Server {
|
pub struct Server {
|
||||||
pub key: SigningKey,
|
pub key: SigningKey,
|
||||||
pub config: Config,
|
pub config: Config,
|
||||||
pub clients: Mutex<HashMap<VerifyingKey, Arc<UserConnections>>>,
|
pub clients: Mutex<HashMap<VerifyingKey, Arc<UserConnections>>>,
|
||||||
|
pub voice_pins: Mutex<HashMap<u64, (VerifyingKey, String)>>,
|
||||||
pub message_store: MessageStore,
|
pub message_store: MessageStore,
|
||||||
pub user_store: UserMetaStore,
|
pub user_store: UserMetaStore,
|
||||||
|
pub voice_socket: OnceCell<UdpSocket>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Server {
|
impl Server {
|
||||||
pub async fn new() -> anyhow::Result<Arc<Self>> {
|
pub async fn new() -> anyhow::Result<Arc<Self>> {
|
||||||
Ok(Arc::new(Self {
|
Ok(Arc::new(Self {
|
||||||
key: crate::signature::get().await?,
|
key: crate::crypto::get().await?,
|
||||||
config: Config::get().await?,
|
config: Config::get().await?,
|
||||||
clients: Mutex::new(HashMap::new()),
|
clients: Mutex::new(HashMap::new()),
|
||||||
|
voice_pins: Mutex::new(HashMap::new()),
|
||||||
message_store: MessageStore::new(PathBuf::from("messages"))?,
|
message_store: MessageStore::new(PathBuf::from("messages"))?,
|
||||||
user_store: UserMetaStore::new(PathBuf::from("users.db"))?,
|
user_store: UserMetaStore::new(PathBuf::from("users.db"))?,
|
||||||
|
voice_socket: OnceCell::new(),
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -49,18 +67,16 @@ impl Server {
|
|||||||
let s = self.clone();
|
let s = self.clone();
|
||||||
|
|
||||||
ws.on_upgrade(move |socket: WebSocket| async move {
|
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)) => {
|
Ok((client, public_key, meta)) => {
|
||||||
if let Err(e) = s
|
if let Err(e) = s
|
||||||
.user_store
|
.user_store
|
||||||
.upsert_user(&crate::signature::to_string(&public_key), &meta)
|
.upsert_user(&crate::crypto::to_string(&public_key), &meta)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
eprintln!("Failed to upsert client: {e}");
|
eprintln!("Failed to upsert client: {e}");
|
||||||
}
|
}
|
||||||
|
|
||||||
let client = Arc::new(Mutex::new(client));
|
|
||||||
|
|
||||||
let mut clients_meta = s.clients.lock().await;
|
let mut clients_meta = s.clients.lock().await;
|
||||||
|
|
||||||
let clients = clients_meta
|
let clients = clients_meta
|
||||||
@@ -71,6 +87,7 @@ impl Server {
|
|||||||
public_key: public_key,
|
public_key: public_key,
|
||||||
counter: AtomicU16::new(0),
|
counter: AtomicU16::new(0),
|
||||||
connections: Mutex::new(HashMap::new()),
|
connections: Mutex::new(HashMap::new()),
|
||||||
|
voice: Mutex::new(None),
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
.clone();
|
.clone();
|
||||||
@@ -136,7 +153,7 @@ impl Server {
|
|||||||
impl UserConnections {
|
impl UserConnections {
|
||||||
pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> {
|
pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> {
|
||||||
for (_, conn) in self.connections.lock().await.iter() {
|
for (_, conn) in self.connections.lock().await.iter() {
|
||||||
send_socket(&mut *conn.lock().await, message).await?;
|
conn.send(message).await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
@@ -144,7 +161,7 @@ impl UserConnections {
|
|||||||
|
|
||||||
pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result<bool> {
|
pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result<bool> {
|
||||||
if let Some(conn) = self.connections.lock().await.get(&id) {
|
if let Some(conn) = self.connections.lock().await.get(&id) {
|
||||||
send_socket(&mut *conn.lock().await, message).await?;
|
conn.send(message).await?;
|
||||||
|
|
||||||
Ok(true)
|
Ok(true)
|
||||||
} else {
|
} else {
|
||||||
|
|||||||
+2
-1
@@ -19,8 +19,9 @@ pub struct ServerMeta {
|
|||||||
#[serde(tag = "kind")]
|
#[serde(tag = "kind")]
|
||||||
#[serde(rename_all = "camelCase")]
|
#[serde(rename_all = "camelCase")]
|
||||||
pub enum ChannelKind {
|
pub enum ChannelKind {
|
||||||
Text,
|
|
||||||
Category { channels: Vec<Channel> },
|
Category { channels: Vec<Channel> },
|
||||||
|
Voice { max_users: u8 },
|
||||||
|
Text,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
|||||||
@@ -0,0 +1,143 @@
|
|||||||
|
use std::{net::SocketAddr, sync::Arc};
|
||||||
|
|
||||||
|
use anyhow::Context;
|
||||||
|
use ed25519_dalek::VerifyingKey;
|
||||||
|
use tokio::net::UdpSocket;
|
||||||
|
|
||||||
|
use crate::{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);
|
||||||
|
|
||||||
|
let (sender_pubkey, channel_id, payload) = {
|
||||||
|
let mut pins = self.voice_pins.lock().await;
|
||||||
|
|
||||||
|
if let Some((pubkey, channel_id)) = pins.remove(&pin) {
|
||||||
|
(pubkey, channel_id, &buf[8..len])
|
||||||
|
} else {
|
||||||
|
drop(pins);
|
||||||
|
|
||||||
|
match self.find_voice_sender(&addr).await {
|
||||||
|
Some((pubkey, channel_id)) => (pubkey, channel_id, &buf[..]),
|
||||||
|
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);
|
||||||
|
|
||||||
|
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(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;
|
||||||
|
}
|
||||||
|
|
||||||
|
let _ = self.udp_send_to(&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, addr: &SocketAddr, payload: &[u8]) -> anyhow::Result<()> {
|
||||||
|
let socket = self.get_voice_socket()?;
|
||||||
|
|
||||||
|
socket.send_to(payload, addr).await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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<SplitSink<WebSocket, Message>>,
|
||||||
|
rx: Mutex<SplitStream<WebSocket>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<Option<ServerMethod>> {
|
||||||
|
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(())
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user