From 39735bad87ec69e1cc961e9cddbbb1c877196d94 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 18:08:16 +0200 Subject: [PATCH 01/16] Pulse wire crate --- Cargo.lock | 9 +++++- Cargo.toml | 9 ++++-- pulse-wire/Cargo.toml | 7 +++++ src/pc.rs => pulse-wire/src/lib.rs | 24 ++------------- src/ptc.rs => pulse-wire/src/terminal.rs | 1 + src/daemon/main.rs | 10 +++---- src/daemon/ptc.rs | 2 -- src/terminal/command.rs | 6 ++-- src/terminal/formatting.rs | 2 +- src/terminal/main.rs | 37 +++++++++++------------- src/terminal/ptc.rs | 2 -- 11 files changed, 49 insertions(+), 60 deletions(-) create mode 100644 pulse-wire/Cargo.toml rename src/pc.rs => pulse-wire/src/lib.rs (77%) rename src/ptc.rs => pulse-wire/src/terminal.rs (99%) delete mode 100644 src/daemon/ptc.rs delete mode 100644 src/terminal/ptc.rs diff --git a/Cargo.lock b/Cargo.lock index 1f87829..9019c07 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -319,8 +319,8 @@ version = "0.1.0-alpha.0" dependencies = [ "chrono", "crossterm", - "pulse-macros", "pulse-ui", + "pulse-wire", "tokio", ] @@ -332,6 +332,13 @@ dependencies = [ "tokio", ] +[[package]] +name = "pulse-wire" +version = "0.1.0-alpha.0" +dependencies = [ + "pulse-macros", +] + [[package]] name = "quote" version = "1.0.46" diff --git a/Cargo.toml b/Cargo.toml index 572e0f1..651f439 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,18 +4,21 @@ version = "0.1.0-alpha.0" edition = "2024" [dependencies] -chrono = "0.4.45" pulse-ui = { workspace = true } -pulse-macros = { workspace = true } +pulse-wire = { workspace = true } + +chrono = "0.4.45" tokio = { workspace = true, features = ["rt-multi-thread", "macros"] } crossterm = { workspace = true } [workspace] -members = ["pulse-macros", "pulse-ui"] +members = ["pulse-macros", "pulse-ui", "pulse-wire"] [workspace.dependencies] pulse-macros = { path = "pulse-macros", version = "0.1.0-alpha.0" } pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" } +pulse-wire = { path = "pulse-wire", version = "0.1.0-alpha.0" } + tokio = "1.52.3" crossterm = "0.29.0" diff --git a/pulse-wire/Cargo.toml b/pulse-wire/Cargo.toml new file mode 100644 index 0000000..449b94f --- /dev/null +++ b/pulse-wire/Cargo.toml @@ -0,0 +1,7 @@ +[package] +name = "pulse-wire" +version = "0.1.0-alpha.0" +edition = "2024" + +[dependencies] +pulse-macros = { workspace = true } diff --git a/src/pc.rs b/pulse-wire/src/lib.rs similarity index 77% rename from src/pc.rs rename to pulse-wire/src/lib.rs index 8c5cb31..a4c1292 100644 --- a/src/pc.rs +++ b/pulse-wire/src/lib.rs @@ -1,3 +1,5 @@ +pub mod terminal; + pub trait PulseCom { fn to_com(&self) -> Vec; fn from_com(_com: &mut Vec) -> Self; @@ -81,25 +83,3 @@ int_com!(usize); int_com!(f64); int_com!(f32); - -#[macro_export] -macro_rules! p_com { - (struct $name:ident { $($n:ident: $v:ty),* $(,)? }) => { - #[derive(Debug, Clone)] - pub struct $name { $(pub $n: $v),* } - - impl PulseCom for $name { - fn to_com(&self) -> Vec { - let mut vec = Vec::new(); - $(vec.extend(self.$n.to_com());)* - vec - } - - fn from_com(com: &mut Vec) -> Self { - Self { - $($n: <$v>::from_com(com),)* - } - } - } - }; -} diff --git a/src/ptc.rs b/pulse-wire/src/terminal.rs similarity index 99% rename from src/ptc.rs rename to pulse-wire/src/terminal.rs index 4dc117c..1972329 100644 --- a/src/ptc.rs +++ b/pulse-wire/src/terminal.rs @@ -1,3 +1,4 @@ +use crate::PulseCom; use pulse_macros::p_com; #[p_com] diff --git a/src/daemon/main.rs b/src/daemon/main.rs index 61b7dd0..acd2ea1 100644 --- a/src/daemon/main.rs +++ b/src/daemon/main.rs @@ -1,12 +1,10 @@ -pub mod ptc; - use std::time::Instant; -use ptc::PulseCom; +use pulse_wire::PulseCom; fn main() { - let input = ptc::EventLog { - kind: ptc::LogKind::Warn, + let input = pulse_wire::terminal::EventLog { + kind: pulse_wire::terminal::LogKind::Warn, name: "Test".to_string(), message: "Hello, WOrld".to_string(), }; @@ -14,7 +12,7 @@ fn main() { let start = Instant::now(); let mut val = input.to_com(); - let out = ptc::EventLog::from_com(&mut val); + let out = pulse_wire::terminal::EventLog::from_com(&mut val); let elapsed = start.elapsed(); diff --git a/src/daemon/ptc.rs b/src/daemon/ptc.rs deleted file mode 100644 index ab16651..0000000 --- a/src/daemon/ptc.rs +++ /dev/null @@ -1,2 +0,0 @@ -include!("../pc.rs"); -include!("../ptc.rs"); diff --git a/src/terminal/command.rs b/src/terminal/command.rs index 8a8a4f6..3eecd8a 100644 --- a/src/terminal/command.rs +++ b/src/terminal/command.rs @@ -1,4 +1,4 @@ -use crate::{PulseTradeApp, ptc::EventLog}; +use crate::PulseTradeApp; impl PulseTradeApp { pub async fn execute_command(&mut self, ctx: &pulse_ui::state::Context, command: &str) { @@ -18,8 +18,8 @@ impl PulseTradeApp { } _ => { - self.logs.lock().await.push(EventLog { - kind: crate::ptc::LogKind::Err, + self.logs.lock().await.push(pulse_wire::terminal::EventLog { + kind: pulse_wire::terminal::LogKind::Err, name: "cmd".to_string(), message: format!("Command '{}' not found", command), }); diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 9b62766..9dfef94 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -1,4 +1,4 @@ -use crate::ptc::{ +use pulse_wire::terminal::{ ActivePosition, Alert, AlertLevel, EventLog, InspectTarget, LogKind, MarketOverview, Signal, SignalKind, Status, WatchListItem, }; diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 1b90a56..0f02f1a 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -1,6 +1,5 @@ pub mod command; pub mod formatting; -pub mod ptc; use std::any::Any; @@ -18,12 +17,10 @@ use pulse_ui::{ }, }; -use crate::{ - formatting::{Formatted, apply_padding}, - ptc::{ - ActivePosition, Alert, EventLog, InspectTarget, MarketOverview, Signal, Status, - WatchListItem, - }, +use crate::formatting::{Formatted, apply_padding}; + +use pulse_wire::terminal::{ + ActivePosition, Alert, EventLog, InspectTarget, MarketOverview, Signal, Status, WatchListItem, }; pub struct PulseTradeApp { @@ -88,30 +85,30 @@ impl App for PulseTradeApp { let mut signals = self.signals.lock().await; signals.push(Signal { - kind: ptc::SignalKind::Buy, + kind: pulse_wire::terminal::SignalKind::Buy, symbol: "BTC".to_string(), - param: ptc::SignalParameter::Lim, + param: pulse_wire::terminal::SignalParameter::Lim, price: 118_800.0, }); signals.push(Signal { - kind: ptc::SignalKind::Buy, + kind: pulse_wire::terminal::SignalKind::Buy, symbol: "BTC".to_string(), - param: ptc::SignalParameter::Tap, + param: pulse_wire::terminal::SignalParameter::Tap, price: 120_000.0, }); signals.push(Signal { - kind: ptc::SignalKind::Buy, + kind: pulse_wire::terminal::SignalKind::Buy, symbol: "BTC".to_string(), - param: ptc::SignalParameter::Stl, + param: pulse_wire::terminal::SignalParameter::Stl, price: 118_000.0, }); let mut logs = self.logs.lock().await; logs.push(EventLog { - kind: ptc::LogKind::Warn, + kind: pulse_wire::terminal::LogKind::Warn, name: "pulse.init".to_string(), message: "We're still not done yet ;)".to_string(), }); @@ -119,17 +116,17 @@ impl App for PulseTradeApp { let mut market_overview = self.market_overview.lock().await; market_overview.alerts.push(Alert { - level: ptc::AlertLevel::High, + level: pulse_wire::terminal::AlertLevel::High, message: "BTC funding rate elevated".to_string(), }); market_overview.alerts.push(Alert { - level: ptc::AlertLevel::Medium, + level: pulse_wire::terminal::AlertLevel::Medium, message: "Market volatility increasing".to_string(), }); market_overview.alerts.push(Alert { - level: ptc::AlertLevel::Low, + level: pulse_wire::terminal::AlertLevel::Low, message: "ETH volatility returning to normal".to_string(), }); } @@ -282,13 +279,13 @@ async fn main() { logs: ctx.use_state(Vec::new()), inspect: ctx.use_state(InspectTarget::None), market_overview: ctx.use_state(MarketOverview { - trend: ptc::MarketTrend::Bullish, - volatility: ptc::Volatility::High, + trend: pulse_wire::terminal::MarketTrend::Bullish, + volatility: pulse_wire::terminal::Volatility::High, pressure: 0.324, alerts: Vec::new(), }), status: ctx.use_state(Status { - feed: ptc::Feed::Connected, + feed: pulse_wire::terminal::Feed::Connected, exchange: "Binance".to_string(), dex: "DEX SCREENER".to_string(), latency: 18, diff --git a/src/terminal/ptc.rs b/src/terminal/ptc.rs deleted file mode 100644 index ab16651..0000000 --- a/src/terminal/ptc.rs +++ /dev/null @@ -1,2 +0,0 @@ -include!("../pc.rs"); -include!("../ptc.rs"); From a641de074f3f609d39decae41be28cc6a0a41387 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 18:30:16 +0200 Subject: [PATCH 02/16] Name refactor --- Cargo.toml | 4 ++-- pulse-macros/src/lib.rs | 12 ++++++------ pulse-wire/src/lib.rs | 8 ++++---- pulse-wire/src/terminal.rs | 34 +++++++++++++++++----------------- src/{daemon => engine}/main.rs | 2 +- 5 files changed, 30 insertions(+), 30 deletions(-) rename src/{daemon => engine}/main.rs (94%) diff --git a/Cargo.toml b/Cargo.toml index 651f439..ec9ed90 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -27,5 +27,5 @@ name = "pulse-trader" path = "src/terminal/main.rs" [[bin]] -name = "pulse-trader-daemon" -path = "src/daemon/main.rs" +name = "pulse-trader-engine" +path = "src/engine/main.rs" diff --git a/pulse-macros/src/lib.rs b/pulse-macros/src/lib.rs index d554b2d..cac53b8 100644 --- a/pulse-macros/src/lib.rs +++ b/pulse-macros/src/lib.rs @@ -3,7 +3,7 @@ use quote::quote; use syn::{Fields, ItemEnum, ItemStruct, parse_macro_input}; #[proc_macro_attribute] -pub fn p_com(_: TokenStream, item: TokenStream) -> TokenStream { +pub fn pwp(_: TokenStream, item: TokenStream) -> TokenStream { let input = parse_macro_input!(item as syn::Item); match input { @@ -49,7 +49,7 @@ fn expand_struct(mut input: ItemStruct) -> TokenStream { #[derive(Debug, Clone)] #input - impl PulseCom for #name { + impl PulseWire for #name { fn to_com(&self) -> Vec { let mut vec = Vec::new(); @@ -63,7 +63,7 @@ fn expand_struct(mut input: ItemStruct) -> TokenStream { fn from_com(com: &mut Vec) -> Self { Self { #( - #field_names2: PulseCom::from_com(com), + #field_names2: PulseWire::from_com(com), )* } } @@ -122,7 +122,7 @@ fn expand_enum(input: ItemEnum) -> TokenStream { Fields::Unnamed(fields) if fields.unnamed.len() == 1 => { quote! { - #tag => Self::#ident(PulseCom::from_com(com)), + #tag => Self::#ident(PulseWire::from_com(com)), } } @@ -132,7 +132,7 @@ fn expand_enum(input: ItemEnum) -> TokenStream { quote! { #tag => Self::#ident { #( - #names: PulseCom::from_com(com), + #names: PulseWire::from_com(com), )* }, } @@ -146,7 +146,7 @@ fn expand_enum(input: ItemEnum) -> TokenStream { #[derive(Debug, Clone)] #input - impl PulseCom for #name { + impl PulseWire for #name { fn to_com(&self) -> Vec { let mut vec = Vec::new(); diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index a4c1292..794082e 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -1,11 +1,11 @@ pub mod terminal; -pub trait PulseCom { +pub trait PulseWire { fn to_com(&self) -> Vec; fn from_com(_com: &mut Vec) -> Self; } -impl PulseCom for Vec { +impl PulseWire for Vec { fn to_com(&self) -> Vec { let mut vec = Vec::new(); @@ -33,7 +33,7 @@ impl PulseCom for Vec { } } -impl PulseCom for String { +impl PulseWire for String { fn to_com(&self) -> Vec { let bytes = self.as_bytes(); let mut out = Vec::with_capacity(4 + bytes.len()); @@ -55,7 +55,7 @@ impl PulseCom for String { macro_rules! int_com { ($t:ty) => { - impl PulseCom for $t { + impl $crate::PulseWire for $t { fn to_com(&self) -> Vec { self.to_le_bytes().to_vec() } diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index 1972329..3247487 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -1,35 +1,35 @@ -use crate::PulseCom; -use pulse_macros::p_com; +use crate::PulseWire; +use pulse_macros::pwp; -#[p_com] +#[pwp] pub struct WatchListItem { symbol: String, price: f64, trend: f64, } -#[p_com] +#[pwp] pub struct ActivePosition { symbol: String, profit: f64, amount: f64, } -#[p_com] +#[pwp] pub enum MarketTrend { Bullish, Bearish, Neutral, } -#[p_com] +#[pwp] pub enum Volatility { Low, Medium, High, } -#[p_com] +#[pwp] pub struct MarketOverview { trend: MarketTrend, volatility: Volatility, @@ -38,7 +38,7 @@ pub struct MarketOverview { alerts: Vec, } -#[p_com] +#[pwp] pub enum Feed { Connected, Disconnected, @@ -46,7 +46,7 @@ pub enum Feed { Failed, } -#[p_com] +#[pwp] pub struct Status { feed: Feed, exchange: String, @@ -54,13 +54,13 @@ pub struct Status { latency: u16, } -#[p_com] +#[pwp] pub enum SignalKind { Buy, Sell, } -#[p_com] +#[pwp] pub enum SignalParameter { Lim, Stl, @@ -68,7 +68,7 @@ pub enum SignalParameter { Chk, } -#[p_com] +#[pwp] pub struct Signal { kind: SignalKind, symbol: String, @@ -76,7 +76,7 @@ pub struct Signal { price: f64, } -#[p_com] +#[pwp] pub enum LogKind { Info, Warn, @@ -84,27 +84,27 @@ pub enum LogKind { Debug, } -#[p_com] +#[pwp] pub struct EventLog { kind: LogKind, name: String, message: String, } -#[p_com] +#[pwp] pub enum AlertLevel { High, Medium, Low, } -#[p_com] +#[pwp] pub struct Alert { level: AlertLevel, message: String, } -#[p_com] +#[pwp] pub enum InspectTarget { None, Symbol(WatchListItem), diff --git a/src/daemon/main.rs b/src/engine/main.rs similarity index 94% rename from src/daemon/main.rs rename to src/engine/main.rs index acd2ea1..4e7f550 100644 --- a/src/daemon/main.rs +++ b/src/engine/main.rs @@ -1,6 +1,6 @@ use std::time::Instant; -use pulse_wire::PulseCom; +use pulse_wire::PulseWire; fn main() { let input = pulse_wire::terminal::EventLog { From 790082674eba5588f2947fff11d2be5d58d9d219 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 19:23:51 +0200 Subject: [PATCH 03/16] Preparing terminal communication --- Cargo.lock | 21 ++++++++++++ Cargo.toml | 2 +- pulse-wire/src/lib.rs | 7 ++++ src/engine/main.rs | 2 ++ src/engine/terminal.rs | 72 ++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 103 insertions(+), 1 deletion(-) create mode 100644 src/engine/terminal.rs diff --git a/Cargo.lock b/Cargo.lock index 9019c07..9caa2a4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -29,6 +29,12 @@ version = "3.20.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" +[[package]] +name = "bytes" +version = "1.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" + [[package]] name = "cc" version = "1.2.67" @@ -446,6 +452,16 @@ version = "1.15.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8ed6a63f02c8539c91a8685a86f4099661ba3da017932f6ebbea6de3f0fa7c90" +[[package]] +name = "socket2" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51" +dependencies = [ + "libc", + "windows-sys", +] + [[package]] name = "syn" version = "2.0.118" @@ -463,8 +479,13 @@ version = "1.52.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8fc7f01b389ac15039e4dc9531aa973a135d7a4135281b12d7c1bc79fd57fffe" dependencies = [ + "bytes", + "libc", + "mio", "pin-project-lite", + "socket2", "tokio-macros", + "windows-sys", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index ec9ed90..ac45c21 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"] } +tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util"] } crossterm = { workspace = true } [workspace] diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index 794082e..5a53401 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -1,5 +1,12 @@ +#[cfg(target_os = "macos")] +use std::path::PathBuf; + pub mod terminal; +pub fn server_path() -> PathBuf { + PathBuf::from("/tmp/pulse-engine.sock") +} + pub trait PulseWire { fn to_com(&self) -> Vec; fn from_com(_com: &mut Vec) -> Self; diff --git a/src/engine/main.rs b/src/engine/main.rs index 4e7f550..70a47b8 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -1,3 +1,5 @@ +pub mod terminal; + use std::time::Instant; use pulse_wire::PulseWire; diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs new file mode 100644 index 0000000..ece2b82 --- /dev/null +++ b/src/engine/terminal.rs @@ -0,0 +1,72 @@ +use tokio::{ + io::AsyncReadExt, + net::{ + UnixListener, + unix::{OwnedReadHalf, OwnedWriteHalf}, + }, +}; + +pub struct TerminalServer { + clients: Vec, +} + +impl TerminalServer { + pub fn new() -> Self { + Self { + clients: Vec::new(), + } + } + + pub async fn run(&mut self) -> tokio::io::Result<()> { + let path = pulse_wire::server_path(); + + if path.exists() { + tokio::fs::remove_file(&path).await?; + } + + let listener = UnixListener::bind(&path)?; + + println!("Terminal server listening on {:?}", path); + + loop { + let (stream, _) = listener.accept().await?; + + let (reader, writer) = stream.into_split(); + + self.clients.push(writer); + + tokio::spawn(async move { + if let Err(err) = Self::handle_client(reader).await { + eprintln!("Terminal connection error: {err}"); + } + }); + } + } + + async fn handle_client(mut reader: OwnedReadHalf) -> tokio::io::Result<()> { + let input = tokio::spawn(async move { + let mut buffer = [0u8; 2048]; + + loop { + let size = reader.read(&mut buffer).await?; + + if size == 0 { + break; + } + + println!("Received {} bytes", size); + } + + Ok::<(), tokio::io::Error>(()) + }); + + let _ = tokio::try_join!(input)?; + + Ok(()) + } + + pub async fn broadcast(&mut self) -> tokio::io::Result<()> { + // for client in &mut self.clients {} + Ok(()) + } +} From 51118e3ee8658ef74bf78f78ae4dbc181009b133 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 19:35:13 +0200 Subject: [PATCH 04/16] Terminal messages --- pulse-wire/src/terminal.rs | 34 ++++++++++++++++++++++++++++++++++ 1 file changed, 34 insertions(+) diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index 3247487..8f598d8 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -1,6 +1,40 @@ use crate::PulseWire; use pulse_macros::pwp; +#[pwp] +pub enum TerminalServerMessage { + // WatchList + SetWatchList(Vec), + + // Positions + SetPositions(Vec), + AddPosition(ActivePosition), + RemovePosition(usize), + + // Overview + SetOverview(MarketOverview), + + // Signals + SetSignals(Vec), + AddSignal(Signal), + RemoveSignal(usize), + + // Inspector + Inspect(InspectTarget), + + // Status + SetStatus(Status), + + // Logs + SetLogs(Vec), + AddLog(EventLog), +} + +#[pwp] +pub enum TerminalClientMessage { + ExecuteCommand(String), +} + #[pwp] pub struct WatchListItem { symbol: String, From 701c345c839efa8b0f0d635e2bc2271ddcd8f356 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 21:29:16 +0200 Subject: [PATCH 05/16] Updated terminal messages --- pulse-wire/src/terminal.rs | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index 8f598d8..11a87de 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -4,20 +4,16 @@ use pulse_macros::pwp; #[pwp] pub enum TerminalServerMessage { // WatchList - SetWatchList(Vec), + WatchListUpdated(Vec), // Positions - SetPositions(Vec), - AddPosition(ActivePosition), - RemovePosition(usize), + PositionsUpdated(Vec), // Overview - SetOverview(MarketOverview), + OverviewUpdated(MarketOverview), // Signals - SetSignals(Vec), - AddSignal(Signal), - RemoveSignal(usize), + SignalsUpdated(Vec), // Inspector Inspect(InspectTarget), From af11afb2557449fa39e0b28f52ec84bcde15b342 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 21:58:32 +0200 Subject: [PATCH 06/16] Broadcasting --- src/engine/terminal.rs | 20 ++++++++++++++++---- 1 file changed, 16 insertions(+), 4 deletions(-) diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index ece2b82..84cb530 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,5 +1,6 @@ +use pulse_wire::PulseWire; use tokio::{ - io::AsyncReadExt, + io::{AsyncReadExt, AsyncWriteExt}, net::{ UnixListener, unix::{OwnedReadHalf, OwnedWriteHalf}, @@ -65,8 +66,19 @@ impl TerminalServer { Ok(()) } - pub async fn broadcast(&mut self) -> tokio::io::Result<()> { - // for client in &mut self.clients {} - Ok(()) + pub async fn broadcast( + &mut 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); + } + } + + res } } From cc9aef33127ca1ea8513f5a3bf7283f531278270 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 22:01:12 +0200 Subject: [PATCH 07/16] Simple terminal server --- src/engine/main.rs | 22 +++++----------------- 1 file changed, 5 insertions(+), 17 deletions(-) diff --git a/src/engine/main.rs b/src/engine/main.rs index 70a47b8..8448dbb 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -1,22 +1,10 @@ pub mod terminal; -use std::time::Instant; +#[tokio::main] +async fn main() -> tokio::io::Result<()> { + let mut terminal_server = terminal::TerminalServer::new(); -use pulse_wire::PulseWire; + terminal_server.run().await?; -fn main() { - let input = pulse_wire::terminal::EventLog { - kind: pulse_wire::terminal::LogKind::Warn, - name: "Test".to_string(), - message: "Hello, WOrld".to_string(), - }; - - let start = Instant::now(); - - let mut val = input.to_com(); - let out = pulse_wire::terminal::EventLog::from_com(&mut val); - - let elapsed = start.elapsed(); - - println!("{:?} {:?}", out, elapsed); + Ok(()) } From 34a4d817b92c1fb8ca8bf7dfe18bad3e4c7a83fe Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 23:19:50 +0200 Subject: [PATCH 08/16] 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 + } +} From 2ff5adfdafd89b4b2a7d0a5e5eded13488ed4d3e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 22 Jul 2026 23:26:49 +0200 Subject: [PATCH 09/16] Fix reading bug --- src/engine/terminal.rs | 8 +++++--- src/terminal/terminal.rs | 2 +- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 4938a68..ee21fe2 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -76,9 +76,11 @@ impl TerminalServer { ) -> tokio::io::Result<()> { let msg = message.to_com(); - 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); + let mut clients = self.clients.lock().await; + + for i in (0..clients.len()).rev() { + if let Err(e) = clients[i].write(&msg).await { + clients.remove(i); println!("{e:?}"); } } diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index ef41a73..7c0f1f6 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -50,7 +50,7 @@ impl TerminalClient { tokio::spawn(async move { loop { - let mut buffer = Vec::new(); + let mut buffer = vec![0u8; 4096]; let len = reader .read(&mut buffer) From e4773db9d83037f3feb6d1b34db05b927ec363ac Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 23 Jul 2026 00:22:20 +0200 Subject: [PATCH 10/16] Scrolling --- pulse-ui/src/widget/mod.rs | 5 +-- pulse-ui/src/widget/scroll.rs | 67 +++++++++++++++++++++++++++++++++++ src/terminal/main.rs | 60 ++++++++++++++++++------------- 3 files changed, 105 insertions(+), 27 deletions(-) create mode 100644 pulse-ui/src/widget/scroll.rs diff --git a/pulse-ui/src/widget/mod.rs b/pulse-ui/src/widget/mod.rs index dc75b5f..b5cc038 100644 --- a/pulse-ui/src/widget/mod.rs +++ b/pulse-ui/src/widget/mod.rs @@ -1,8 +1,9 @@ pub mod align; -pub mod outline; -pub mod spaced; pub mod dynamic; pub mod input; +pub mod outline; +pub mod scroll; +pub mod spaced; use crate::render::RenderScope; diff --git a/pulse-ui/src/widget/scroll.rs b/pulse-ui/src/widget/scroll.rs new file mode 100644 index 0000000..ef93d34 --- /dev/null +++ b/pulse-ui/src/widget/scroll.rs @@ -0,0 +1,67 @@ +use crate::widget::Widget; + +pub struct ScrollState(pub usize, pub [usize; N]); + +pub struct ScrollText { + pub scroll: usize, + pub title: String, + pub text: String, +} + +impl Widget for ScrollText { + fn render(&self, scope: &mut crate::render::RenderScope) { + scope.draw_text(0, &self.title); + + let title_lines = self.title.lines().count(); + + for (y, line) in self + .text + .lines() + .skip(self.scroll) + .take(scope.rect.height as usize - title_lines) + .enumerate() + { + scope.draw_text((0, (y + title_lines) as u16), line); + } + } +} + +impl ScrollState { + pub fn get_selected(&self, index: usize) -> &'static str { + if index == self.0 { "\x1b[32m" } else { "" } + } + + pub fn scroll(&self, index: usize, title: String, text: String) -> ScrollText { + ScrollText { + scroll: self.1[index], + title, + text, + } + } + + pub fn up(&mut self) { + if self.1[self.0] > 0 { + self.1[self.0] -= 1; + } + } + + pub fn down(&mut self) { + self.1[self.0] += 1; + } + + pub fn back_tab(&mut self) { + if self.0 > 0 { + self.0 -= 1; + } else { + self.0 = N - 1; + } + } + + pub fn tab(&mut self) { + if self.0 + 1 < N { + self.0 += 1; + } else { + self.0 = 0; + } + } +} diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 135bff1..848bf9f 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -14,6 +14,7 @@ use pulse_ui::{ align::End, input::{Input, InputState}, outline::{Outline, VLine}, + scroll::{ScrollState, ScrollText}, spaced::SpacedColumns, }, }; @@ -26,6 +27,7 @@ use pulse_wire::terminal::{ pub struct PulseTradeApp { command: State, + scroll: State>, watch_list: State>, active_positions: State>, logs: State>, @@ -149,6 +151,14 @@ impl App for PulseTradeApp { drop(command); self.execute_command(ctx, command_text.trim()).await; + } else if key.code.is_up() { + self.scroll.value.lock().await.up(); + } else if key.code.is_down() { + self.scroll.value.lock().await.down(); + } else if key.code.is_back_tab() { + self.scroll.value.lock().await.back_tab(); + } else if key.code.is_tab() { + self.scroll.value.lock().await.tab(); } } } @@ -207,18 +217,14 @@ impl App for PulseTradeApp { SpacedColumns(vec![ ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " WATCHLIST\n{}", - apply_padding(self.watch_list.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 0, "WATCH LIST", &self.watch_list).await), ), ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " ACTIVE POSITIONS\n{}", - apply_padding(self.active_positions.lock().await.get_formatted()) - .join("\n") - )), + Box::new( + advanced_draw(&self.scroll, 1, "ACTIVE POSITIONS", &self.active_positions) + .await, + ), ), (LayoutItem::Widget(Size::Flex(1)), { let mo = self.market_overview.lock().await; @@ -236,34 +242,22 @@ impl App for PulseTradeApp { SpacedColumns(vec![ ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " SIGNALS\n{}", - apply_padding(self.signals.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 3, "SIGNALS", &self.signals).await), ), ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " INSPECTOR\n{}", - apply_padding(self.inspect.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 4, "INSPECTOR", &self.inspect).await), ), ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " STATUS\n{}", - apply_padding(self.status.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 5, "STATUS", &self.status).await), ), ]), ); layout.draw( 5, - format!( - " EVENT LOGS\n{}", - apply_padding(self.logs.lock().await.get_formatted()).join("\n") - ), + advanced_draw(&self.scroll, 6, "EVENT LOGS", &self.logs).await, ); layout.draw(6, Input(" > ", &*self.command.lock().await)); @@ -283,6 +277,7 @@ async fn main() -> tokio::io::Result<()> { pulse_ui::run(|ctx| { client.use_app(PulseTradeApp { command: ctx.use_state(InputState::new()), + scroll: ctx.use_state(ScrollState(0, [0; 7])), watch_list: ctx.use_state(Vec::new()), active_positions: ctx.use_state(Vec::new()), signals: ctx.use_state(Vec::new()), @@ -306,3 +301,18 @@ async fn main() -> tokio::io::Result<()> { Ok(()) } + +pub async fn advanced_draw( + scroll: &State>, + index: usize, + title: &'static str, + state: &State, +) -> ScrollText { + let scroll = scroll.lock().await; + + scroll.scroll( + index, + format!("{}{title}\x1b[0m", scroll.get_selected(index)), + apply_padding(state.lock().await.get_formatted()).join("\n"), + ) +} From b76a1f9aeda13c4cd3354b9b63a58f0ab1fa0eef Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 23 Jul 2026 00:30:57 +0200 Subject: [PATCH 11/16] Fix market overview --- src/terminal/formatting.rs | 11 +++++++++-- src/terminal/main.rs | 17 ++++++++--------- 2 files changed, 17 insertions(+), 11 deletions(-) diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 9dfef94..dd431b2 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -137,7 +137,7 @@ struct Property(&'static str, String); impl Formatted for MarketOverview { fn get_formatted(&self) -> Vec { - vec![ + let mut o = vec![ Property("TREND", format!("{}", self.trend)), Property("VOLATILITY", format!("{}", self.volatility)), Property( @@ -149,7 +149,14 @@ impl Formatted for MarketOverview { }, ), ] - .get_formatted() + .get_formatted(); + + o.push(format!( + "\n{}", + apply_padding(self.alerts.get_formatted()).join("\n") + )); + + o } } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 848bf9f..3dfc51c 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -226,14 +226,13 @@ impl App for PulseTradeApp { .await, ), ), - (LayoutItem::Widget(Size::Flex(1)), { - let mo = self.market_overview.lock().await; - Box::new(format!( - " MARKET OVERVIEW\n{}\n\n{}", - apply_padding(mo.get_formatted()).join("\n"), - apply_padding(mo.alerts.get_formatted()).join("\n") - )) - }), + ( + LayoutItem::Widget(Size::Flex(1)), + Box::new( + advanced_draw(&self.scroll, 2, "MARKET OVERVIEW", &self.market_overview) + .await, + ), + ), ]), ); @@ -312,7 +311,7 @@ pub async fn advanced_draw( scroll.scroll( index, - format!("{}{title}\x1b[0m", scroll.get_selected(index)), + format!(" {}{title}\x1b[0m", scroll.get_selected(index)), apply_padding(state.lock().await.get_formatted()).join("\n"), ) } From ad84145f12e86a7019768232969a80138aeab5e5 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 23 Jul 2026 00:40:58 +0200 Subject: [PATCH 12/16] Scroll improvements --- pulse-ui/src/widget/scroll.rs | 2 +- src/terminal/main.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pulse-ui/src/widget/scroll.rs b/pulse-ui/src/widget/scroll.rs index ef93d34..be6511e 100644 --- a/pulse-ui/src/widget/scroll.rs +++ b/pulse-ui/src/widget/scroll.rs @@ -28,7 +28,7 @@ impl Widget for ScrollText { impl ScrollState { pub fn get_selected(&self, index: usize) -> &'static str { - if index == self.0 { "\x1b[32m" } else { "" } + if index == self.0 { "\x1b[4m\x1b[1m" } else { "\x1b[1m" } } pub fn scroll(&self, index: usize, title: String, text: String) -> ScrollText { diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 3dfc51c..1bc6008 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -276,7 +276,7 @@ async fn main() -> tokio::io::Result<()> { pulse_ui::run(|ctx| { client.use_app(PulseTradeApp { command: ctx.use_state(InputState::new()), - scroll: ctx.use_state(ScrollState(0, [0; 7])), + scroll: ctx.use_state(ScrollState(1, [0; 7])), watch_list: ctx.use_state(Vec::new()), active_positions: ctx.use_state(Vec::new()), signals: ctx.use_state(Vec::new()), From e3ddaa43750f1d51e310b689c04fc7d3eda90f5f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 23 Jul 2026 00:54:17 +0200 Subject: [PATCH 13/16] Using pwp --- src/engine/main.rs | 83 +++++++++++++++++++++++++++++++++++ src/engine/terminal.rs | 3 ++ src/terminal/main.rs | 98 +----------------------------------------- 3 files changed, 88 insertions(+), 96 deletions(-) diff --git a/src/engine/main.rs b/src/engine/main.rs index 27ed309..85c45ff 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -1,3 +1,5 @@ +use pulse_wire::terminal::{ActivePosition, Signal, WatchListItem}; + pub mod terminal; #[tokio::main] @@ -17,6 +19,86 @@ async fn main() -> tokio::io::Result<()> { loop { tokio::time::sleep(tokio::time::Duration::from_millis(1500)).await; + + terminal_server + .broadcast( + pulse_wire::terminal::TerminalServerMessage::WatchListUpdated(vec![ + WatchListItem { + symbol: "BTC".to_string(), + price: 118_402.12, + trend: 0.82, + }, + WatchListItem { + symbol: "ETH".to_string(), + price: 3_912.48, + trend: -0.41, + }, + WatchListItem { + symbol: "SOL".to_string(), + price: 182.91, + trend: 2.18, + }, + WatchListItem { + symbol: "XRP".to_string(), + price: 2.84, + trend: 1.22, + }, + ]), + ) + .await?; + + terminal_server + .broadcast( + pulse_wire::terminal::TerminalServerMessage::PositionsUpdated(vec![ + ActivePosition { + symbol: "BTC".to_string(), + profit: 125.50, + amount: 0.25, + }, + ActivePosition { + symbol: "SOL".to_string(), + profit: 84.20, + amount: 5.0, + }, + ActivePosition { + symbol: "ETH".to_string(), + profit: -32.75, + amount: 1.0, + }, + ActivePosition { + symbol: "XRP".to_string(), + profit: 12.30, + amount: 0.75, + }, + ]), + ) + .await?; + + terminal_server + .broadcast(pulse_wire::terminal::TerminalServerMessage::SignalsUpdated( + vec![ + Signal { + kind: pulse_wire::terminal::SignalKind::Buy, + symbol: "BTC".to_string(), + param: pulse_wire::terminal::SignalParameter::Lim, + price: 118_800.0, + }, + Signal { + kind: pulse_wire::terminal::SignalKind::Buy, + symbol: "BTC".to_string(), + param: pulse_wire::terminal::SignalParameter::Tap, + price: 120_000.0, + }, + Signal { + kind: pulse_wire::terminal::SignalKind::Buy, + symbol: "BTC".to_string(), + param: pulse_wire::terminal::SignalParameter::Stl, + price: 118_000.0, + }, + ], + )) + .await?; + terminal_server .broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog( pulse_wire::terminal::EventLog { @@ -26,6 +108,7 @@ async fn main() -> tokio::io::Result<()> { }, )) .await?; + println!("Ok"); } } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index ee21fe2..467afe4 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -82,6 +82,9 @@ impl TerminalServer { if let Err(e) = clients[i].write(&msg).await { clients.remove(i); println!("{e:?}"); + } else if let Err(e) = clients[i].flush().await { + clients.remove(i); + println!("{e:?}"); } } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 1bc6008..13f9366 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -22,7 +22,7 @@ use pulse_ui::{ use crate::formatting::{Formatted, apply_padding}; use pulse_wire::terminal::{ - ActivePosition, Alert, EventLog, InspectTarget, MarketOverview, Signal, Status, WatchListItem, + ActivePosition, EventLog, InspectTarget, MarketOverview, Signal, Status, WatchListItem, }; pub struct PulseTradeApp { @@ -38,101 +38,7 @@ pub struct PulseTradeApp { } impl App for PulseTradeApp { - async fn init(&mut self, _ctx: &pulse_ui::state::Context) { - let mut watch_list = self.watch_list.lock().await; - - watch_list.push(WatchListItem { - symbol: "BTC".to_string(), - price: 118_402.12, - trend: 0.82, - }); - watch_list.push(WatchListItem { - symbol: "ETH".to_string(), - price: 3_912.48, - trend: -0.41, - }); - watch_list.push(WatchListItem { - symbol: "SOL".to_string(), - price: 182.91, - trend: 2.18, - }); - watch_list.push(WatchListItem { - symbol: "XRP".to_string(), - price: 2.84, - trend: 1.22, - }); - - let mut active_positions = self.active_positions.lock().await; - - active_positions.push(ActivePosition { - symbol: "BTC".to_string(), - profit: 125.50, - amount: 0.25, - }); - active_positions.push(ActivePosition { - symbol: "ETH".to_string(), - profit: -32.75, - amount: 1.0, - }); - active_positions.push(ActivePosition { - symbol: "SOL".to_string(), - profit: 84.20, - amount: 5.0, - }); - active_positions.push(ActivePosition { - symbol: "XRP".to_string(), - profit: 12.30, - amount: 0.75, - }); - - let mut signals = self.signals.lock().await; - - signals.push(Signal { - kind: pulse_wire::terminal::SignalKind::Buy, - symbol: "BTC".to_string(), - param: pulse_wire::terminal::SignalParameter::Lim, - price: 118_800.0, - }); - - signals.push(Signal { - kind: pulse_wire::terminal::SignalKind::Buy, - symbol: "BTC".to_string(), - param: pulse_wire::terminal::SignalParameter::Tap, - price: 120_000.0, - }); - - signals.push(Signal { - kind: pulse_wire::terminal::SignalKind::Buy, - symbol: "BTC".to_string(), - param: pulse_wire::terminal::SignalParameter::Stl, - price: 118_000.0, - }); - - let mut logs = self.logs.lock().await; - - logs.push(EventLog { - kind: pulse_wire::terminal::LogKind::Warn, - name: "pulse.init".to_string(), - message: "We're still not done yet ;)".to_string(), - }); - - let mut market_overview = self.market_overview.lock().await; - - market_overview.alerts.push(Alert { - level: pulse_wire::terminal::AlertLevel::High, - message: "BTC funding rate elevated".to_string(), - }); - - market_overview.alerts.push(Alert { - level: pulse_wire::terminal::AlertLevel::Medium, - message: "Market volatility increasing".to_string(), - }); - - market_overview.alerts.push(Alert { - level: pulse_wire::terminal::AlertLevel::Low, - message: "ETH volatility returning to normal".to_string(), - }); - } + async fn init(&mut self, _ctx: &pulse_ui::state::Context) {} async fn update(&mut self, ctx: &pulse_ui::state::Context, event: Box) { if let Some(Refresh) = event.downcast_ref() { From d74d94af8963cdfe59176a2952ff6fa1ef742e9f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 23 Jul 2026 01:14:06 +0200 Subject: [PATCH 14/16] PWP Header --- src/engine/terminal.rs | 32 +++++++++++++++++++++++--------- src/terminal/terminal.rs | 25 +++++++++++++++---------- 2 files changed, 38 insertions(+), 19 deletions(-) diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 467afe4..76c480f 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -50,16 +50,21 @@ impl TerminalServer { async fn handle_client(mut reader: OwnedReadHalf) -> tokio::io::Result<()> { let input = tokio::spawn(async move { - let mut buffer = [0u8; 2048]; - loop { - let size = reader.read(&mut buffer).await?; + let mut len_buf = [0u8; size_of::()]; + let size = reader.read_exact(&mut len_buf).await?; - if size == 0 { + let len = usize::from_le_bytes(len_buf); + + if size == 0 || len == 0 { break; } - println!("Received {} bytes", size); + let mut buffer = vec![0u8; len]; + + reader.read_exact(&mut buffer).await?; + + println!("Received {} bytes", len); } Ok::<(), tokio::io::Error>(()) @@ -79,10 +84,7 @@ impl TerminalServer { let mut clients = self.clients.lock().await; for i in (0..clients.len()).rev() { - if let Err(e) = clients[i].write(&msg).await { - clients.remove(i); - println!("{e:?}"); - } else if let Err(e) = clients[i].flush().await { + if let Err(e) = Self::send_to_client(&mut clients, i, &msg).await { clients.remove(i); println!("{e:?}"); } @@ -90,4 +92,16 @@ impl TerminalServer { Ok(()) } + + pub async fn send_to_client( + clients: &mut tokio::sync::MutexGuard<'_, Vec>, + i: usize, + msg: &[u8], + ) -> tokio::io::Result<()> { + clients[i].write(&msg.len().to_le_bytes()).await?; + clients[i].write(msg).await?; + clients[i].flush().await?; + + Ok(()) + } } diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index 7c0f1f6..c58aea4 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -28,7 +28,10 @@ impl TerminalClient { &mut self, message: pulse_wire::terminal::TerminalClientMessage, ) -> tokio::io::Result<()> { - self.writer.write(&message.to_com()).await?; + let msg = message.to_com(); + self.writer.write(&msg.len().to_le_bytes()).await?; + self.writer.write(&msg).await?; + self.writer.flush().await?; Ok(()) } @@ -50,19 +53,21 @@ impl TerminalClient { tokio::spawn(async move { loop { - let mut buffer = vec![0u8; 4096]; + let mut len_buf = [0u8; size_of::()]; + reader + .read_exact(&mut len_buf) + .await + .expect("Failed to get header length"); - let len = reader - .read(&mut buffer) + let len = usize::from_le_bytes(len_buf); + + let mut buffer = vec![0u8; len]; + + reader + .read_exact(&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; From 8e83c5cd5bf0790ef3ee14e602badff150321eb6 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 23 Jul 2026 01:34:04 +0200 Subject: [PATCH 15/16] Command processing --- pulse-ui/src/state.rs | 10 +++++++- src/engine/terminal.rs | 55 ++++++++++++++++++++++++++--------------- src/terminal/command.rs | 16 ++---------- src/terminal/main.rs | 4 +-- 4 files changed, 47 insertions(+), 38 deletions(-) diff --git a/pulse-ui/src/state.rs b/pulse-ui/src/state.rs index f84f952..4d1d536 100644 --- a/pulse-ui/src/state.rs +++ b/pulse-ui/src/state.rs @@ -83,7 +83,15 @@ impl<'a, T> Drop for StateGuard<'a, T> { } impl Context { + pub async fn event(&self, event: E) { + self.tx.send(Box::new(event)).await.unwrap() + } + pub async fn close(&self) { - self.tx.send(Box::new(Close)).await.unwrap() + self.event(Close).await + } + + pub async fn refresh(&self) { + self.event(Refresh).await } } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 76c480f..9d2b2cd 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,6 +1,6 @@ use std::sync::Arc; -use pulse_wire::PulseWire; +use pulse_wire::{PulseWire, terminal::TerminalClientMessage}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::{ @@ -40,37 +40,52 @@ impl TerminalServer { self.clients.lock().await.push(writer); + let s = self.clone(); + tokio::spawn(async move { - if let Err(err) = Self::handle_client(reader).await { + if let Err(err) = s.handle_client(reader).await { eprintln!("Terminal connection error: {err}"); } }); } } - async fn handle_client(mut reader: OwnedReadHalf) -> tokio::io::Result<()> { - let input = tokio::spawn(async move { - loop { - let mut len_buf = [0u8; size_of::()]; - let size = reader.read_exact(&mut len_buf).await?; + async fn handle_client(&self, mut reader: OwnedReadHalf) -> tokio::io::Result<()> { + loop { + let mut len_buf = [0u8; size_of::()]; + let size = reader.read_exact(&mut len_buf).await?; - let len = usize::from_le_bytes(len_buf); + let len = usize::from_le_bytes(len_buf); - if size == 0 || len == 0 { - break; - } - - let mut buffer = vec![0u8; len]; - - reader.read_exact(&mut buffer).await?; - - println!("Received {} bytes", len); + if size == 0 || len == 0 { + break; } - Ok::<(), tokio::io::Error>(()) - }); + let mut buffer = vec![0u8; len]; - let _ = tokio::try_join!(input)?; + reader.read_exact(&mut buffer).await?; + + match TerminalClientMessage::from_com(&mut buffer) { + TerminalClientMessage::ExecuteCommand(command) => { + let command = command.as_str(); + + let (command, _args) = if let Some((command, args)) = command.split_once(" ") { + (command, args.split(" ").collect()) + } else { + (command, Vec::new()) + }; + + 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), + }, + )) + .await?; + } + } + } Ok(()) } diff --git a/src/terminal/command.rs b/src/terminal/command.rs index 3eecd8a..86c0c99 100644 --- a/src/terminal/command.rs +++ b/src/terminal/command.rs @@ -6,24 +6,12 @@ impl PulseTradeApp { return; } - let (command, _args) = if let Some((command, args)) = command.split_once(" ") { - (command, args.split(" ").collect()) - } else { - (command, Vec::new()) - }; - - match command { + match command.trim() { "exit" | "quit" | "leave" => { ctx.close().await; } - _ => { - self.logs.lock().await.push(pulse_wire::terminal::EventLog { - kind: pulse_wire::terminal::LogKind::Err, - name: "cmd".to_string(), - message: format!("Command '{}' not found", command), - }); - } + _ => {} } } } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 13f9366..4bd019a 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -43,9 +43,7 @@ impl App for PulseTradeApp { async fn update(&mut self, ctx: &pulse_ui::state::Context, event: Box) { if let Some(Refresh) = event.downcast_ref() { return; - } - - if let Some(event) = event.downcast_ref() { + } else if let Some(event) = event.downcast_ref() { if self.command.value.lock().await.handle_event(event) { return; } else if let crossterm::event::Event::Key(key) = event { From df8b5b5c419fd5859ce927ae82da424c7727fe70 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 23 Jul 2026 01:42:03 +0200 Subject: [PATCH 16/16] Socket state --- src/terminal/command.rs | 14 +++++++++++++- src/terminal/main.rs | 11 ++++------- src/terminal/terminal.rs | 4 +++- 3 files changed, 20 insertions(+), 9 deletions(-) diff --git a/src/terminal/command.rs b/src/terminal/command.rs index 86c0c99..b8849df 100644 --- a/src/terminal/command.rs +++ b/src/terminal/command.rs @@ -11,7 +11,19 @@ impl PulseTradeApp { ctx.close().await; } - _ => {} + _ => { + if let Err(e) = self + .sock + .as_mut() + .unwrap() + .send(pulse_wire::terminal::TerminalClientMessage::ExecuteCommand( + command.to_string(), + )) + .await + { + eprintln!("{e}"); + } + } } } } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 4bd019a..c0d3faf 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -26,6 +26,8 @@ use pulse_wire::terminal::{ }; pub struct PulseTradeApp { + sock: Option, + command: State, scroll: State>, watch_list: State>, @@ -169,16 +171,11 @@ impl App for PulseTradeApp { #[tokio::main] async fn main() -> tokio::io::Result<()> { - let mut client = terminal::TerminalClient::new().await?; - - client - .send(pulse_wire::terminal::TerminalClientMessage::ExecuteCommand( - "Hello".to_string(), - )) - .await?; + let client = terminal::TerminalClient::new().await?; pulse_ui::run(|ctx| { client.use_app(PulseTradeApp { + sock: None, command: ctx.use_state(InputState::new()), scroll: ctx.use_state(ScrollState(1, [0; 7])), watch_list: ctx.use_state(Vec::new()), diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index c58aea4..617979d 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -36,7 +36,7 @@ impl TerminalClient { Ok(()) } - pub fn use_app(&mut self, app: PulseTradeApp) -> PulseTradeApp { + pub fn use_app(mut self, mut app: PulseTradeApp) -> PulseTradeApp { let mut reader = None; std::mem::swap(&mut self.reader, &mut reader); @@ -104,6 +104,8 @@ impl TerminalClient { } }); + app.sock = Some(self); + app } }