From 8bf9a5d4d2f3bd004904da41a0b866cbddd49b83 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 15:55:15 +0200 Subject: [PATCH] Working logger --- Cargo.lock | 1 + Cargo.toml | 10 ++++- src/engine/terminal.rs | 98 +++++++++++++++++++++++++++++------------- 3 files changed, 78 insertions(+), 31 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d61f03c..a60d9a7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3887,6 +3887,7 @@ dependencies = [ "hypersdk", "pulse-ui", "pulse-wire", + "rand 0.8.7", "serde_json", "tokio", ] diff --git a/Cargo.toml b/Cargo.toml index d516a74..22ba238 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,10 +8,18 @@ pulse-ui = { workspace = true } pulse-wire = { workspace = true } chrono = "0.4.45" -tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time"] } +tokio = { workspace = true, features = [ + "rt-multi-thread", + "macros", + "net", + "fs", + "io-util", + "time", +] } crossterm = { workspace = true } hypersdk = "0.2.14" serde_json = "1" +rand = "0.8.7" [workspace] members = ["pulse-macros", "pulse-ui", "pulse-wire"] diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index c046d85..e4273fc 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,4 +1,4 @@ -use std::sync::Arc; +use std::{collections::HashMap, sync::Arc}; use pulse_wire::{ PulseWire, @@ -15,14 +15,14 @@ use tokio::{ #[derive(Debug)] pub struct TerminalServer { - clients: Mutex>, + clients: Mutex>, logs: Mutex>, } impl TerminalServer { pub fn new() -> Arc { Arc::new(Self { - clients: Mutex::new(Vec::new()), + clients: Mutex::new(HashMap::new()), logs: Mutex::new(Vec::new()), }) } @@ -43,19 +43,31 @@ impl TerminalServer { let (reader, writer) = stream.into_split(); - self.clients.lock().await.push(writer); + let id = rand::random(); + + self.clients.lock().await.insert(id, writer); let s = self.clone(); tokio::spawn(async move { - if let Err(err) = s.handle_client(reader).await { + if let Err(err) = s.handle_client(&id, reader).await { eprintln!("Terminal connection error: {err}"); } }); } } - async fn handle_client(self: &Arc, mut reader: OwnedReadHalf) -> tokio::io::Result<()> { + async fn handle_client( + self: &Arc, + id: &usize, + mut reader: OwnedReadHalf, + ) -> tokio::io::Result<()> { + self.send_to( + id, + pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), + ) + .await?; + loop { let mut len_buf = [0u8; size_of::()]; let size = reader.read_exact(&mut len_buf).await?; @@ -82,13 +94,10 @@ impl TerminalServer { match command { _ => { - self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog( - pulse_wire::terminal::EventLog { - kind: pulse_wire::terminal::LogKind::Err, - name: "Command executor".to_string(), - message: format!("Command '{}' not found", command), - }, - )) + self.error( + "Command executor", + &format!("Command '{}' not found", command), + ) .await?; } } @@ -107,49 +116,78 @@ impl TerminalServer { let mut clients = self.clients.lock().await; - for i in (0..clients.len()).rev() { - if let Err(e) = Self::send_to_client(&mut clients, i, &msg).await { - clients.remove(i); + let mut remove_clients = Vec::new(); + + for (id, client) in clients.iter_mut() { + if let Err(e) = Self::send_to_client(client, &msg).await { + remove_clients.push(*id); println!("{e:?}"); } } + for id in remove_clients { + clients.remove(&id); + } + Ok(()) } - pub async fn send_to_client( - clients: &mut tokio::sync::MutexGuard<'_, Vec>, - i: usize, - msg: &[u8], + pub async fn send_to( + self: &Arc, + id: &usize, + message: pulse_wire::terminal::TerminalServerMessage, ) -> tokio::io::Result<()> { - clients[i].write(&msg.len().to_le_bytes()).await?; - clients[i].write(msg).await?; - clients[i].flush().await?; + Self::send_to_client( + self.clients.lock().await.get_mut(id).ok_or_else(|| { + tokio::io::Error::new( + std::io::ErrorKind::Other, + format!("Client({id}) does not exist"), + ) + })?, + &message.to_com(), + ) + .await + } + + pub async fn send_to_client(client: &mut OwnedWriteHalf, msg: &[u8]) -> tokio::io::Result<()> { + client.write(&msg.len().to_le_bytes()).await?; + client.write(msg).await?; + client.flush().await?; Ok(()) } - pub async fn log(self: &Arc, kind: LogKind, name: &str, message: &str) { - self.logs.lock().await.push(EventLog { + pub async fn log( + self: &Arc, + kind: LogKind, + name: &str, + message: &str, + ) -> tokio::io::Result<()> { + let log = EventLog { kind, name: name.to_string(), message: message.to_string(), - }); + }; + + self.logs.lock().await.push(log.clone()); + + self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog(log)) + .await } - pub async fn info(self: &Arc, name: &str, message: &str) { + pub async fn info(self: &Arc, name: &str, message: &str) -> tokio::io::Result<()> { self.log(LogKind::Info, name, message).await } - pub async fn warn(self: &Arc, name: &str, message: &str) { + pub async fn warn(self: &Arc, name: &str, message: &str) -> tokio::io::Result<()> { self.log(LogKind::Warn, name, message).await } - pub async fn error(self: &Arc, name: &str, message: &str) { + pub async fn error(self: &Arc, name: &str, message: &str) -> tokio::io::Result<()> { self.log(LogKind::Err, name, message).await } - pub async fn debug(self: &Arc, name: &str, message: &str) { + pub async fn debug(self: &Arc, name: &str, message: &str) -> tokio::io::Result<()> { self.log(LogKind::Debug, name, message).await } }