From 1127bac9b76885988fde25042ffc028e1b53131c Mon Sep 17 00:00:00 2001 From: Claas Date: Sun, 21 Dec 2025 17:49:39 +0100 Subject: [PATCH] Improve logging --- delivery-service/src/actor/client.rs | 33 ++++++++++++++++++----- delivery-service/src/actor/switchboard.rs | 14 +++++++++- delivery-service/src/actor/web_socket.rs | 14 ++++++---- 3 files changed, 48 insertions(+), 13 deletions(-) diff --git a/delivery-service/src/actor/client.rs b/delivery-service/src/actor/client.rs index e1ddd33..b57b520 100644 --- a/delivery-service/src/actor/client.rs +++ b/delivery-service/src/actor/client.rs @@ -1,4 +1,4 @@ -use std::{collections::HashMap, sync::Arc}; +use std::{collections::HashMap, ops::ControlFlow, sync::Arc}; use tokio::{sync::mpsc, task::JoinSet}; @@ -24,11 +24,16 @@ impl Client { } } - async fn handle_message(&mut self, message: Message) { + async fn handle_message(&mut self, message: Message) -> ControlFlow<()> { match message { Message::Send(data) => { let mut sends = JoinSet::new(); for socket in self.sockets.values().cloned() { + tracing::debug!( + "[Clients/{}] Sending message to socket {}", + self.id, + socket.id + ); let data = data.clone(); sends .spawn(async move { (socket.id.clone(), socket.send_message(data).await) }); @@ -40,19 +45,29 @@ impl Client { continue; }; - tracing::debug!("[{}] Removing closed websocket receiver {}", self.id, id); + tracing::debug!( + "[Clients/{}] Removing closed websocket receiver {}", + self.id, + id + ); self.sockets.remove(&id); } // If there are no more sockets, remove the client if self.sockets.is_empty() { - tracing::debug!("[{}] No more sockets. Stopping client actor", self.id); - return; + tracing::debug!( + "[Clients/{}] No more sockets. Stopping client actor", + self.id + ); + return ControlFlow::Break(()); } + + ControlFlow::Continue(()) } Message::AddConnection(handle) => { self.sockets.insert(handle.id.clone(), handle); + ControlFlow::Continue(()) } } } @@ -60,9 +75,13 @@ impl Client { async fn run(mut self) { loop { match self.receiver.recv().await { - Some(message) => self.handle_message(message).await, + Some(message) => { + if self.handle_message(message).await.is_break() { + return; + } + } None => { - tracing::debug!("[{}] Actor send channel closed", self.id); + tracing::debug!("[Clients/{}] Actor send channel closed", self.id); return; } } diff --git a/delivery-service/src/actor/switchboard.rs b/delivery-service/src/actor/switchboard.rs index f6521a3..edd1e32 100644 --- a/delivery-service/src/actor/switchboard.rs +++ b/delivery-service/src/actor/switchboard.rs @@ -67,6 +67,11 @@ impl Switchboard { let Err(client::HandleError::Closed) = client.add_socket(socket.clone()).await else { + tracing::debug!( + "[Switchboard] Client {} added socket {}", + socket.id, + client_id + ); return; }; @@ -75,11 +80,18 @@ impl Switchboard { client_id ); let client = client.rebirth(); - self.clients.insert(client_id, client.clone()); + self.clients.insert(client_id.clone(), client.clone()); + let socket_id = socket.id.clone(); if let Err(client::HandleError::Closed) = client.add_socket(socket).await { tracing::error!("[Switchboard] Just rebirthed client actor is already dead") } + + tracing::debug!( + "[Switchboard] Client {} added socket {}", + client_id, + socket_id + ); } } } diff --git a/delivery-service/src/actor/web_socket.rs b/delivery-service/src/actor/web_socket.rs index e087f43..b121d7b 100644 --- a/delivery-service/src/actor/web_socket.rs +++ b/delivery-service/src/actor/web_socket.rs @@ -33,11 +33,15 @@ impl WebSocket { fn handle_socket_message(&mut self, message: ws::Message) -> ControlFlow<()> { match message { ws::Message::Close(_) => { - tracing::debug!("[{}] Websocket closed", self.id); + tracing::debug!("[WebSockets/{}] Websocket closed", self.id); ControlFlow::Break(()) } other => { - tracing::debug!("[{}] Unexpected Websocket message: {:?}", self.id, other); + tracing::debug!( + "[WebSockets/{}] Unexpected Websocket message: {:?}", + self.id, + other + ); ControlFlow::Continue(()) } } @@ -49,18 +53,18 @@ impl WebSocket { message = self.receiver.recv() => match message { Some(message) => self.handle_message(message).await?, None => { - tracing::debug!("[{}] Actor send channel closed", self.id); + tracing::debug!("[WebSockets/{}] Actor send channel closed", self.id); return Ok(()); }, }, message = self.socket.recv() => match message { Some(Ok(message)) => if self.handle_socket_message(message).is_break() { return Ok(()); }, Some(Err(error)) => { - tracing::debug!("[{}] websocket error: {}", self.id, error); + tracing::debug!("[WebSockets/{}] websocket error: {}", self.id, error); return Err(error.into()); }, None => { - tracing::debug!("[{}] client closed websocket", self.id); + tracing::debug!("[WebSockets/{}] client closed websocket", self.id); return Ok(()); } }, -- 2.51.2