From 14b3d12d0342574199b4c6127ae197aa213e024c Mon Sep 17 00:00:00 2001 From: Claas Date: Sun, 21 Dec 2025 17:24:53 +0100 Subject: [PATCH] Integrate actors --- .../src/actor/{user.rs => client.rs} | 10 +-- delivery-service/src/actor/mod.rs | 2 +- delivery-service/src/actor/switchboard.rs | 72 ++++++++++--------- delivery-service/src/main.rs | 38 ++++++---- 4 files changed, 68 insertions(+), 54 deletions(-) rename delivery-service/src/actor/{user.rs => client.rs} (94%) diff --git a/delivery-service/src/actor/user.rs b/delivery-service/src/actor/client.rs similarity index 94% rename from delivery-service/src/actor/user.rs rename to delivery-service/src/actor/client.rs index 392fff8..e1ddd33 100644 --- a/delivery-service/src/actor/user.rs +++ b/delivery-service/src/actor/client.rs @@ -9,13 +9,13 @@ enum Message { AddConnection(web_socket::Handle), } -struct User { +struct Client { id: Arc, receiver: mpsc::Receiver, sockets: HashMap, web_socket::Handle>, } -impl User { +impl Client { fn new(id: Arc, receiver: mpsc::Receiver) -> Self { Self { id, @@ -45,9 +45,9 @@ impl User { self.sockets.remove(&id); } - // If there are no more sockets, remove the user + // If there are no more sockets, remove the client if self.sockets.is_empty() { - tracing::debug!("[{}] No more sockets. Stopping user actor", self.id); + tracing::debug!("[{}] No more sockets. Stopping client actor", self.id); return; } } @@ -84,7 +84,7 @@ impl Handle { pub(in crate::actor) fn new(id: Arc) -> Self { let (sender, receiver) = mpsc::channel(8); - let actor = User::new(id, receiver); + let actor = Client::new(id, receiver); let id = actor.id.clone(); tokio::spawn(actor.run()); diff --git a/delivery-service/src/actor/mod.rs b/delivery-service/src/actor/mod.rs index a901a9c..bd8187e 100644 --- a/delivery-service/src/actor/mod.rs +++ b/delivery-service/src/actor/mod.rs @@ -1,3 +1,3 @@ +mod client; pub(crate) mod switchboard; -mod user; mod web_socket; diff --git a/delivery-service/src/actor/switchboard.rs b/delivery-service/src/actor/switchboard.rs index 286ccbd..f6521a3 100644 --- a/delivery-service/src/actor/switchboard.rs +++ b/delivery-service/src/actor/switchboard.rs @@ -3,81 +3,82 @@ use std::{collections::HashMap, sync::Arc}; use axum::extract::ws; use tokio::sync::{mpsc, oneshot}; -use crate::actor::{user, web_socket}; +use crate::actor::{client, web_socket}; -struct UserNotFound; +struct ClientNotFound; enum Message { Send { - user_id: Arc, + client_id: Arc, message: Arc<[u8]>, - response: oneshot::Sender>, + response: oneshot::Sender>, }, Connect { - user_id: Arc, + client_id: Arc, web_socket: ws::WebSocket, }, } struct Switchboard { receiver: mpsc::Receiver, - users: HashMap, user::Handle>, + clients: HashMap, client::Handle>, } impl Switchboard { fn new(receiver: mpsc::Receiver) -> Self { Self { receiver, - users: HashMap::new(), + clients: HashMap::new(), } } async fn handle_message(&mut self, message: Message) { match message { Message::Send { - user_id, + client_id, message, response, } => { - //TODO might need to store the message for the user to receive later - let Some(user) = self.users.get_mut(&user_id) else { + //TODO might need to store the message for the client to receive later + let Some(client) = self.clients.get_mut(&client_id) else { // We don't care if they stopped waiting for the response - _ = response.send(Err(UserNotFound)); + _ = response.send(Err(ClientNotFound)); return; }; - let Err(user::HandleError::Closed) = user.send_message(message).await else { + let Err(client::HandleError::Closed) = client.send_message(message).await else { _ = response.send(Ok(())); return; }; - // TODO user is not active. Store message for the user to receive later or send push notification + // TODO client is not active. Store message for the client to receive later or send push notification _ = response.send(Ok(())); } Message::Connect { - user_id, + client_id, web_socket, } => { let socket = web_socket::Handle::new(web_socket); - let user = self - .users - .entry(user_id.clone()) - .or_insert_with(|| user::Handle::new(user_id.clone())); + let client = self + .clients + .entry(client_id.clone()) + .or_insert_with(|| client::Handle::new(client_id.clone())); - let Err(user::HandleError::Closed) = user.add_socket(socket.clone()).await else { + let Err(client::HandleError::Closed) = client.add_socket(socket.clone()).await + else { return; }; tracing::debug!( - "[Switchboard] User actor {} is dead. Rebirthing user.", - user_id + "[Switchboard] Client actor {} is dead. Rebirthing client.", + client_id ); - let user = user.rebirth(); - self.users.insert(user_id, user.clone()); + let client = client.rebirth(); + self.clients.insert(client_id, client.clone()); - if let Err(user::HandleError::Closed) = user.add_socket(socket).await { - tracing::error!("[Switchboard] Just rebirthed user actor is already dead") + if let Err(client::HandleError::Closed) = client.add_socket(socket).await { + tracing::error!("[Switchboard] Just rebirthed client actor is already dead") } } } @@ -103,16 +104,17 @@ pub(crate) struct Handle { pub(crate) enum SendMessageError { Closed, - UserNotFound, + ClientNotFound, } +#[derive(Debug)] pub(crate) enum ConnectError { Closed, } -impl From for SendMessageError { - fn from(_: UserNotFound) -> Self { - Self::UserNotFound +impl From for SendMessageError { + fn from(_: ClientNotFound) -> Self { + Self::ClientNotFound } } @@ -131,15 +133,15 @@ impl Handle { Self { sender } } - pub(in crate::actor) async fn send_message( + pub(crate) async fn send_message( &self, - user_id: Arc, + client_id: Arc, message: Arc<[u8]>, ) -> Result<(), SendMessageError> { let (sender, receiver) = oneshot::channel(); self.sender .send(Message::Send { - user_id, + client_id, message, response: sender, }) @@ -150,14 +152,14 @@ impl Handle { Ok(()) } - pub async fn add_connection( + pub(crate) async fn add_connection( &self, - user_id: Arc, + client_id: Arc, web_socket: ws::WebSocket, ) -> Result<(), ConnectError> { self.sender .send(Message::Connect { - user_id, + client_id, web_socket, }) .await diff --git a/delivery-service/src/main.rs b/delivery-service/src/main.rs index 181fadc..0ea46ca 100644 --- a/delivery-service/src/main.rs +++ b/delivery-service/src/main.rs @@ -24,7 +24,7 @@ use tower_http::{ }; use tracing::Level; -use crate::actor::switchboard; +use crate::actor::switchboard::{self, SendMessageError}; mod actor; mod extractor; @@ -130,24 +130,26 @@ async fn create_message( Path(to): Path>, bytes: Bytes, ) -> impl IntoResponse { - let channels = state.channels.lock().await; //TODO think about not leaking if they exist or not //TODO think about leaking data through timings - debug!("New message for client {}", to); - let Some(sender) = channels.get(to.as_ref()) else { - error!("Client {} not found", to); - return StatusCode::NOT_FOUND; - }; - - if let Err(error) = sender.send(bytes).await { - tracing::error!("Error sending message {:?}", error); + let data = bytes.as_ref(); + let data = Arc::from(data); + match state.switchboard.send_message(to, data).await { + Ok(()) => StatusCode::CREATED, + Err(SendMessageError::ClientNotFound) => StatusCode::NOT_FOUND, + Err(SendMessageError::Closed) => { + tracing::error!("Switchboard is closed unexpectedly"); + StatusCode::INTERNAL_SERVER_ERROR + } } - - StatusCode::CREATED } #[tracing::instrument(skip(socket, state))] -async fn handle_socket(mut socket: WebSocket, State(state): State, client_id: Arc) { +async fn handle_socket_old( + mut socket: WebSocket, + State(state): State, + client_id: Arc, +) { debug!("[{}] connected", client_id); let _droppochino = Droppochino(client_id.to_string()); @@ -239,3 +241,13 @@ async fn subscribe_messages( ) -> impl IntoResponse { websocket.on_upgrade(|socket| handle_socket(socket, state, client_id)) } + +async fn handle_socket(web_socket: WebSocket, State(state): State, client_id: Arc) { + if let Err(error) = state + .switchboard + .add_connection(client_id, web_socket) + .await + { + tracing::error!("Error adding connection to switchboard {:?}", error) + }; +} -- 2.51.2