Improved multi client architecture
This commit is contained in:
+13
-18
@@ -4,28 +4,29 @@ use std::{
|
|||||||
};
|
};
|
||||||
|
|
||||||
use axum::extract::ws::WebSocket;
|
use axum::extract::ws::WebSocket;
|
||||||
use ed25519_dalek::Signer;
|
use ed25519_dalek::{Signer, VerifyingKey};
|
||||||
|
|
||||||
use crate::server::Server;
|
use crate::server::Server;
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
|
use crate::server::UserConnections;
|
||||||
|
|
||||||
impl super::Client {
|
impl UserConnections {
|
||||||
pub async fn initialize(
|
pub async fn initialize(
|
||||||
server: &Arc<Server>,
|
server: &Arc<Server>,
|
||||||
mut socket: WebSocket,
|
mut socket: WebSocket,
|
||||||
) -> anyhow::Result<(Self, ClientMeta)> {
|
) -> anyhow::Result<(WebSocket, 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,
|
||||||
}) = Client::read_socket(&mut socket).await?
|
}) = read_socket(&mut socket).await?
|
||||||
else {
|
else {
|
||||||
Client::send_socket(
|
send_socket(
|
||||||
&mut socket,
|
&mut socket,
|
||||||
ClientMethod::Error {
|
&ClientMethod::Error {
|
||||||
error: Cow::Borrowed("Initialization required"),
|
error: Cow::Borrowed("Initialization required"),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -51,9 +52,9 @@ impl super::Client {
|
|||||||
{
|
{
|
||||||
let timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs();
|
let timestamp = SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs();
|
||||||
|
|
||||||
Client::send_socket(
|
send_socket(
|
||||||
&mut socket,
|
&mut socket,
|
||||||
ClientMethod::Initialized {
|
&ClientMethod::Initialized {
|
||||||
public_key: crate::signature::to_string(&server.key.verifying_key()),
|
public_key: crate::signature::to_string(&server.key.verifying_key()),
|
||||||
signature: server
|
signature: server
|
||||||
.key
|
.key
|
||||||
@@ -67,10 +68,10 @@ impl super::Client {
|
|||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
let Some(ServerMethod::Meta(meta)) = Client::read_socket(&mut socket).await? else {
|
let Some(ServerMethod::Meta(meta)) = read_socket(&mut socket).await? else {
|
||||||
Client::send_socket(
|
send_socket(
|
||||||
&mut socket,
|
&mut socket,
|
||||||
ClientMethod::Error {
|
&ClientMethod::Error {
|
||||||
error: Cow::Borrowed("Expected meta"),
|
error: Cow::Borrowed("Expected meta"),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -81,12 +82,6 @@ impl super::Client {
|
|||||||
));
|
));
|
||||||
};
|
};
|
||||||
|
|
||||||
Ok((
|
Ok((socket, public_key, meta))
|
||||||
Self {
|
|
||||||
socket: Arc::new(Mutex::new(socket)),
|
|
||||||
public_key,
|
|
||||||
},
|
|
||||||
meta,
|
|
||||||
))
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+12
-25
@@ -1,18 +1,11 @@
|
|||||||
use std::{borrow::Cow, sync::Arc};
|
use std::{borrow::Cow, sync::Arc};
|
||||||
|
|
||||||
use axum::extract::ws::{Message, Utf8Bytes, WebSocket};
|
use axum::extract::ws::{Message, Utf8Bytes, WebSocket};
|
||||||
use ed25519_dalek::VerifyingKey;
|
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use tokio::sync::Mutex;
|
use tokio::sync::Mutex;
|
||||||
|
|
||||||
pub mod initialize;
|
pub mod initialize;
|
||||||
|
|
||||||
#[derive(Clone)]
|
|
||||||
pub struct Client {
|
|
||||||
pub socket: Arc<Mutex<WebSocket>>,
|
|
||||||
pub public_key: VerifyingKey,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
pub struct ClientMeta {}
|
pub struct ClientMeta {}
|
||||||
|
|
||||||
@@ -56,17 +49,20 @@ pub enum ServerMethod {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Client {
|
pub async fn read_loop(socket: &Arc<Mutex<WebSocket>>) -> anyhow::Result<()> {
|
||||||
pub async fn read_loop(&mut self) -> anyhow::Result<()> {
|
while let Some(message) = read_socket(&mut *socket.lock().await).await? {
|
||||||
while let Some(message) = self.read().await? {
|
|
||||||
match message {
|
match message {
|
||||||
ServerMethod::Initialize { .. } => {
|
ServerMethod::Initialize { .. } => {
|
||||||
self.send(ClientMethod::Error {
|
send_socket(
|
||||||
|
&mut *socket.lock().await,
|
||||||
|
&ClientMethod::Error {
|
||||||
error: Cow::Borrowed("Already initialized"),
|
error: Cow::Borrowed("Already initialized"),
|
||||||
})
|
},
|
||||||
|
)
|
||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[allow(unused_variables)]
|
||||||
ServerMethod::Meta(meta) => {}
|
ServerMethod::Meta(meta) => {}
|
||||||
|
|
||||||
ServerMethod::Error { error } => {
|
ServerMethod::Error { error } => {
|
||||||
@@ -84,9 +80,9 @@ impl Client {
|
|||||||
if let Ok(msg) = serde_json::from_str(&text.to_string()) {
|
if let Ok(msg) = serde_json::from_str(&text.to_string()) {
|
||||||
Ok(Some(msg))
|
Ok(Some(msg))
|
||||||
} else {
|
} else {
|
||||||
Client::send_socket(
|
send_socket(
|
||||||
socket,
|
socket,
|
||||||
ClientMethod::Error {
|
&ClientMethod::Error {
|
||||||
error: Cow::Borrowed("Unable to parse message: {text}"),
|
error: Cow::Borrowed("Unable to parse message: {text}"),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
@@ -108,21 +104,12 @@ impl Client {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn send_socket(socket: &mut WebSocket, message: ClientMethod) -> anyhow::Result<()> {
|
pub async fn send_socket(socket: &mut WebSocket, message: &ClientMethod) -> anyhow::Result<()> {
|
||||||
socket
|
socket
|
||||||
.send(Message::Text(Utf8Bytes::from(serde_json::to_string(
|
.send(Message::Text(Utf8Bytes::from(serde_json::to_string(
|
||||||
&message,
|
message,
|
||||||
)?)))
|
)?)))
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn read(&mut self) -> anyhow::Result<Option<ServerMethod>> {
|
|
||||||
Self::read_socket(&mut *self.socket.lock().await).await
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn send(&mut self, message: ClientMethod) -> anyhow::Result<()> {
|
|
||||||
Self::send_socket(&mut *self.socket.lock().await, message).await
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
+33
-9
@@ -12,19 +12,20 @@ use tokio::sync::Mutex;
|
|||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
config::Config,
|
config::Config,
|
||||||
protocol::{Client, ClientMeta},
|
protocol::{ClientMeta, ClientMethod, read_loop, send_socket},
|
||||||
};
|
};
|
||||||
|
|
||||||
pub struct OnlineClientMeta {
|
pub struct UserConnections {
|
||||||
pub meta: ClientMeta,
|
pub meta: ClientMeta,
|
||||||
pub counter: AtomicU16,
|
pub counter: AtomicU16,
|
||||||
pub connections: HashMap<u16, Client>,
|
pub public_key: VerifyingKey,
|
||||||
|
pub connections: HashMap<u16, Arc<Mutex<WebSocket>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct Server {
|
pub struct Server {
|
||||||
pub key: SigningKey,
|
pub key: SigningKey,
|
||||||
pub config: Config,
|
pub config: Config,
|
||||||
pub clients: Mutex<HashMap<VerifyingKey, OnlineClientMeta>>,
|
pub clients: Mutex<HashMap<VerifyingKey, UserConnections>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Server {
|
impl Server {
|
||||||
@@ -42,15 +43,18 @@ 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 Client::initialize(&s, socket).await {
|
match UserConnections::initialize(&s, socket).await {
|
||||||
Ok((mut client, meta)) => {
|
Ok((client, public_key, meta)) => {
|
||||||
|
let client = Arc::new(Mutex::new(client));
|
||||||
|
|
||||||
let mut clients_meta = s.clients.lock().await;
|
let mut clients_meta = s.clients.lock().await;
|
||||||
|
|
||||||
let client_meta =
|
let client_meta =
|
||||||
clients_meta
|
clients_meta
|
||||||
.entry(client.public_key)
|
.entry(public_key)
|
||||||
.or_insert_with(|| OnlineClientMeta {
|
.or_insert_with(|| UserConnections {
|
||||||
meta,
|
meta,
|
||||||
|
public_key: public_key,
|
||||||
counter: AtomicU16::new(0),
|
counter: AtomicU16::new(0),
|
||||||
connections: HashMap::new(),
|
connections: HashMap::new(),
|
||||||
});
|
});
|
||||||
@@ -62,7 +66,7 @@ impl Server {
|
|||||||
client.clone(),
|
client.clone(),
|
||||||
);
|
);
|
||||||
|
|
||||||
if let Err(e) = client.read_loop().await {
|
if let Err(e) = read_loop(&client).await {
|
||||||
eprintln!("Failed to handle client: {e}");
|
eprintln!("Failed to handle client: {e}");
|
||||||
} else {
|
} else {
|
||||||
println!("Client connection closed")
|
println!("Client connection closed")
|
||||||
@@ -76,3 +80,23 @@ impl Server {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl UserConnections {
|
||||||
|
pub async fn send(&self, message: &ClientMethod) -> anyhow::Result<()> {
|
||||||
|
for (_, conn) in &self.connections {
|
||||||
|
send_socket(&mut *conn.lock().await, message).await?;
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn send_to(&self, id: u16, message: &ClientMethod) -> anyhow::Result<bool> {
|
||||||
|
if let Some(conn) = self.connections.get(&id) {
|
||||||
|
send_socket(&mut *conn.lock().await, message).await?;
|
||||||
|
|
||||||
|
Ok(true)
|
||||||
|
} else {
|
||||||
|
Ok(false)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user