From a7cc3cc29012b99e34cf1ce3b095c91a04ec1af8 Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Sun, 1 Feb 2026 16:51:59 +0100 Subject: [PATCH] [QOL] Modified Logger --- .../anonymous_client_connection.rs | 10 +- src/calls/call_util.rs | 7 +- src/data/communication.rs | 2 + src/main.rs | 66 ++++++++++--- src/omega/omega_connection.rs | 28 +++--- src/rho/client_connection.rs | 28 +++++- src/rho/iota_connection.rs | 58 ++++++++--- src/rho/rho_manager.rs | 2 +- src/util/logger.rs | 97 ++----------------- 9 files changed, 165 insertions(+), 133 deletions(-) diff --git a/src/anonymous_clients/anonymous_client_connection.rs b/src/anonymous_clients/anonymous_client_connection.rs index 6e7d6ac..b91776b 100644 --- a/src/anonymous_clients/anonymous_client_connection.rs +++ b/src/anonymous_clients/anonymous_client_connection.rs @@ -87,6 +87,7 @@ impl AnonymousClientConnection { .await { log_out!( + self.get_user_id().await, PrintType::Client, "Failed to send message to anonymous client: {}", e, @@ -98,13 +99,19 @@ impl AnonymousClientConnection { pub async fn send_message(self: Arc, cv: &CommunicationValue) { if !*self.is_open.read().await { log_out!( + self.get_user_id().await, PrintType::Client, "Attempted to send message to a closed connection." ); return; } if !cv.is_type(CommunicationType::pong) { - log_out!(PrintType::Client, "{}", &cv.to_json().to_string()); + log_out!( + self.get_user_id().await, + PrintType::Client, + "{}", + &cv.to_json().to_string() + ); } self.send_message_str(&cv.to_json().to_string()).await; } @@ -118,6 +125,7 @@ impl AnonymousClientConnection { return; } log_in!( + self.get_user_id().await, PrintType::Client, "Anonymous: {}", &cv.to_json().to_string() diff --git a/src/calls/call_util.rs b/src/calls/call_util.rs index 4ab5f0c..796fcdb 100644 --- a/src/calls/call_util.rs +++ b/src/calls/call_util.rs @@ -14,21 +14,21 @@ pub fn get_livekit() -> Result<(String, String, String), ()> { let hostname = match env::var("LIVEKI_HOSTNAME") { Ok(secret) => secret, Err(_) => { - log_err!(PrintType::General, "LIVEKI_HOSTNAME not set!"); + log_err!(0, PrintType::General, "LIVEKI_HOSTNAME not set!"); return Err(()); } }; let api_key = match env::var("LIVEKIT_API_KEY") { Ok(key) => key, Err(_) => { - log_err!(PrintType::General, "LIVEKIT_API_KEY not set!"); + log_err!(0, PrintType::General, "LIVEKIT_API_KEY not set!"); return Err(()); } }; let api_secret = match env::var("LIVEKIT_API_SECRET") { Ok(secret) => secret, Err(_) => { - log_err!(PrintType::General, "LIVEKIT_API_SECRET not set!"); + log_err!(0, PrintType::General, "LIVEKIT_API_SECRET not set!"); return Err(()); } }; @@ -139,6 +139,7 @@ pub async fn clean_calls(room_service: RoomClient) { let size_post = CALL_GROUPS.len(); if size_pre - size_post != 0 { log!( + 0, PrintType::Call, "Cleaned {} calls, {} remaining", size_pre - size_post, diff --git a/src/data/communication.rs b/src/data/communication.rs index 330f8d0..f6b2fe3 100644 --- a/src/data/communication.rs +++ b/src/data/communication.rs @@ -53,6 +53,7 @@ pub enum DataTypes { signature, signed, message, + message_state, last_ping, ping_iota, ping_clients, @@ -135,6 +136,7 @@ pub enum CommunicationType { settings_load, settings_list, message, + message_state, message_send, message_live, message_other_iota, diff --git a/src/main.rs b/src/main.rs index ac99074..e77d70e 100644 --- a/src/main.rs +++ b/src/main.rs @@ -51,6 +51,7 @@ async fn main() { let listener = TcpListener::bind(&address).await.unwrap(); log!( + 0, PrintType::General, "WebSocket server listening on {}", address, @@ -69,13 +70,13 @@ async fn main() { let ws_stream = match accept_hdr_async(stream.compat(), callback).await { Ok(ws) => ws, Err(e) => { - log!(PrintType::General, "WebSocket upgrade failed: {}", e,); + log!(0, PrintType::General, "WebSocket upgrade failed: {}", e,); return; } }; let (sender, receiver) = ws_stream.split(); if path == "/ws/client/" { - log_in!(PrintType::Client, "New Client connection"); + log_in!(0, PrintType::Client, "New Client connection"); let client_conn: Arc = Arc::from(ClientConnection::new(sender, receiver)); loop { @@ -90,25 +91,38 @@ async fn main() { let text = msg.into_text().unwrap(); client_conn.clone().handle_message(text).await; } else if msg.is_close() { - log_in!(PrintType::Client, "Client disconnected"); + log_in!( + client_conn.get_user_id().await, + PrintType::Client, + "Client disconnected" + ); client_conn.handle_close().await; return; } } Some(Err(e)) => { - log_err!(PrintType::Client, "WebSocket error: {}", e); + log_err!( + client_conn.get_user_id().await, + PrintType::Client, + "WebSocket error: {}", + e + ); client_conn.handle_close().await; return; } _ => { - log_in!(PrintType::Client, "Client stream ended"); + log_in!( + client_conn.get_user_id().await, + PrintType::Client, + "Client stream ended" + ); client_conn.handle_close().await; return; } } } } else if path == "/ws/anonymous_client/" { - log_in!(PrintType::Client, "New Anonymous Client connection"); + log_in!(0, PrintType::Client, "New Anonymous Client connection"); let client_conn: Arc = Arc::from(AnonymousClientConnection::new(sender, receiver)); anonymous_manager::add_anonymous_user(client_conn.clone()).await; @@ -124,7 +138,11 @@ async fn main() { let text = msg.into_text().unwrap(); client_conn.clone().handle_message(text).await; } else if msg.is_close() { - log_in!(PrintType::Client, "Anonymous Client disconnected"); + log_in!( + client_conn.get_user_id().await, + PrintType::Client, + "Anonymous Client disconnected" + ); anonymous_manager::remove_anonymous_user( client_conn.get_user_id().await, ) @@ -134,7 +152,12 @@ async fn main() { } } Some(Err(e)) => { - log_err!(PrintType::Client, "WebSocket error: {}", e); + log_err!( + client_conn.get_user_id().await, + PrintType::Client, + "WebSocket error: {}", + e + ); anonymous_manager::remove_anonymous_user( client_conn.get_user_id().await, ) @@ -143,7 +166,11 @@ async fn main() { return; } _ => { - log_in!(PrintType::Client, "Anonymous Client stream ended"); + log_in!( + client_conn.get_user_id().await, + PrintType::Client, + "Anonymous Client stream ended" + ); anonymous_manager::remove_anonymous_user( client_conn.get_user_id().await, ) @@ -154,7 +181,7 @@ async fn main() { } } } else if path == "/ws/iota/" { - log_in!(PrintType::Iota, "New Iota connection"); + log_in!(0, PrintType::Iota, "New Iota connection"); let iota_conn: Arc = Arc::from(IotaConnection::new(sender, receiver)); loop { @@ -169,19 +196,32 @@ async fn main() { let text = msg.into_text().unwrap(); iota_conn.clone().handle_message(text).await; } else if msg.is_close() { - log_in!(PrintType::Iota, "Iota disconnected"); + log_in!( + iota_conn.get_iota_id().await, + PrintType::Iota, + "Iota disconnected" + ); iota_conn.handle_close().await; return; } } Some(Err(e)) => { - log_err!(PrintType::Iota, "WebSocket error: {}", e); + log_err!( + iota_conn.get_iota_id().await, + PrintType::Iota, + "WebSocket error: {}", + e + ); iota_conn.handle_close().await; return; } _ => { // Stream ended - log_in!(PrintType::Iota, "Iota stream ended"); + log_in!( + iota_conn.get_iota_id().await, + PrintType::Iota, + "Iota stream ended" + ); iota_conn.handle_close().await; return; } diff --git a/src/omega/omega_connection.rs b/src/omega/omega_connection.rs index bc30563..a1c5a7d 100755 --- a/src/omega/omega_connection.rs +++ b/src/omega/omega_connection.rs @@ -95,7 +95,11 @@ impl OmegaConnection { async fn connect_internal(self: Arc, mut retry: usize) { loop { if retry > 5 { - log_err!(PrintType::Omega, "Max retry attempts reached, giving up."); + log_err!( + 0, + PrintType::Omega, + "Max retry attempts reached, giving up." + ); return; } @@ -104,7 +108,7 @@ impl OmegaConnection { match connect_async(&url_str).await { Ok((ws_stream, _)) => { *self.is_connected.write().await = true; - log_in!(PrintType::Omega, "WebSocket connected to {}", url_str); + log_in!(0, PrintType::Omega, "WebSocket connected to {}", url_str); retry = 0; let (write, read) = ws_stream.split(); *self.read.write().await = Some(read); @@ -135,7 +139,7 @@ impl OmegaConnection { id, Box::new(|selfc, cv| { if cv.is_type(CommunicationType::error_not_found) { - log_err!( + log_err!(0, PrintType::Omega, "Identification failed: Omikron ID not found on Omega.", ); @@ -186,7 +190,7 @@ impl OmegaConnection { if !final_cv .is_type(CommunicationType::identification_response) { - log_err!( + log_err!(0, PrintType::Omega, "Expected identification_response, got something else.", ); @@ -195,11 +199,11 @@ impl OmegaConnection { if let Some(accepted) = final_cv.get_data(DataTypes::accepted).and_then(|v| v.as_bool()) { if !accepted { - log_err!(PrintType::Omega, "Omega did not accept identification."); + log_err!(0, PrintType::Omega, "Omega did not accept identification."); return false; } } else { - log_err!(PrintType::Omega, "Omega response did not contain 'accepted' field."); + log_err!(0, PrintType::Omega, "Omega response did not contain 'accepted' field."); return false; } @@ -227,7 +231,7 @@ impl OmegaConnection { selfc.send_message(&sync_msg).await; }); - log!( + log!(0, PrintType::Omega, "Successfully identified with Omega.", ); @@ -241,7 +245,7 @@ impl OmegaConnection { }; if let Err(e) = task.await { - log_err!(PrintType::Omega, "{}", &e); + log_err!(0, PrintType::Omega, "{}", &e); } }); @@ -270,12 +274,13 @@ impl OmegaConnection { } *self.read.write().await = None; *self.write.write().await = None; - log_err!(PrintType::Omega, "Connection lost. Retrying..."); + log_err!(0, PrintType::Omega, "Connection lost. Retrying..."); retry += 1; sleep(Duration::from_secs(2)).await; } Err(e) => { log_err!( + 0, PrintType::Omega, "WebSocket connection failed (attempt {}): {}", retry + 1, @@ -308,7 +313,7 @@ impl OmegaConnection { continue; } let msg_id = cv.get_id(); - log_in!(PrintType::Omega, "{}", &cv.to_json().to_string()); + log_in!(0, PrintType::Omega, "{}", &cv.to_json().to_string()); // Handle waiting tasks if let Some(task) = WAITING_TASKS.remove(&msg_id) { if (task.1)(self.clone(), cv.clone()) { @@ -333,7 +338,7 @@ impl OmegaConnection { let mut guard = self.write.write().await; if let Some(ws) = guard.as_mut() { if !cv.is_type(CommunicationType::ping) { - log_out!(PrintType::Omega, "{}", &cv.to_json().to_string()); + log_out!(0, PrintType::Omega, "{}", &cv.to_json().to_string()); } let _ = ws .send(Message::Text(cv.to_json().to_string().into())) @@ -438,6 +443,7 @@ impl OmegaConnection { tokio::spawn(async move { if let Err(e) = inner_tx.send(response_cv).await { log_err!( + 0, PrintType::Omega, "Failed to send response back to awaiter: {}", e diff --git a/src/rho/client_connection.rs b/src/rho/client_connection.rs index c384bec..7112e5e 100644 --- a/src/rho/client_connection.rs +++ b/src/rho/client_connection.rs @@ -90,7 +90,12 @@ impl ClientConnection { .send(Message::Text(Utf8Bytes::from(message.to_string()))) .await { - log_out!(PrintType::Client, "Failed to send message to client: {}", e,); + log_out!( + self.get_user_id().await, + PrintType::Client, + "Failed to send message to client: {}", + e, + ); } } @@ -98,13 +103,19 @@ impl ClientConnection { pub async fn send_message(self: Arc, cv: &CommunicationValue) { if !*self.is_open.read().await { log_out!( + self.get_user_id().await, PrintType::Client, "Attempted to send message to a closed connection." ); return; } if !cv.is_type(CommunicationType::pong) { - log_out!(PrintType::Client, "{}", &cv.to_json().to_string()); + log_out!( + self.get_user_id().await, + PrintType::Client, + "{}", + &cv.to_json().to_string() + ); } self.send_message_str(&cv.to_json().to_string()).await; } @@ -117,7 +128,12 @@ impl ClientConnection { self.handle_ping(cv).await; return; } - log_in!(PrintType::Client, "{}", &cv.to_json().to_string()); + log_in!( + self.get_user_id().await, + PrintType::Client, + "{}", + &cv.to_json().to_string() + ); let identified = *self.identified.read().await; let challenged = *self.challenged.read().await; @@ -128,7 +144,11 @@ impl ClientConnection { .and_then(|v| v.as_i64()) .unwrap_or(0); if user_id == 0 { - log_out!(PrintType::Client, "Invalid USER ID"); + log_out!( + self.get_user_id().await, + PrintType::Client, + "Invalid USER ID" + ); self.clone() .send_error_response(&cv.get_id(), CommunicationType::error_invalid_data) .await; diff --git a/src/rho/iota_connection.rs b/src/rho/iota_connection.rs index 5f96b5c..36e9fe5 100755 --- a/src/rho/iota_connection.rs +++ b/src/rho/iota_connection.rs @@ -49,7 +49,7 @@ pub struct IotaConnection { pub ping: Arc>, pub_key: Arc>>>, pub waiting_tasks: - DashMap, CommunicationValue) -> bool + Send + Sync>>, + DashMap, CommunicationValue) -> bool + Send + Sync>>, pub rho_connection: Arc>>>, } @@ -125,14 +125,24 @@ impl IotaConnection { .send(Message::Text(Utf8Bytes::from(message.to_string()))) .await { - log_err!(PrintType::Iota, "Failed to send WebSocket message: {:?}", e,); + log_err!( + self.get_iota_id().await, + PrintType::Iota, + "Failed to send WebSocket message: {:?}", + e, + ); } } /// Send a CommunicationValue to the Iota pub async fn send_message(&self, cv: &CommunicationValue) { if !cv.is_type(CommunicationType::pong) { - log_out!(PrintType::Iota, "{}", cv.to_json().to_string()); + log_out!( + self.get_iota_id().await, + PrintType::Iota, + "{}", + cv.to_json().to_string() + ); } self.send_message_str(&cv.to_json().to_string()).await; } @@ -146,7 +156,12 @@ impl IotaConnection { return; } - log_in!(PrintType::Iota, "{}", cv.to_json().to_string()); + log_in!( + self.get_iota_id().await, + PrintType::Iota, + "{}", + cv.to_json().to_string() + ); let identified = *self.identified.read().await; let challenged = *self.challenged.read().await; @@ -157,7 +172,7 @@ impl IotaConnection { .and_then(|v| v.as_i64()) .unwrap_or(0); if iota_id == 0 { - log_out!(PrintType::Iota, "Invalid IOTA ID"); + log_out!(self.get_iota_id().await, PrintType::Iota, "Invalid IOTA ID"); self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_data) .await; self.close().await; @@ -332,6 +347,7 @@ impl IotaConnection { if let Ok(iota_users_cv) = iota_users_cv { if !iota_users_cv.is_type(CommunicationType::iota_user_data) { log_err!( + self.get_iota_id().await, PrintType::Omikron, "Invalid communication type {:?}", iota_users_cv.get_type() @@ -351,9 +367,18 @@ impl IotaConnection { _ => {} } } else { - log_err!(PrintType::Omikron, "Failed to retrieve user IDs"); + log_err!( + self.get_iota_id().await, + PrintType::Omikron, + "Failed to retrieve user IDs" + ); } - log_in!(PrintType::General, "User IDs: {:?}", user_ids.clone()); + log_in!( + self.get_iota_id().await, + PrintType::General, + "User IDs: {:?}", + user_ids.clone() + ); *self.user_ids.write().await = user_ids.clone(); let rho_connection = @@ -537,10 +562,14 @@ impl IotaConnection { // Process contacts and add call information let enriched_contacts = if *empty { if let Some(contacts_data) = cv.get_data(DataTypes::user_ids) { - log_in!(PrintType::Call, "Call empty"); + log_in!(self.get_iota_id().await, PrintType::Call, "Call empty"); contacts_data.clone() } else { - log_in!(PrintType::Call, "Call empty No Data"); + log_in!( + self.get_iota_id().await, + PrintType::Call, + "Call empty No Data" + ); JsonValue::new_array() } } else { @@ -583,7 +612,11 @@ impl IotaConnection { let updated_cv = cv.with_sender(self.get_iota_id().await); rho_conn.message_to_client(updated_cv).await; } else { - log_err!(PrintType::General, "Failed to forward message to client"); + log_err!( + self.get_iota_id().await, + PrintType::General, + "Failed to forward message to client" + ); } } @@ -596,7 +629,7 @@ impl IotaConnection { } pub async fn await_response( - &self, + self: Arc, cv: &CommunicationValue, timeout_duration: Option, ) -> Result { @@ -606,11 +639,12 @@ impl IotaConnection { let task_tx = tx.clone(); self.waiting_tasks.insert( msg_id, - Box::new(move |_, response_cv| { + Box::new(move |io, response_cv| { let inner_tx = task_tx.clone(); tokio::spawn(async move { if let Err(e) = inner_tx.send(response_cv).await { log_err!( + io.get_iota_id().await, PrintType::Iota, "Failed to send response back to awaiter: {}", e diff --git a/src/rho/rho_manager.rs b/src/rho/rho_manager.rs index 2334ca6..55f67af 100644 --- a/src/rho/rho_manager.rs +++ b/src/rho/rho_manager.rs @@ -12,9 +12,9 @@ pub static RHO_CONNECTIONS: LazyLock> pub async fn get_rho_con_for_user(user_id: i64) -> Option> { let connections = RHO_CONNECTIONS.read().await; - log_in!(PrintType::Client, "Checking user ID: {:?}", user_id,); for rho_connection in connections.values() { log_in!( + user_id, PrintType::Client, "Comparing user IDs: {:?}", rho_connection.get_user_ids().to_vec() diff --git a/src/util/logger.rs b/src/util/logger.rs index 360e520..6fdcacd 100644 --- a/src/util/logger.rs +++ b/src/util/logger.rs @@ -30,8 +30,6 @@ struct LogMessage { message: String, } -/// Initialize the logging subsystem. -/// Must be called exactly once during startup. pub fn startup() { let (tx, rx) = mpsc::channel::(); LOGGER.set(tx).expect("Logger already initialized"); @@ -61,10 +59,8 @@ pub fn startup() { let line = format!("{} {} {} {}", ts, sender, msg.prefix, msg.message); - // Console (ANSI-colored) println!("{}", colorize(msg.kind, msg.is_error).paint(&line)); - // File (plain text) let _ = writeln!(file, "{}", line); } }); @@ -95,16 +91,15 @@ fn fixed_box(content: &str, width: usize) -> String { } } -/** Internal async logging entry point. -* Not exposed publicly; all access goes through macros. -*/ pub fn log_internal( - sender: Option, + sender: i64, kind: PrintType, prefix: &'static str, is_error: bool, message: String, ) { + let sender = if sender == 0 { None } else { Some(sender) }; + if let Some(tx) = LOGGER.get() { let _ = tx.send(LogMessage { timestamp_ms: SystemTime::now() @@ -119,101 +114,27 @@ pub fn log_internal( }); } } -/// Log a general informational message. #[macro_export] macro_rules! log { - - // actor only - ($kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(None, $kind, "", false, format!($($arg)*)) - }; - - // sender + actor - ($sender:expr, $kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(Some($sender), $kind, "", false, format!($($arg)*)) - }; - - // plain - ($($arg:tt)*) => { - $crate::util::logger::log_internal( - None, - $crate::util::logger::PrintType::General, - "", - false, - format!($($arg)*) - ) + ($sender: expr, $kind:expr, $($arg:tt)*) => { + $crate::util::logger::log_internal($sender, $kind, "", false, format!($($arg)*)) }; } -/// Log an inbound message (`>`). #[macro_export] macro_rules! log_in { - // actor only - ($kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(None, $kind, ">", false, format!($($arg)*)) - }; - - // sender + actor - ($sender:expr, $kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(Some($sender), $kind, ">", false, format!($($arg)*)) - }; - - // plain - ($($arg:tt)*) => { - $crate::util::logger::log_internal( - None, - $crate::util::logger::PrintType::General, - ">", - false, - format!($($arg)*) - ) + ($sender: expr, $kind:expr, $($arg:tt)*) => { + $crate::util::logger::log_internal($sender, $kind, ">", false, format!($($arg)*)) }; } -/// Log an outbound message (`<`). #[macro_export] macro_rules! log_out { - // actor only - ($kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(None, $kind, "<", false, format!($($arg)*)) - }; - - // sender + actor ($sender:expr, $kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(Some($sender), $kind, "<", false, format!($($arg)*)) - }; - - // plain - ($($arg:tt)*) => { - $crate::util::logger::log_internal( - None, - $crate::util::logger::PrintType::General, - "<", - false, - format!($($arg)*) - ) + $crate::util::logger::log_internal($sender, $kind, "<", false, format!($($arg)*)) }; } -/// Log an error message (`>>`). #[macro_export] macro_rules! log_err { - - // actor only - ($kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(None, $kind, ">>", true, format!($($arg)*)) - }; - - // sender + actor ($sender:expr, $kind:expr, $($arg:tt)*) => { - $crate::util::logger::log_internal(Some($sender), $kind, ">>", true, format!($($arg)*)) - }; - - // plain - ($($arg:tt)*) => { - $crate::util::logger::log_internal( - None, - $crate::util::logger::PrintType::General, - ">>", - true, - format!($($arg)*) - ) + $crate::util::logger::log_internal($sender, $kind, ">>", true, format!($($arg)*)) }; }