From 34a4d817b92c1fb8ca8bf7dfe18bad3e4c7a83fe Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 23:19:50 +0200 Subject: [PATCH] Terminal client --- Cargo.toml | 2 +- src/engine/main.rs | 27 ++++++++-- src/engine/terminal.rs | 24 +++++---- src/terminal/main.rs | 53 ++++++++++++-------- src/terminal/terminal.rs | 104 +++++++++++++++++++++++++++++++++++++++ 5 files changed, 176 insertions(+), 34 deletions(-) create mode 100644 src/terminal/terminal.rs diff --git a/Cargo.toml b/Cargo.toml index ac45c21..60ec14c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,7 +8,7 @@ pulse-ui = { workspace = true } pulse-wire = { workspace = true } chrono = "0.4.45" -tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util"] } +tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time"] } crossterm = { workspace = true } [workspace] diff --git a/src/engine/main.rs b/src/engine/main.rs index 8448dbb..27ed309 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -2,9 +2,30 @@ pub mod terminal; #[tokio::main] async fn main() -> tokio::io::Result<()> { - let mut terminal_server = terminal::TerminalServer::new(); + let terminal_server = terminal::TerminalServer::new(); - terminal_server.run().await?; + { + let terminal_server = terminal_server.clone(); - Ok(()) + tokio::spawn(async move { + terminal_server + .run() + .await + .expect("Failed to run terminal server"); + }); + } + + loop { + tokio::time::sleep(tokio::time::Duration::from_millis(1500)).await; + terminal_server + .broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog( + pulse_wire::terminal::EventLog { + kind: pulse_wire::terminal::LogKind::Debug, + name: "Engine".to_string(), + message: "Hello, world!".to_string(), + }, + )) + .await?; + println!("Ok"); + } } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 84cb530..4938a68 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,3 +1,5 @@ +use std::sync::Arc; + use pulse_wire::PulseWire; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, @@ -5,20 +7,22 @@ use tokio::{ UnixListener, unix::{OwnedReadHalf, OwnedWriteHalf}, }, + sync::Mutex, }; +#[derive(Debug, Clone)] pub struct TerminalServer { - clients: Vec, + clients: Arc>>, } impl TerminalServer { pub fn new() -> Self { Self { - clients: Vec::new(), + clients: Arc::new(Mutex::new(Vec::new())), } } - pub async fn run(&mut self) -> tokio::io::Result<()> { + pub async fn run(&self) -> tokio::io::Result<()> { let path = pulse_wire::server_path(); if path.exists() { @@ -34,7 +38,7 @@ impl TerminalServer { let (reader, writer) = stream.into_split(); - self.clients.push(writer); + self.clients.lock().await.push(writer); tokio::spawn(async move { if let Err(err) = Self::handle_client(reader).await { @@ -67,18 +71,18 @@ impl TerminalServer { } pub async fn broadcast( - &mut self, + &self, message: pulse_wire::terminal::TerminalServerMessage, ) -> tokio::io::Result<()> { - let mut res = Ok(()); let msg = message.to_com(); - for client in &mut self.clients { - if let Err(e) = client.write(&msg).await { - res = Err(e); + for i in (0..self.clients.lock().await.len()).rev() { + if let Err(e) = self.clients.lock().await[i].write(&msg).await { + self.clients.lock().await.remove(i); + println!("{e:?}"); } } - res + Ok(()) } } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 0f02f1a..135bff1 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -1,5 +1,6 @@ pub mod command; pub mod formatting; +pub mod terminal; use std::any::Any; @@ -270,26 +271,38 @@ impl App for PulseTradeApp { } #[tokio::main] -async fn main() { - pulse_ui::run(|ctx| PulseTradeApp { - command: ctx.use_state(InputState::new()), - watch_list: ctx.use_state(Vec::new()), - active_positions: ctx.use_state(Vec::new()), - signals: ctx.use_state(Vec::new()), - logs: ctx.use_state(Vec::new()), - inspect: ctx.use_state(InspectTarget::None), - market_overview: ctx.use_state(MarketOverview { - trend: pulse_wire::terminal::MarketTrend::Bullish, - volatility: pulse_wire::terminal::Volatility::High, - pressure: 0.324, - alerts: Vec::new(), - }), - status: ctx.use_state(Status { - feed: pulse_wire::terminal::Feed::Connected, - exchange: "Binance".to_string(), - dex: "DEX SCREENER".to_string(), - latency: 18, - }), +async fn main() -> tokio::io::Result<()> { + let mut client = terminal::TerminalClient::new().await?; + + client + .send(pulse_wire::terminal::TerminalClientMessage::ExecuteCommand( + "Hello".to_string(), + )) + .await?; + + pulse_ui::run(|ctx| { + client.use_app(PulseTradeApp { + command: ctx.use_state(InputState::new()), + watch_list: ctx.use_state(Vec::new()), + active_positions: ctx.use_state(Vec::new()), + signals: ctx.use_state(Vec::new()), + logs: ctx.use_state(Vec::new()), + inspect: ctx.use_state(InspectTarget::None), + market_overview: ctx.use_state(MarketOverview { + trend: pulse_wire::terminal::MarketTrend::Bullish, + volatility: pulse_wire::terminal::Volatility::High, + pressure: 0.324, + alerts: Vec::new(), + }), + status: ctx.use_state(Status { + feed: pulse_wire::terminal::Feed::Connected, + exchange: "Binance".to_string(), + dex: "DEX SCREENER".to_string(), + latency: 18, + }), + }) }) .await; + + Ok(()) } diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs new file mode 100644 index 0000000..ef41a73 --- /dev/null +++ b/src/terminal/terminal.rs @@ -0,0 +1,104 @@ +use pulse_wire::{PulseWire, server_path, terminal::TerminalServerMessage}; +use tokio::{ + io::{AsyncReadExt, AsyncWriteExt}, + net::{ + UnixStream, + unix::{OwnedReadHalf, OwnedWriteHalf}, + }, +}; + +use crate::PulseTradeApp; + +pub struct TerminalClient { + writer: OwnedWriteHalf, + reader: Option, +} + +impl TerminalClient { + pub async fn new() -> tokio::io::Result { + let (reader, writer) = UnixStream::connect(server_path()).await?.into_split(); + + Ok(Self { + reader: Some(reader), + writer, + }) + } + + pub async fn send( + &mut self, + message: pulse_wire::terminal::TerminalClientMessage, + ) -> tokio::io::Result<()> { + self.writer.write(&message.to_com()).await?; + + Ok(()) + } + + pub fn use_app(&mut self, app: PulseTradeApp) -> PulseTradeApp { + let mut reader = None; + + std::mem::swap(&mut self.reader, &mut reader); + + let mut reader = reader.expect("Reader failed to swap"); + + let watch_list = app.watch_list.clone(); + let active_positions = app.active_positions.clone(); + let logs = app.logs.clone(); + let signals = app.signals.clone(); + let market_overview = app.market_overview.clone(); + let status = app.status.clone(); + let inspect = app.inspect.clone(); + + tokio::spawn(async move { + loop { + let mut buffer = Vec::new(); + + let len = reader + .read(&mut buffer) + .await + .expect("Failed to read socket"); + + if len == 0 { + break; + } + + buffer.truncate(len); + + match TerminalServerMessage::from_com(&mut buffer) { + TerminalServerMessage::WatchListUpdated(v) => { + *watch_list.lock().await = v; + } + + TerminalServerMessage::PositionsUpdated(v) => { + *active_positions.lock().await = v; + } + + TerminalServerMessage::OverviewUpdated(v) => { + *market_overview.lock().await = v; + } + + TerminalServerMessage::SignalsUpdated(v) => { + *signals.lock().await = v; + } + + TerminalServerMessage::Inspect(v) => { + *inspect.lock().await = v; + } + + TerminalServerMessage::SetStatus(v) => { + *status.lock().await = v; + } + + TerminalServerMessage::SetLogs(v) => { + *logs.lock().await = v; + } + + TerminalServerMessage::AddLog(v) => { + logs.lock().await.push(v); + } + } + } + }); + + app + } +}