From 7e595a7d0ff7c6166e3455532427099203b60274 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 15:36:56 +0200 Subject: [PATCH 1/8] Logger --- src/engine/engine.rs | 4 +++- src/engine/terminal.rs | 49 +++++++++++++++++++++++++++++++++--------- 2 files changed, 42 insertions(+), 11 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index ed734b6..a309542 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -1,9 +1,11 @@ +use std::sync::Arc; + use crate::terminal::TerminalServer; const WATCH_LIST_SYMBOLS: &[&str] = &["BTC", "ETH", "SOL", "XRP"]; pub struct Engine { - pub terminal_server: TerminalServer, + pub terminal_server: Arc, } impl Engine { diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 85a9fb5..c046d85 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,6 +1,9 @@ use std::sync::Arc; -use pulse_wire::{PulseWire, terminal::TerminalClientMessage}; +use pulse_wire::{ + PulseWire, + terminal::{EventLog, LogKind, TerminalClientMessage}, +}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::{ @@ -10,19 +13,21 @@ use tokio::{ sync::Mutex, }; -#[derive(Debug, Clone)] +#[derive(Debug)] pub struct TerminalServer { - clients: Arc>>, + clients: Mutex>, + logs: Mutex>, } impl TerminalServer { - pub fn new() -> Self { - Self { - clients: Arc::new(Mutex::new(Vec::new())), - } + pub fn new() -> Arc { + Arc::new(Self { + clients: Mutex::new(Vec::new()), + logs: Mutex::new(Vec::new()), + }) } - pub async fn run(&self) -> tokio::io::Result<()> { + pub async fn run(self: &Arc) -> tokio::io::Result<()> { let path = pulse_wire::server_path(); if path.exists() { @@ -50,7 +55,7 @@ impl TerminalServer { } } - async fn handle_client(&self, mut reader: OwnedReadHalf) -> tokio::io::Result<()> { + async fn handle_client(self: &Arc, mut reader: OwnedReadHalf) -> tokio::io::Result<()> { loop { let mut len_buf = [0u8; size_of::()]; let size = reader.read_exact(&mut len_buf).await?; @@ -95,7 +100,7 @@ impl TerminalServer { } pub async fn broadcast( - &self, + self: &Arc, message: pulse_wire::terminal::TerminalServerMessage, ) -> tokio::io::Result<()> { let msg = message.to_com(); @@ -123,4 +128,28 @@ impl TerminalServer { Ok(()) } + + pub async fn log(self: &Arc, kind: LogKind, name: &str, message: &str) { + self.logs.lock().await.push(EventLog { + kind, + name: name.to_string(), + message: message.to_string(), + }); + } + + pub async fn info(self: &Arc, name: &str, message: &str) { + self.log(LogKind::Info, name, message).await + } + + pub async fn warn(self: &Arc, name: &str, message: &str) { + self.log(LogKind::Warn, name, message).await + } + + pub async fn error(self: &Arc, name: &str, message: &str) { + self.log(LogKind::Err, name, message).await + } + + pub async fn debug(self: &Arc, name: &str, message: &str) { + self.log(LogKind::Debug, name, message).await + } } From 8bf9a5d4d2f3bd004904da41a0b866cbddd49b83 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 15:55:15 +0200 Subject: [PATCH 2/8] 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 } } From bba94e244c2dacaaeb1c4a7f4c6f6c29c9364e05 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 16:36:44 +0200 Subject: [PATCH 3/8] Config --- Cargo.lock | 32 ++++++++++++++++++++++++++ Cargo.toml | 2 ++ src/engine/config.rs | 54 ++++++++++++++++++++++++++++++++++++++++++++ src/engine/main.rs | 1 + 4 files changed, 89 insertions(+) create mode 100644 src/engine/config.rs diff --git a/Cargo.lock b/Cargo.lock index a60d9a7..fc34e9b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3888,8 +3888,10 @@ dependencies = [ "pulse-ui", "pulse-wire", "rand 0.8.7", + "serde", "serde_json", "tokio", + "toml", ] [[package]] @@ -4729,6 +4731,15 @@ dependencies = [ "zmij", ] +[[package]] +name = "serde_spanned" +version = "1.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6662b5879511e06e8999a8a235d848113e942c9124f211511b16466ee2995f26" +dependencies = [ + "serde_core", +] + [[package]] name = "serde_urlencoded" version = "0.7.1" @@ -5240,6 +5251,21 @@ dependencies = [ "tokio", ] +[[package]] +name = "toml" +version = "1.1.3+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "53c96ecdfa941c8fc4fcaed14f99ada8ebed502eef533015095a07e3301d4c3c" +dependencies = [ + "indexmap 2.14.0", + "serde_core", + "serde_spanned", + "toml_datetime", + "toml_parser", + "toml_writer", + "winnow", +] + [[package]] name = "toml_datetime" version = "1.1.1+spec-1.1.0" @@ -5270,6 +5296,12 @@ dependencies = [ "winnow", ] +[[package]] +name = "toml_writer" +version = "1.1.2+spec-1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7d56353a2a665ad0f41a421187180aab746c8c325620617ad883a99a1cbe66d2" + [[package]] name = "tonic" version = "0.14.6" diff --git a/Cargo.toml b/Cargo.toml index 22ba238..f70095e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,6 +20,8 @@ crossterm = { workspace = true } hypersdk = "0.2.14" serde_json = "1" rand = "0.8.7" +toml = "1.1.3" +serde = { version = "1.0.229", features = ["serde_derive"] } [workspace] members = ["pulse-macros", "pulse-ui", "pulse-wire"] diff --git a/src/engine/config.rs b/src/engine/config.rs new file mode 100644 index 0000000..3380131 --- /dev/null +++ b/src/engine/config.rs @@ -0,0 +1,54 @@ +use std::path::PathBuf; + +pub fn home_dir() -> tokio::io::Result { + std::env::home_dir().ok_or_else(|| { + tokio::io::Error::new(std::io::ErrorKind::NotFound, "Unable to get home directory") + }) +} + +pub fn pulse_directory() -> tokio::io::Result { + Ok(home_dir()?.join(".config/pulse-trader")) +} + +pub fn pulse_config_directory() -> tokio::io::Result { + Ok(home_dir()?.join(".config/pulse-trader/config.toml")) +} + +#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] +pub struct WatchList { + symbols: Vec, +} + +#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] +pub struct Config { + watchlist: WatchList, +} + +impl Config { + pub async fn new() -> tokio::io::Result { + let path = pulse_config_directory()?; + + if !path.exists() { + let default = Self::default(); + + tokio::fs::create_dir_all(path.parent().unwrap()).await?; + tokio::fs::write(&path, default.to_string()?).await?; + + return Ok(default); + } + + let output = tokio::fs::read_to_string(path).await?; + + Self::from_str(&output) + } + + pub fn from_str(s: &str) -> tokio::io::Result { + toml::from_str(s) + .map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string())) + } + + pub fn to_string(&self) -> tokio::io::Result { + toml::to_string(self) + .map_err(|v| tokio::io::Error::new(std::io::ErrorKind::Other, v.to_string())) + } +} diff --git a/src/engine/main.rs b/src/engine/main.rs index bb7f06f..b1a7454 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -1,3 +1,4 @@ +pub mod config; pub mod engine; pub mod fetch; pub mod terminal; From ae37c02adcdca940d08ddba9dca5b65906550dcc Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 16:50:36 +0200 Subject: [PATCH 4/8] Using config in engine --- src/engine/config.rs | 4 ++-- src/engine/engine.rs | 21 +++++++++++++-------- src/engine/fetch.rs | 4 ++-- src/engine/main.rs | 2 +- 4 files changed, 18 insertions(+), 13 deletions(-) diff --git a/src/engine/config.rs b/src/engine/config.rs index 3380131..8555d75 100644 --- a/src/engine/config.rs +++ b/src/engine/config.rs @@ -16,12 +16,12 @@ pub fn pulse_config_directory() -> tokio::io::Result { #[derive(Debug, Default, serde::Serialize, serde::Deserialize)] pub struct WatchList { - symbols: Vec, + pub symbols: Vec, } #[derive(Debug, Default, serde::Serialize, serde::Deserialize)] pub struct Config { - watchlist: WatchList, + pub watchlist: WatchList, } impl Config { diff --git a/src/engine/engine.rs b/src/engine/engine.rs index a309542..e1dd561 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -1,18 +1,20 @@ use std::sync::Arc; -use crate::terminal::TerminalServer; +use tokio::sync::Mutex; -const WATCH_LIST_SYMBOLS: &[&str] = &["BTC", "ETH", "SOL", "XRP"]; +use crate::{config::Config, terminal::TerminalServer}; pub struct Engine { pub terminal_server: Arc, + pub config: Arc>, } impl Engine { - pub fn new() -> Self { - Self { + pub async fn new() -> tokio::io::Result { + Ok(Self { terminal_server: TerminalServer::new(), - } + config: Arc::new(Mutex::new(Config::new().await?)), + }) } pub fn spawn_terminal_server(&self) { @@ -26,8 +28,9 @@ impl Engine { }); } - pub fn spawn_broadcaster(&mut self) { + pub fn spawn_broadcaster(&self) { let terminal_server = self.terminal_server.clone(); + let config = self.config.clone(); tokio::spawn(async move { let mut refresh = tokio::time::interval(tokio::time::Duration::from_secs(5)); @@ -35,7 +38,9 @@ impl Engine { loop { refresh.tick().await; - match crate::fetch::fetch_watch_list(WATCH_LIST_SYMBOLS).await { + let watch_list = &config.lock().await.watchlist.symbols; + + match crate::fetch::fetch_watch_list(watch_list).await { Ok(watch_list) => { if let Err(error) = terminal_server .broadcast( @@ -54,7 +59,7 @@ impl Engine { }); } - pub async fn run_engine(&mut self) -> tokio::io::Result<()> { + pub async fn run_engine(&self) -> tokio::io::Result<()> { loop { tokio::time::sleep(tokio::time::Duration::from_millis(5000)).await; } diff --git a/src/engine/fetch.rs b/src/engine/fetch.rs index 1c6ae18..9f7b271 100644 --- a/src/engine/fetch.rs +++ b/src/engine/fetch.rs @@ -12,7 +12,7 @@ fn number(value: &Value, field: &str) -> Result { .map_err(|error| format!("could not parse asset context field {field} ({raw}): {error}")) } -pub async fn fetch_watch_list(symbols: &[&str]) -> Result, String> { +pub async fn fetch_watch_list(symbols: &[String]) -> Result, String> { let response = hypersdk::hypercore::mainnet() .meta_and_asset_ctxs(None) .await @@ -73,7 +73,7 @@ pub async fn fetch_watch_list(symbols: &[&str]) -> Result, St .iter() .map(|symbol| { by_symbol - .remove(*symbol) + .remove(symbol.as_str()) .ok_or_else(|| format!("{symbol} is not in the Hyperliquid perpetual universe")) }) .collect() diff --git a/src/engine/main.rs b/src/engine/main.rs index b1a7454..1b50670 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -5,7 +5,7 @@ pub mod terminal; #[tokio::main] async fn main() -> tokio::io::Result<()> { - let mut engine = engine::Engine::new(); + let engine = engine::Engine::new().await?; engine.spawn_terminal_server(); From f036d04d7b2c2388d9f2f1eb46232dcd8cb67a33 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 16:56:16 +0200 Subject: [PATCH 5/8] Improved error handling and spawn structure --- src/engine/engine.rs | 60 +++++++++++++++++++++++--------------------- src/engine/main.rs | 4 +-- 2 files changed, 33 insertions(+), 31 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index e1dd561..e98ba9f 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -4,6 +4,7 @@ use tokio::sync::Mutex; use crate::{config::Config, terminal::TerminalServer}; +#[derive(Debug, Clone)] pub struct Engine { pub terminal_server: Arc, pub config: Arc>, @@ -17,46 +18,47 @@ impl Engine { }) } - pub fn spawn_terminal_server(&self) { + pub async fn spawn_terminal_server(&self) -> tokio::io::Result<()> { let terminal_server = self.terminal_server.clone(); - tokio::spawn(async move { - terminal_server - .run() - .await - .expect("Failed to run terminal server"); - }); + tokio::spawn(async move { terminal_server.run().await }) + .await + .expect("Failed to spawn terminal server") } - pub fn spawn_broadcaster(&self) { - let terminal_server = self.terminal_server.clone(); - let config = self.config.clone(); + pub async fn spawn_broadcaster(&self) -> tokio::io::Result<()> { + let s = self.clone(); - tokio::spawn(async move { - let mut refresh = tokio::time::interval(tokio::time::Duration::from_secs(5)); + tokio::spawn(async move { s.run_broadcaster().await }) + .await + .expect("Failed to spawn broadcaster") + } - loop { - refresh.tick().await; + pub async fn run_broadcaster(&self) -> tokio::io::Result<()> { + let mut refresh = tokio::time::interval(tokio::time::Duration::from_secs(5)); - let watch_list = &config.lock().await.watchlist.symbols; + loop { + refresh.tick().await; - match crate::fetch::fetch_watch_list(watch_list).await { - Ok(watch_list) => { - if let Err(error) = terminal_server - .broadcast( - pulse_wire::terminal::TerminalServerMessage::WatchListUpdated( - watch_list, - ), - ) - .await - { - eprintln!("Failed to broadcast Hyperliquid watch list: {error}"); - } + let watch_list = &self.config.lock().await.watchlist.symbols; + + match crate::fetch::fetch_watch_list(watch_list).await { + Ok(watch_list) => { + if let Err(error) = self + .terminal_server + .broadcast( + pulse_wire::terminal::TerminalServerMessage::WatchListUpdated( + watch_list, + ), + ) + .await + { + eprintln!("Failed to broadcast Hyperliquid watch list: {error}"); } - Err(error) => eprintln!("Failed to refresh Hyperliquid watch list: {error}"), } + Err(error) => eprintln!("Failed to refresh Hyperliquid watch list: {error}"), } - }); + } } pub async fn run_engine(&self) -> tokio::io::Result<()> { diff --git a/src/engine/main.rs b/src/engine/main.rs index 1b50670..0caae94 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -7,9 +7,9 @@ pub mod terminal; async fn main() -> tokio::io::Result<()> { let engine = engine::Engine::new().await?; - engine.spawn_terminal_server(); + engine.spawn_terminal_server().await?; - engine.spawn_broadcaster(); + engine.spawn_broadcaster().await?; engine.run_engine().await?; From 2eb80394907c2203e726c08c1b927e3b5a50ffe7 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 17:36:53 +0200 Subject: [PATCH 6/8] Engine as command executor --- src/engine/engine.rs | 27 ++++++++++++++++++++++----- src/engine/terminal.rs | 29 +++++++++++++++++------------ 2 files changed, 39 insertions(+), 17 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index e98ba9f..7506360 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -11,11 +11,13 @@ pub struct Engine { } impl Engine { - pub async fn new() -> tokio::io::Result { - Ok(Self { - terminal_server: TerminalServer::new(), - config: Arc::new(Mutex::new(Config::new().await?)), - }) + pub async fn new() -> tokio::io::Result> { + let config = Arc::new(Mutex::new(Config::new().await?)); + + Ok(Arc::new_cyclic(|engine| Self { + terminal_server: TerminalServer::new(engine.clone()), + config, + })) } pub async fn spawn_terminal_server(&self) -> tokio::io::Result<()> { @@ -66,4 +68,19 @@ impl Engine { tokio::time::sleep(tokio::time::Duration::from_millis(5000)).await; } } + + pub async fn execute_command(&self, command: &str, _args: Vec<&str>) -> tokio::io::Result<()> { + match command { + _ => { + self.terminal_server + .error( + "Command executor", + &format!("Command '{}' not found", command), + ) + .await?; + } + } + + Ok(()) + } } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index e4273fc..4e3a503 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,4 +1,7 @@ -use std::{collections::HashMap, sync::Arc}; +use std::{ + collections::HashMap, + sync::{Arc, Weak}, +}; use pulse_wire::{ PulseWire, @@ -13,17 +16,21 @@ use tokio::{ sync::Mutex, }; +use crate::engine::Engine; + #[derive(Debug)] pub struct TerminalServer { clients: Mutex>, logs: Mutex>, + engine: Weak, } impl TerminalServer { - pub fn new() -> Arc { + pub fn new(engine: Weak) -> Arc { Arc::new(Self { clients: Mutex::new(HashMap::new()), logs: Mutex::new(Vec::new()), + engine, }) } @@ -86,21 +93,13 @@ impl TerminalServer { TerminalClientMessage::ExecuteCommand(command) => { let command = command.as_str(); - let (command, _args) = if let Some((command, args)) = command.split_once(" ") { + let (command, args) = if let Some((command, args)) = command.split_once(" ") { (command, args.split(" ").collect()) } else { (command, Vec::new()) }; - match command { - _ => { - self.error( - "Command executor", - &format!("Command '{}' not found", command), - ) - .await?; - } - } + self.get_engine().execute_command(command, args).await?; } } } @@ -190,4 +189,10 @@ impl TerminalServer { pub async fn debug(self: &Arc, name: &str, message: &str) -> tokio::io::Result<()> { self.log(LogKind::Debug, name, message).await } + + pub fn get_engine(&self) -> Arc { + self.engine + .upgrade() + .expect("Failed to upgrade engine(Weak) to Arc") + } } From ef5f88516a14d57d473764a065ebbe8cb130a0ac Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 17:47:06 +0200 Subject: [PATCH 7/8] Fix spawning --- src/engine/engine.rs | 10 +++------- src/engine/main.rs | 9 +++++++-- 2 files changed, 10 insertions(+), 9 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index 7506360..3581add 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -1,6 +1,6 @@ use std::sync::Arc; -use tokio::sync::Mutex; +use tokio::{sync::Mutex, task::JoinHandle}; use crate::{config::Config, terminal::TerminalServer}; @@ -20,20 +20,16 @@ impl Engine { })) } - pub async fn spawn_terminal_server(&self) -> tokio::io::Result<()> { + pub async fn spawn_terminal_server(&self) -> JoinHandle> { let terminal_server = self.terminal_server.clone(); tokio::spawn(async move { terminal_server.run().await }) - .await - .expect("Failed to spawn terminal server") } - pub async fn spawn_broadcaster(&self) -> tokio::io::Result<()> { + pub async fn spawn_broadcaster(&self) -> JoinHandle> { let s = self.clone(); tokio::spawn(async move { s.run_broadcaster().await }) - .await - .expect("Failed to spawn broadcaster") } pub async fn run_broadcaster(&self) -> tokio::io::Result<()> { diff --git a/src/engine/main.rs b/src/engine/main.rs index 0caae94..2e7cd52 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -7,11 +7,16 @@ pub mod terminal; async fn main() -> tokio::io::Result<()> { let engine = engine::Engine::new().await?; - engine.spawn_terminal_server().await?; + let terminal_server = engine.spawn_terminal_server().await; - engine.spawn_broadcaster().await?; + let broadcaster = engine.spawn_broadcaster().await; engine.run_engine().await?; + let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster); + + terminal_server??; + broadcaster??; + Ok(()) } From c84b26b2ec18d3211043905e2bba5ffff2ac2ce9 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 24 Jul 2026 18:02:26 +0200 Subject: [PATCH 8/8] Config reload command --- src/engine/engine.rs | 22 +++++++++++++++++++++- 1 file changed, 21 insertions(+), 1 deletion(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index 3581add..0aa96c7 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -65,8 +65,28 @@ impl Engine { } } - pub async fn execute_command(&self, command: &str, _args: Vec<&str>) -> tokio::io::Result<()> { + pub async fn execute_command(&self, command: &str, args: Vec<&str>) -> tokio::io::Result<()> { match command { + "config" => { + const MESSAGE: &str = "Invalid command arguments, usage: config "; + + if args.len() != 1 { + self.terminal_server.error("config", MESSAGE).await?; + + return Ok(()); + } + + match args[0] { + "reload" => { + *self.config.lock().await = Config::new().await?; + } + + _ => { + self.terminal_server.error("config", MESSAGE).await?; + } + } + } + _ => { self.terminal_server .error(