From 7e79390052d339e8965a6baf20452d8833847ea2 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 25 Jul 2026 19:06:11 +0200 Subject: [PATCH 01/25] Using percent for trend --- src/terminal/formatting.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 43a6ffc..ed76ef9 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -79,7 +79,7 @@ impl Formatted for MarketItem { self.price.to_string(), self.volume_24h.to_string(), format!( - "{} {}", + "{} {}%", if self.trend.is_sign_positive() { "\x1b[32m▲\x1b[0m" } else { From 90876c83a78f03ae04dd12d1f122e585fc10ee06 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 25 Jul 2026 19:23:04 +0200 Subject: [PATCH 02/25] Strategy messaging --- pulse-wire/src/lib.rs | 18 ++++++++++ pulse-wire/src/strategy.rs | 73 +++++++++++++++++++++++++++++++++++++- 2 files changed, 90 insertions(+), 1 deletion(-) diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index 689db77..ff1a69c 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -52,6 +52,24 @@ impl PulseWire for Vec { } } +impl PulseWire for Option { + fn to_com(&self) -> Vec { + if let Some(v) = self { + v.to_com() + } else { + vec![0] + } + } + + fn from_com(com: &mut Vec) -> Self { + if com[0] > 0 { + Some(T::from_com(com)) + } else { + None + } + } +} + impl PulseWire for String { fn to_com(&self) -> Vec { let bytes = self.as_bytes(); diff --git a/pulse-wire/src/strategy.rs b/pulse-wire/src/strategy.rs index 757ac70..07d2419 100644 --- a/pulse-wire/src/strategy.rs +++ b/pulse-wire/src/strategy.rs @@ -1,4 +1,9 @@ -use crate::{PulseWire, units::TimeFrame}; +use crate::{ + PulseWire, + general::{EventLog, Signal}, + terminal::MarketItem, + units::{Direction, TimeFrame}, +}; use pulse_macros::pwp; #[pwp] @@ -23,3 +28,69 @@ pub struct RiskManifest { max_loss: u8, cooldown: TimeFrame, } + +#[pwp] +pub enum StrategyMessage { + RequestOHLC { + symbol: String, + timeframe: TimeFrame, + count: u32, + }, + + SubscribeCandle { + symbol: String, + timeframe: TimeFrame, + }, + + Signal(StrategySignal), + + Log(EventLog), +} + +#[pwp] +pub enum StrategyEngineMessage { + Initialize { + watchlist: Vec, + }, + + CandleUpdate { + symbol: String, + candle: String, + }, + + OHLC { + symbol: String, + timeframe: TimeFrame, + candles: Vec, + }, + + Start, + + Stop, +} + +#[pwp] +pub enum RiskMessage { + Approve(Signal), + + Reject { reason: String }, + + Log(EventLog), +} + +#[pwp] +pub enum RiskEngineMessage { + Initialize { strategy: String }, + + Signal(StrategySignal), + + MarketUpdate { symbol: String, price: f64 }, +} + +#[pwp] +pub struct StrategySignal { + pub symbol: String, + pub side: Direction, + pub confidence: f32, + pub price: Option, +} From 1b3025ee34c372744d04daf2d857f6cc2d4802f1 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 25 Jul 2026 19:45:40 +0200 Subject: [PATCH 03/25] Plugin --- pulse-wire/src/lib.rs | 4 ++-- pulse-wire/src/{strategy.rs => plugin.rs} | 0 pulse-wire/src/terminal.rs | 2 +- src/engine/store/mod.rs | 9 +++++++++ src/engine/store/plugin.rs | 1 + 5 files changed, 13 insertions(+), 3 deletions(-) rename pulse-wire/src/{strategy.rs => plugin.rs} (100%) create mode 100644 src/engine/store/plugin.rs diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index ff1a69c..0d8b6d8 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -2,7 +2,7 @@ use std::path::PathBuf; pub mod general; -pub mod strategy; +pub mod plugin; pub mod terminal; pub mod units; @@ -10,7 +10,7 @@ pub mod prelude { pub use crate::PulseWire; pub use crate::general::*; pub use crate::server_path; - pub use crate::strategy::*; + pub use crate::plugin::*; pub use crate::terminal::*; pub use crate::units::*; } diff --git a/pulse-wire/src/strategy.rs b/pulse-wire/src/plugin.rs similarity index 100% rename from pulse-wire/src/strategy.rs rename to pulse-wire/src/plugin.rs diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index 5b26c2c..d59b3fc 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -1,7 +1,7 @@ use crate::{ PulseWire, general::{EventLog, MarketTrend, Position, Signal}, - strategy::{RiskManifest, StrategyManifest}, + plugin::{RiskManifest, StrategyManifest}, units::{Symbol, TimeFrame, USD, Volatility}, }; use pulse_macros::pwp; diff --git a/src/engine/store/mod.rs b/src/engine/store/mod.rs index 5299a77..73b6dee 100644 --- a/src/engine/store/mod.rs +++ b/src/engine/store/mod.rs @@ -2,6 +2,7 @@ use std::path::PathBuf; pub mod accounts; pub mod config; +pub mod plugin; pub fn home_dir() -> tokio::io::Result { std::env::home_dir().ok_or_else(|| { @@ -13,6 +14,14 @@ pub fn pulse_directory() -> tokio::io::Result { Ok(home_dir()?.join(".config").join("pulse-trader")) } +pub fn pulse_plugins_directory() -> tokio::io::Result { + Ok(pulse_directory()?.join("plugins")) +} + +pub fn pulse_plugin(id: &str) -> tokio::io::Result { + Ok(pulse_directory()?.join("plugins").join(id)) +} + pub fn pulse_config_file() -> tokio::io::Result { Ok(pulse_directory()?.join("config.toml")) } diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs new file mode 100644 index 0000000..5fd4ca3 --- /dev/null +++ b/src/engine/store/plugin.rs @@ -0,0 +1 @@ +pub struct Plugin(); From 87e9f62c8db82fd06db42a52004751ab3959a565 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 25 Jul 2026 20:09:18 +0200 Subject: [PATCH 04/25] Plugin child process --- Cargo.lock | 1 + Cargo.toml | 9 +----- src/engine/store/plugin.rs | 60 +++++++++++++++++++++++++++++++++++++- 3 files changed, 61 insertions(+), 9 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8495645..3cf15f6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5200,6 +5200,7 @@ dependencies = [ "libc", "mio", "pin-project-lite", + "signal-hook-registry", "socket2", "tokio-macros", "windows-sys 0.61.2", diff --git a/Cargo.toml b/Cargo.toml index 8780172..8b87ecb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,14 +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", - "time", -] } +tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time", "process"] } crossterm = { workspace = true } hypersdk = "0.2.14" serde_json = "1" diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 5fd4ca3..57d03de 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1 +1,59 @@ -pub struct Plugin(); +use std::marker::PhantomData; + +use pulse_wire::{ + PulseWire, + plugin::{RiskEngineMessage, RiskMessage, StrategyEngineMessage, StrategyMessage}, +}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; + +pub struct Plugin { + pub process: tokio::process::Child, + pub _p: (PhantomData, PhantomData), +} + +impl Plugin { + pub fn new(child: tokio::process::Child) -> Self { + Self { + process: child, + _p: (PhantomData, PhantomData), + } + } + + pub async fn recv(&mut self) -> tokio::io::Result> { + let stderr = self.process.stderr.as_mut().unwrap(); + + let mut len_buf = [0u8; size_of::()]; + let size = stderr.read_exact(&mut len_buf).await?; + + let len = usize::from_le_bytes(len_buf); + + if size == 0 || len == 0 { + return Ok(None); + } + + let mut buffer = vec![0u8; len]; + + stderr.read_exact(&mut buffer).await?; + + Ok(Some(R::from_com(&mut buffer))) + } + + pub async fn send(&mut self, msg: &S) -> tokio::io::Result<()> { + self.send_raw(&msg.to_com()).await + } + + pub async fn send_raw(&mut self, msg: &[u8]) -> tokio::io::Result<()> { + let stdin = self.process.stdin.as_mut().unwrap(); + + stdin.write(&msg.len().to_le_bytes()).await?; + stdin.write(msg).await?; + stdin.flush().await?; + + Ok(()) + } +} + +pub struct StrategyPair { + pub strategy: Plugin, + pub risk: Plugin, +} From ff5704ec88052c19a2300479490f8bd8678bbcac Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 25 Jul 2026 23:48:17 +0200 Subject: [PATCH 05/25] Strategy spawning and maniest --- src/engine/engine.rs | 10 ++++++-- src/engine/store/config.rs | 2 ++ src/engine/store/plugin.rs | 47 ++++++++++++++++++++++++++++++++++++-- 3 files changed, 55 insertions(+), 4 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index 438b06c..52a5f2c 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -1,5 +1,5 @@ use crate::{ - store::{accounts::AccountList, config::Config}, + store::{accounts::AccountList, config::Config, plugin::StrategyPair}, terminal::TerminalServer, }; use pulse_wire::prelude::*; @@ -11,17 +11,23 @@ pub struct Engine { pub terminal_server: Arc, pub config: Arc>, pub accounts: Arc>, + pub strategy: Arc>, } impl Engine { pub async fn new() -> tokio::io::Result> { - let config = Arc::new(Mutex::new(Config::new().await?)); + let config = Config::new().await?; let accounts = Arc::new(Mutex::new(AccountList::new().await?)); + let strategy = Arc::new(Mutex::new( + StrategyPair::new(&config.strategy, &config.risk).await?, + )); + let config = Arc::new(Mutex::new(config)); Ok(Arc::new_cyclic(|engine| Self { terminal_server: TerminalServer::new(engine.clone()), config, accounts, + strategy, })) } diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index 2ead6ae..afaf487 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -6,6 +6,8 @@ pub struct WatchList { #[derive(Debug, Default, serde::Serialize, serde::Deserialize)] pub struct Config { pub watchlist: WatchList, + pub strategy: String, + pub risk: String, } impl Config { diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 57d03de..b014315 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -2,10 +2,20 @@ use std::marker::PhantomData; use pulse_wire::{ PulseWire, - plugin::{RiskEngineMessage, RiskMessage, StrategyEngineMessage, StrategyMessage}, + plugin::{ + RiskEngineMessage, RiskManifest, RiskMessage, StrategyEngineMessage, StrategyManifest, + StrategyMessage, + }, +}; +use tokio::{ + fs, + io::{AsyncReadExt, AsyncWriteExt}, + process::Command, }; -use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use crate::store::pulse_plugin; + +#[derive(Debug)] pub struct Plugin { pub process: tokio::process::Child, pub _p: (PhantomData, PhantomData), @@ -53,7 +63,40 @@ impl Plugin { } } +#[derive(Debug)] pub struct StrategyPair { pub strategy: Plugin, pub risk: Plugin, + + pub strategy_manifest: StrategyManifest, + pub risk_manifest: RiskManifest, +} + +impl StrategyPair { + pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { + let strategy = pulse_plugin(strategy_id)?; + let risk = pulse_plugin(risk_id)?; + + Ok(Self { + strategy: Plugin::new( + Command::new("bash") + .arg(strategy.join("strategy.bash")) + .current_dir(&strategy) + .spawn()?, + ), + risk: Plugin::new( + Command::new("bash") + .arg(strategy.join("risk.bash")) + .current_dir(&strategy) + .spawn()?, + ), + strategy_manifest: toml::from_slice(&fs::read(strategy.join("strategy.toml")).await?) + .map_err(|v| { + tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()) + })?, + risk_manifest: toml::from_slice(&fs::read(risk.join("risk.toml")).await?).map_err( + |v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()), + )?, + }) + } } From 270a0f18339039b355fcdfdfdc204b71858f6115 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 01:56:46 +0200 Subject: [PATCH 06/25] Strategy engine --- src/engine/engine.rs | 8 +++----- src/engine/store/plugin.rs | 40 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 43 insertions(+), 5 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index 52a5f2c..c3ce821 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -1,5 +1,5 @@ use crate::{ - store::{accounts::AccountList, config::Config, plugin::StrategyPair}, + store::{accounts::AccountList, config::Config, plugin::StrategyEngine}, terminal::TerminalServer, }; use pulse_wire::prelude::*; @@ -11,16 +11,14 @@ pub struct Engine { pub terminal_server: Arc, pub config: Arc>, pub accounts: Arc>, - pub strategy: Arc>, + pub strategy: Arc, } impl Engine { pub async fn new() -> tokio::io::Result> { let config = Config::new().await?; let accounts = Arc::new(Mutex::new(AccountList::new().await?)); - let strategy = Arc::new(Mutex::new( - StrategyPair::new(&config.strategy, &config.risk).await?, - )); + let strategy = Arc::new(StrategyEngine::new(&config.strategy, &config.risk).await?); let config = Arc::new(Mutex::new(config)); Ok(Arc::new_cyclic(|engine| Self { diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index b014315..a589dfe 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -7,10 +7,13 @@ use pulse_wire::{ StrategyMessage, }, }; + use tokio::{ fs, io::{AsyncReadExt, AsyncWriteExt}, process::Command, + sync::Mutex, + task::JoinHandle, }; use crate::store::pulse_plugin; @@ -100,3 +103,40 @@ impl StrategyPair { }) } } + +#[derive(Debug)] +pub struct StrategyEngine { + pub pair: Mutex, + + pub strategy_handle: Mutex>>, + pub risk_handle: Mutex>>, +} + +impl StrategyEngine { + pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { + Ok(Self { + pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), + strategy_handle: Mutex::new(None), + risk_handle: Mutex::new(None), + }) + } + + pub async fn reload(&self, strategy_id: &str, risk_id: &str) -> tokio::io::Result<()> { + if let Some(handle) = &*self.strategy_handle.lock().await { + handle.abort(); + } + + if let Some(handle) = &*self.risk_handle.lock().await { + handle.abort(); + } + + let mut pair = self.pair.lock().await; + + pair.strategy.process.kill().await?; + pair.risk.process.kill().await?; + + *pair = StrategyPair::new(strategy_id, risk_id).await?; + + Ok(()) + } +} From 73a47c43f73b2cc3512c523d1778819116275526 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 02:38:51 +0200 Subject: [PATCH 07/25] Running strategy engine --- src/engine/engine.rs | 8 +++--- src/engine/main.rs | 15 +++++++++++ src/engine/store/plugin.rs | 53 +++++++++++++++++++++++++++++++++++--- 3 files changed, 69 insertions(+), 7 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index c3ce821..798ebb3 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -9,23 +9,25 @@ use tokio::{sync::Mutex, task::JoinHandle}; #[derive(Debug, Clone)] pub struct Engine { pub terminal_server: Arc, + pub strategy: Arc, pub config: Arc>, pub accounts: Arc>, - pub strategy: Arc, } impl Engine { pub async fn new() -> tokio::io::Result> { let config = Config::new().await?; + + let strategy = StrategyEngine::new(&config.strategy, &config.risk).await?; + let accounts = Arc::new(Mutex::new(AccountList::new().await?)); - let strategy = Arc::new(StrategyEngine::new(&config.strategy, &config.risk).await?); let config = Arc::new(Mutex::new(config)); Ok(Arc::new_cyclic(|engine| Self { terminal_server: TerminalServer::new(engine.clone()), + strategy: strategy.initialize(engine.clone()), config, accounts, - strategy, })) } diff --git a/src/engine/main.rs b/src/engine/main.rs index f4f7f64..184a811 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -11,6 +11,8 @@ async fn main() -> tokio::io::Result<()> { let broadcaster = engine.spawn_broadcaster().await; + engine.strategy.spawn().await; + engine.run_engine().await?; let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster); @@ -18,5 +20,18 @@ async fn main() -> tokio::io::Result<()> { terminal_server??; broadcaster??; + let strategy = engine.strategy.strategy_handle.lock().await.take(); + let risk = engine.strategy.risk_handle.lock().await.take(); + + drop(engine); + + if let Some(strategy) = strategy { + strategy.into_future().await??; + } + + if let Some(risk) = risk { + risk.into_future().await??; + } + Ok(()) } diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index a589dfe..de0fe97 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1,4 +1,7 @@ -use std::marker::PhantomData; +use std::{ + marker::PhantomData, + sync::{Arc, Weak}, +}; use pulse_wire::{ PulseWire, @@ -16,7 +19,7 @@ use tokio::{ task::JoinHandle, }; -use crate::store::pulse_plugin; +use crate::{engine::Engine, store::pulse_plugin}; #[derive(Debug)] pub struct Plugin { @@ -108,8 +111,10 @@ impl StrategyPair { pub struct StrategyEngine { pub pair: Mutex, - pub strategy_handle: Mutex>>, - pub risk_handle: Mutex>>, + pub strategy_handle: Mutex>>>, + pub risk_handle: Mutex>>>, + + pub engine: Weak, } impl StrategyEngine { @@ -118,9 +123,49 @@ impl StrategyEngine { pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), strategy_handle: Mutex::new(None), risk_handle: Mutex::new(None), + engine: Weak::new(), }) } + pub fn initialize(mut self, engine: Weak) -> Arc { + self.engine = engine; + + Arc::new(self) + } + + pub async fn run_strategy(&self) -> tokio::io::Result<()> { + let engine = self + .engine + .upgrade() + .expect("Failed to upgrade engine (StrategyEngine)"); + + loop {} + + Ok(()) + } + + pub async fn run_risk(&self) -> tokio::io::Result<()> { + let engine = self + .engine + .upgrade() + .expect("Failed to upgrade engine (StrategyEngine)"); + + loop {} + + Ok(()) + } + + pub async fn spawn(self: &Arc) { + let engine = self.clone(); + + *self.strategy_handle.lock().await = + Some(tokio::spawn(async move { engine.run_strategy().await })); + + let engine = self.clone(); + + *self.risk_handle.lock().await = Some(tokio::spawn(async move { engine.run_risk().await })); + } + pub async fn reload(&self, strategy_id: &str, risk_id: &str) -> tokio::io::Result<()> { if let Some(handle) = &*self.strategy_handle.lock().await { handle.abort(); From 4bf47afaa404075aef7b62c16e8e7ab214519c57 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 02:56:08 +0200 Subject: [PATCH 08/25] Fix deadlock --- src/engine/store/plugin.rs | 85 ++++++++++++++++++++++++++------------ 1 file changed, 58 insertions(+), 27 deletions(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index de0fe97..3b2bcb7 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -71,11 +71,11 @@ impl Plugin { #[derive(Debug)] pub struct StrategyPair { - pub strategy: Plugin, - pub risk: Plugin, + pub strategy: Mutex>, + pub risk: Mutex>, - pub strategy_manifest: StrategyManifest, - pub risk_manifest: RiskManifest, + pub strategy_manifest: Mutex, + pub risk_manifest: Mutex, } impl StrategyPair { @@ -84,32 +84,35 @@ impl StrategyPair { let risk = pulse_plugin(risk_id)?; Ok(Self { - strategy: Plugin::new( + strategy: Mutex::new(Plugin::new( Command::new("bash") .arg(strategy.join("strategy.bash")) .current_dir(&strategy) .spawn()?, - ), - risk: Plugin::new( + )), + risk: Mutex::new(Plugin::new( Command::new("bash") .arg(strategy.join("risk.bash")) .current_dir(&strategy) .spawn()?, + )), + strategy_manifest: Mutex::new( + toml::from_slice(&fs::read(strategy.join("strategy.toml")).await?).map_err( + |v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()), + )?, + ), + risk_manifest: Mutex::new( + toml::from_slice(&fs::read(risk.join("risk.toml")).await?).map_err(|v| { + tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()) + })?, ), - strategy_manifest: toml::from_slice(&fs::read(strategy.join("strategy.toml")).await?) - .map_err(|v| { - tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()) - })?, - risk_manifest: toml::from_slice(&fs::read(risk.join("risk.toml")).await?).map_err( - |v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()), - )?, }) } } #[derive(Debug)] pub struct StrategyEngine { - pub pair: Mutex, + pub pair: StrategyPair, pub strategy_handle: Mutex>>>, pub risk_handle: Mutex>>>, @@ -120,7 +123,7 @@ pub struct StrategyEngine { impl StrategyEngine { pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { Ok(Self { - pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), + pair: StrategyPair::new(strategy_id, risk_id).await?, strategy_handle: Mutex::new(None), risk_handle: Mutex::new(None), engine: Weak::new(), @@ -166,21 +169,49 @@ impl StrategyEngine { *self.risk_handle.lock().await = Some(tokio::spawn(async move { engine.run_risk().await })); } - pub async fn reload(&self, strategy_id: &str, risk_id: &str) -> tokio::io::Result<()> { - if let Some(handle) = &*self.strategy_handle.lock().await { - handle.abort(); + pub async fn reload( + self: &Arc, + strategy_id: &str, + risk_id: &str, + ) -> tokio::io::Result<()> { + { + if let Some(handle) = &*self.strategy_handle.lock().await { + handle.abort(); + } + + if let Some(handle) = &*self.risk_handle.lock().await { + handle.abort(); + } + + self.pair.strategy.lock().await.process.kill().await?; + self.pair.risk.lock().await.process.kill().await?; } - if let Some(handle) = &*self.risk_handle.lock().await { - handle.abort(); + { + let strategy = StrategyPair::new(strategy_id, risk_id).await?; + + std::mem::swap( + &mut *self.pair.strategy.lock().await, + &mut *strategy.strategy.lock().await, + ); + + std::mem::swap( + &mut *self.pair.strategy_manifest.lock().await, + &mut *strategy.strategy_manifest.lock().await, + ); + + std::mem::swap( + &mut *self.pair.risk.lock().await, + &mut *strategy.risk.lock().await, + ); + + std::mem::swap( + &mut *self.pair.risk_manifest.lock().await, + &mut *strategy.risk_manifest.lock().await, + ); } - let mut pair = self.pair.lock().await; - - pair.strategy.process.kill().await?; - pair.risk.process.kill().await?; - - *pair = StrategyPair::new(strategy_id, risk_id).await?; + self.spawn().await; Ok(()) } From 7df1902a694059c9cb2ed73f301c25cf2dee2b56 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 04:29:33 +0200 Subject: [PATCH 09/25] Hopefully fixed deadlocks in strategy engine --- src/engine/store/plugin.rs | 77 ++++++++++++++------------------------ 1 file changed, 29 insertions(+), 48 deletions(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 3b2bcb7..5d0698f 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -14,7 +14,7 @@ use pulse_wire::{ use tokio::{ fs, io::{AsyncReadExt, AsyncWriteExt}, - process::Command, + process::{Child, ChildStdout, Command}, sync::Mutex, task::JoinHandle, }; @@ -23,23 +23,25 @@ use crate::{engine::Engine, store::pulse_plugin}; #[derive(Debug)] pub struct Plugin { - pub process: tokio::process::Child, + pub stdout: Mutex, + pub process: Mutex, pub _p: (PhantomData, PhantomData), } impl Plugin { - pub fn new(child: tokio::process::Child) -> Self { + pub fn new(mut child: Child) -> Self { Self { - process: child, + stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")), + process: Mutex::new(child), _p: (PhantomData, PhantomData), } } - pub async fn recv(&mut self) -> tokio::io::Result> { - let stderr = self.process.stderr.as_mut().unwrap(); + pub async fn recv(&self) -> tokio::io::Result> { + let mut stdout = self.stdout.lock().await; let mut len_buf = [0u8; size_of::()]; - let size = stderr.read_exact(&mut len_buf).await?; + let size = stdout.read_exact(&mut len_buf).await?; let len = usize::from_le_bytes(len_buf); @@ -49,7 +51,7 @@ impl Plugin { let mut buffer = vec![0u8; len]; - stderr.read_exact(&mut buffer).await?; + stdout.read_exact(&mut buffer).await?; Ok(Some(R::from_com(&mut buffer))) } @@ -59,7 +61,8 @@ impl Plugin { } pub async fn send_raw(&mut self, msg: &[u8]) -> tokio::io::Result<()> { - let stdin = self.process.stdin.as_mut().unwrap(); + let mut process = self.process.lock().await; + let stdin = process.stdin.as_mut().unwrap(); stdin.write(&msg.len().to_le_bytes()).await?; stdin.write(msg).await?; @@ -71,31 +74,31 @@ impl Plugin { #[derive(Debug)] pub struct StrategyPair { - pub strategy: Mutex>, - pub risk: Mutex>, + pub strategy: Plugin, + pub risk: Plugin, pub strategy_manifest: Mutex, pub risk_manifest: Mutex, } impl StrategyPair { - pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { + pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result> { let strategy = pulse_plugin(strategy_id)?; let risk = pulse_plugin(risk_id)?; - Ok(Self { - strategy: Mutex::new(Plugin::new( + Ok(Arc::new(Self { + strategy: Plugin::new( Command::new("bash") .arg(strategy.join("strategy.bash")) .current_dir(&strategy) .spawn()?, - )), - risk: Mutex::new(Plugin::new( + ), + risk: Plugin::new( Command::new("bash") .arg(strategy.join("risk.bash")) .current_dir(&strategy) .spawn()?, - )), + ), strategy_manifest: Mutex::new( toml::from_slice(&fs::read(strategy.join("strategy.toml")).await?).map_err( |v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()), @@ -106,13 +109,13 @@ impl StrategyPair { tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()) })?, ), - }) + })) } } #[derive(Debug)] pub struct StrategyEngine { - pub pair: StrategyPair, + pub pair: Mutex>, pub strategy_handle: Mutex>>>, pub risk_handle: Mutex>>>, @@ -123,7 +126,7 @@ pub struct StrategyEngine { impl StrategyEngine { pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { Ok(Self { - pair: StrategyPair::new(strategy_id, risk_id).await?, + pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), strategy_handle: Mutex::new(None), risk_handle: Mutex::new(None), engine: Weak::new(), @@ -143,8 +146,6 @@ impl StrategyEngine { .expect("Failed to upgrade engine (StrategyEngine)"); loop {} - - Ok(()) } pub async fn run_risk(&self) -> tokio::io::Result<()> { @@ -154,8 +155,6 @@ impl StrategyEngine { .expect("Failed to upgrade engine (StrategyEngine)"); loop {} - - Ok(()) } pub async fn spawn(self: &Arc) { @@ -183,32 +182,14 @@ impl StrategyEngine { handle.abort(); } - self.pair.strategy.lock().await.process.kill().await?; - self.pair.risk.lock().await.process.kill().await?; - } + { + let pair = self.pair.lock().await.clone(); - { - let strategy = StrategyPair::new(strategy_id, risk_id).await?; + pair.strategy.process.lock().await.kill().await?; + pair.risk.process.lock().await.kill().await?; + } - std::mem::swap( - &mut *self.pair.strategy.lock().await, - &mut *strategy.strategy.lock().await, - ); - - std::mem::swap( - &mut *self.pair.strategy_manifest.lock().await, - &mut *strategy.strategy_manifest.lock().await, - ); - - std::mem::swap( - &mut *self.pair.risk.lock().await, - &mut *strategy.risk.lock().await, - ); - - std::mem::swap( - &mut *self.pair.risk_manifest.lock().await, - &mut *strategy.risk_manifest.lock().await, - ); + *self.pair.lock().await = StrategyPair::new(strategy_id, risk_id).await?; } self.spawn().await; From 785bcd9399d6293853b3598518e5d9041b887312 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 04:48:04 +0200 Subject: [PATCH 10/25] Testing recv --- src/engine/store/plugin.rs | 18 +++++++++++++++++- src/engine/terminal.rs | 7 +++++++ 2 files changed, 24 insertions(+), 1 deletion(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 5d0698f..e1d876f 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -145,7 +145,23 @@ impl StrategyEngine { .upgrade() .expect("Failed to upgrade engine (StrategyEngine)"); - loop {} + let pair = self.pair.lock().await.clone(); + + loop { + match pair.strategy.recv().await? { + Some(StrategyMessage::Log(log)) => { + engine.terminal_server.log_raw(log).await?; + } + Some(StrategyMessage::RequestOHLC { + symbol, + timeframe, + count, + }) => {} + Some(StrategyMessage::Signal(signal)) => {} + Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} + None => {} + } + } } pub async fn run_risk(&self) -> tokio::io::Result<()> { diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 7de2c44..62170fa 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -169,6 +169,13 @@ impl TerminalServer { .await } + pub async fn log_raw(self: &Arc, log: EventLog) -> tokio::io::Result<()> { + 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) -> tokio::io::Result<()> { self.log(LogKind::Info, name, message).await } From ec7c3e0c3f5b63951203f62146dd3858df0d549c Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 05:16:20 +0200 Subject: [PATCH 11/25] Fix std error --- src/engine/store/plugin.rs | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index e1d876f..9ab4e0c 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1,5 +1,6 @@ use std::{ marker::PhantomData, + process::Stdio, sync::{Arc, Weak}, }; @@ -91,12 +92,18 @@ impl StrategyPair { Command::new("bash") .arg(strategy.join("strategy.bash")) .current_dir(&strategy) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::inherit()) .spawn()?, ), risk: Plugin::new( Command::new("bash") - .arg(strategy.join("risk.bash")) - .current_dir(&strategy) + .arg(risk.join("risk.bash")) + .current_dir(&risk) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::inherit()) .spawn()?, ), strategy_manifest: Mutex::new( @@ -170,7 +177,9 @@ impl StrategyEngine { .upgrade() .expect("Failed to upgrade engine (StrategyEngine)"); - loop {} + // loop {} + + Ok(()) } pub async fn spawn(self: &Arc) { From 1429133b1e8a80e055479254260b3813972f0dfa Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 14:39:16 +0200 Subject: [PATCH 12/25] Send strategy manifest to clients --- src/engine/terminal.rs | 31 ++++++++++++++++++++++++++----- src/terminal/main.rs | 27 +-------------------------- 2 files changed, 27 insertions(+), 31 deletions(-) diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 62170fa..24849e1 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -59,17 +59,38 @@ impl TerminalServer { } } - async fn handle_client( - self: &Arc, - id: &usize, - mut reader: OwnedReadHalf, - ) -> tokio::io::Result<()> { + async fn initialize_client(self: &Arc, id: &usize) -> tokio::io::Result<()> { + let engine = self.get_engine(); + let pair = engine.strategy.pair.lock().await; + + self.send_to( + id, + pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { + strategy: pair.strategy_manifest.lock().await.clone(), + risk: pair.risk_manifest.lock().await.clone(), + mode: Mode::Auto, + state: ItemState::Running, + cooldown: TimeFrame::M15, + }), + ) + .await?; + self.send_to( id, pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), ) .await?; + Ok(()) + } + + async fn handle_client( + self: &Arc, + id: &usize, + mut reader: OwnedReadHalf, + ) -> tokio::io::Result<()> { + self.initialize_client(id).await?; + loop { let mut len_buf = [0u8; size_of::()]; let size = reader.read_exact(&mut len_buf).await?; diff --git a/src/terminal/main.rs b/src/terminal/main.rs index da08d43..e1ccf47 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -181,32 +181,7 @@ async fn main() -> tokio::io::Result<()> { signals: ctx.use_state(Vec::new()), logs: ctx.use_state(Vec::new()), inspect: ctx.use_state(InspectTarget::None), - strategy: ctx.use_state(Some(Strategy { - strategy: StrategyManifest { - name: "Liquidity Sweep".to_string(), - description: - "Detects liquidity grabs around key support and resistance levels." - .to_string(), - author: "Klesty Selimaj".to_string(), - version: "1.0.0".to_string(), - - timeframes: vec![TimeFrame::M5, TimeFrame::M15, TimeFrame::H1], - }, - - risk: RiskManifest { - name: "Aggressive".to_string(), - description: "High-risk profile with larger position sizing.".to_string(), - author: "Klesty Selimaj".to_string(), - version: "1.0.0".to_string(), - - max_loss: 5, - cooldown: TimeFrame::M15, - }, - - mode: Mode::Auto, - state: ItemState::Running, - cooldown: TimeFrame::M15, - })), + strategy: ctx.use_state(None), status: ctx.use_state(None), }) }) From 0d1944258dbde8d1a2b16daa2a27c2d6bdca8bec Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 15:31:29 +0200 Subject: [PATCH 13/25] Improved reload --- src/engine/store/plugin.rs | 111 ++++++++++++++++++++++--------------- 1 file changed, 66 insertions(+), 45 deletions(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 9ab4e0c..b891aae 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1,5 +1,6 @@ use std::{ marker::PhantomData, + path::PathBuf, process::Stdio, sync::{Arc, Weak}, }; @@ -12,6 +13,7 @@ use pulse_wire::{ }, }; +use serde::Deserialize; use tokio::{ fs, io::{AsyncReadExt, AsyncWriteExt}, @@ -71,6 +73,36 @@ impl Plugin { Ok(()) } + + pub async fn reload(&self, mut child: Child) -> tokio::io::Result<()> { + let mut process = self.process.lock().await; + + process.kill().await?; + + *self.stdout.lock().await = child.stdout.take().expect("Failed to obtain child stdout"); + + *process = child; + + Ok(()) + } +} + +pub fn get_manifest_plugin_pair<'de, M: Deserialize<'de>>( + plugin_dir: &PathBuf, + plugin_path: &PathBuf, + manifest: &'de [u8], +) -> tokio::io::Result<(Child, M)> { + Ok(( + Command::new("bash") + .arg(plugin_path) + .current_dir(plugin_dir) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::inherit()) + .spawn()?, + toml::from_slice(manifest) + .map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))?, + )) } #[derive(Debug)] @@ -87,35 +119,23 @@ impl StrategyPair { let strategy = pulse_plugin(strategy_id)?; let risk = pulse_plugin(risk_id)?; + let (strategy, strategy_manifest) = get_manifest_plugin_pair( + &strategy, + &strategy.join("strategy.bash"), + &fs::read(strategy.join("strategy.toml")).await?, + )?; + + let (risk, risk_manifest) = get_manifest_plugin_pair( + &risk, + &risk.join("risk.bash"), + &fs::read(risk.join("risk.toml")).await?, + )?; + Ok(Arc::new(Self { - strategy: Plugin::new( - Command::new("bash") - .arg(strategy.join("strategy.bash")) - .current_dir(&strategy) - .stdin(Stdio::piped()) - .stdout(Stdio::piped()) - .stderr(Stdio::inherit()) - .spawn()?, - ), - risk: Plugin::new( - Command::new("bash") - .arg(risk.join("risk.bash")) - .current_dir(&risk) - .stdin(Stdio::piped()) - .stdout(Stdio::piped()) - .stderr(Stdio::inherit()) - .spawn()?, - ), - strategy_manifest: Mutex::new( - toml::from_slice(&fs::read(strategy.join("strategy.toml")).await?).map_err( - |v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()), - )?, - ), - risk_manifest: Mutex::new( - toml::from_slice(&fs::read(risk.join("risk.toml")).await?).map_err(|v| { - tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()) - })?, - ), + strategy: Plugin::new(strategy), + risk: Plugin::new(risk), + strategy_manifest: Mutex::new(strategy_manifest), + risk_manifest: Mutex::new(risk_manifest), })) } } @@ -156,17 +176,21 @@ impl StrategyEngine { loop { match pair.strategy.recv().await? { + None => {} + Some(StrategyMessage::Log(log)) => { engine.terminal_server.log_raw(log).await?; } + + Some(StrategyMessage::Signal(signal)) => {} + + Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} + Some(StrategyMessage::RequestOHLC { symbol, timeframe, count, }) => {} - Some(StrategyMessage::Signal(signal)) => {} - Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} - None => {} } } } @@ -193,28 +217,25 @@ impl StrategyEngine { *self.risk_handle.lock().await = Some(tokio::spawn(async move { engine.run_risk().await })); } - pub async fn reload( - self: &Arc, - strategy_id: &str, - risk_id: &str, - ) -> tokio::io::Result<()> { + pub async fn reload_strategy(self: &Arc, id: &str) -> tokio::io::Result<()> { { if let Some(handle) = &*self.strategy_handle.lock().await { handle.abort(); } - if let Some(handle) = &*self.risk_handle.lock().await { - handle.abort(); - } + let pair = self.pair.lock().await; - { - let pair = self.pair.lock().await.clone(); + let strategy = pulse_plugin(id)?; - pair.strategy.process.lock().await.kill().await?; - pair.risk.process.lock().await.kill().await?; - } + let (child, manifest) = get_manifest_plugin_pair( + &strategy, + &strategy.join("strategy.bash"), + &fs::read(strategy.join("strategy.toml")).await?, + )?; - *self.pair.lock().await = StrategyPair::new(strategy_id, risk_id).await?; + pair.risk.reload(child).await?; + + *pair.strategy_manifest.lock().await = manifest; } self.spawn().await; From ee0ba0fe6913c5c92ce86e0175cf6abeb19888df Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 15:44:47 +0200 Subject: [PATCH 14/25] Improved reloading --- src/engine/main.rs | 13 --------- src/engine/store/plugin.rs | 55 ++++++++++++++------------------------ src/engine/terminal.rs | 4 +-- 3 files changed, 22 insertions(+), 50 deletions(-) diff --git a/src/engine/main.rs b/src/engine/main.rs index 184a811..c3c6dc6 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -20,18 +20,5 @@ async fn main() -> tokio::io::Result<()> { terminal_server??; broadcaster??; - let strategy = engine.strategy.strategy_handle.lock().await.take(); - let risk = engine.strategy.risk_handle.lock().await.take(); - - drop(engine); - - if let Some(strategy) = strategy { - strategy.into_future().await??; - } - - if let Some(risk) = risk { - risk.into_future().await??; - } - Ok(()) } diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index b891aae..112ee0d 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -19,23 +19,26 @@ use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, process::{Child, ChildStdout, Command}, sync::Mutex, - task::JoinHandle, }; use crate::{engine::Engine, store::pulse_plugin}; #[derive(Debug)] -pub struct Plugin { +pub struct Plugin Deserialize<'de>> { + pub manifest: Mutex, pub stdout: Mutex, pub process: Mutex, + pub _p: (PhantomData, PhantomData), } -impl Plugin { - pub fn new(mut child: Child) -> Self { +impl Deserialize<'de>> Plugin { + pub fn new(mut child: Child, manifest: M) -> Self { Self { stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")), process: Mutex::new(child), + manifest: Mutex::new(manifest), + _p: (PhantomData, PhantomData), } } @@ -74,13 +77,13 @@ impl Plugin { Ok(()) } - pub async fn reload(&self, mut child: Child) -> tokio::io::Result<()> { + pub async fn reload(&self, mut child: Child, manifest: M) -> tokio::io::Result<()> { let mut process = self.process.lock().await; process.kill().await?; *self.stdout.lock().await = child.stdout.take().expect("Failed to obtain child stdout"); - + *self.manifest.lock().await = manifest; *process = child; Ok(()) @@ -107,11 +110,8 @@ pub fn get_manifest_plugin_pair<'de, M: Deserialize<'de>>( #[derive(Debug)] pub struct StrategyPair { - pub strategy: Plugin, - pub risk: Plugin, - - pub strategy_manifest: Mutex, - pub risk_manifest: Mutex, + pub strategy: Plugin, + pub risk: Plugin, } impl StrategyPair { @@ -132,10 +132,8 @@ impl StrategyPair { )?; Ok(Arc::new(Self { - strategy: Plugin::new(strategy), - risk: Plugin::new(risk), - strategy_manifest: Mutex::new(strategy_manifest), - risk_manifest: Mutex::new(risk_manifest), + strategy: Plugin::new(strategy, strategy_manifest), + risk: Plugin::new(risk, risk_manifest), })) } } @@ -143,10 +141,6 @@ impl StrategyPair { #[derive(Debug)] pub struct StrategyEngine { pub pair: Mutex>, - - pub strategy_handle: Mutex>>>, - pub risk_handle: Mutex>>>, - pub engine: Weak, } @@ -154,8 +148,6 @@ impl StrategyEngine { pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { Ok(Self { pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), - strategy_handle: Mutex::new(None), - risk_handle: Mutex::new(None), engine: Weak::new(), }) } @@ -209,33 +201,26 @@ impl StrategyEngine { pub async fn spawn(self: &Arc) { let engine = self.clone(); - *self.strategy_handle.lock().await = - Some(tokio::spawn(async move { engine.run_strategy().await })); + tokio::spawn(async move { engine.run_strategy().await }); let engine = self.clone(); - *self.risk_handle.lock().await = Some(tokio::spawn(async move { engine.run_risk().await })); + tokio::spawn(async move { engine.run_risk().await }); } pub async fn reload_strategy(self: &Arc, id: &str) -> tokio::io::Result<()> { { - if let Some(handle) = &*self.strategy_handle.lock().await { - handle.abort(); - } - let pair = self.pair.lock().await; - let strategy = pulse_plugin(id)?; + let plugin = pulse_plugin(id)?; let (child, manifest) = get_manifest_plugin_pair( - &strategy, - &strategy.join("strategy.bash"), - &fs::read(strategy.join("strategy.toml")).await?, + &plugin, + &plugin.join("strategy.bash"), + &fs::read(plugin.join("strategy.toml")).await?, )?; - pair.risk.reload(child).await?; - - *pair.strategy_manifest.lock().await = manifest; + pair.strategy.reload(child, manifest).await?; } self.spawn().await; diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 24849e1..13aec15 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -66,8 +66,8 @@ impl TerminalServer { self.send_to( id, pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { - strategy: pair.strategy_manifest.lock().await.clone(), - risk: pair.risk_manifest.lock().await.clone(), + strategy: pair.strategy.manifest.lock().await.clone(), + risk: pair.risk.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, cooldown: TimeFrame::M15, From 782ea4f92dc8fd6d3bfa9aa75f81e3dd9cefc15d Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 15:47:19 +0200 Subject: [PATCH 15/25] Improved reloading --- src/engine/store/plugin.rs | 24 +++++++++++++++++++++++- 1 file changed, 23 insertions(+), 1 deletion(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 112ee0d..a681b6c 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -223,7 +223,29 @@ impl StrategyEngine { pair.strategy.reload(child, manifest).await?; } - self.spawn().await; + let engine = self.clone(); + tokio::spawn(async move { engine.run_strategy().await }); + + Ok(()) + } + + pub async fn reload_risk(self: &Arc, id: &str) -> tokio::io::Result<()> { + { + let pair = self.pair.lock().await; + + let plugin = pulse_plugin(id)?; + + let (child, manifest) = get_manifest_plugin_pair( + &plugin, + &plugin.join("risk.bash"), + &fs::read(plugin.join("risk.toml")).await?, + )?; + + pair.risk.reload(child, manifest).await?; + } + + let engine = self.clone(); + tokio::spawn(async move { engine.run_risk().await }); Ok(()) } From 033f576785d95904fab7fe1e03d5c0f6613c46ff Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 17:06:41 +0200 Subject: [PATCH 16/25] Signal forwarding --- src/engine/store/plugin.rs | 27 ++++++++++++++++++++++----- 1 file changed, 22 insertions(+), 5 deletions(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index a681b6c..3a83bd9 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -62,11 +62,11 @@ impl Deserialize<'de>> Plugin { Ok(Some(R::from_com(&mut buffer))) } - pub async fn send(&mut self, msg: &S) -> tokio::io::Result<()> { + pub async fn send(&self, msg: &S) -> tokio::io::Result<()> { self.send_raw(&msg.to_com()).await } - pub async fn send_raw(&mut self, msg: &[u8]) -> tokio::io::Result<()> { + pub async fn send_raw(&self, msg: &[u8]) -> tokio::io::Result<()> { let mut process = self.process.lock().await; let stdin = process.stdin.as_mut().unwrap(); @@ -170,11 +170,14 @@ impl StrategyEngine { match pair.strategy.recv().await? { None => {} - Some(StrategyMessage::Log(log)) => { + Some(StrategyMessage::Log(mut log)) => { + log.name.insert_str(0, "strategy::"); engine.terminal_server.log_raw(log).await?; } - Some(StrategyMessage::Signal(signal)) => {} + Some(StrategyMessage::Signal(signal)) => { + pair.risk.send(&RiskEngineMessage::Signal(signal)).await?; + } Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} @@ -193,7 +196,21 @@ impl StrategyEngine { .upgrade() .expect("Failed to upgrade engine (StrategyEngine)"); - // loop {} + let pair = self.pair.lock().await.clone(); + + loop { + match pair.risk.recv().await? { + None => {} + + Some(RiskMessage::Log(mut log)) => { + log.name.insert_str(0, "strategy::"); + engine.terminal_server.log_raw(log).await?; + } + + Some(RiskMessage::Approve(signal)) => {} + Some(RiskMessage::Reject { reason }) => {} + } + } Ok(()) } From d703bdda4ec2a78429c500497bf175bb2b01830e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 17:18:48 +0200 Subject: [PATCH 17/25] Improved engine architecture, and plugin architecture --- src/engine/{engine.rs => engine/mod.rs} | 5 +- src/engine/engine/plugin.rs | 174 +++++++++++++++++++++ src/engine/store/plugin.rs | 198 +----------------------- src/engine/terminal.rs | 5 +- 4 files changed, 183 insertions(+), 199 deletions(-) rename src/engine/{engine.rs => engine/mod.rs} (98%) create mode 100644 src/engine/engine/plugin.rs diff --git a/src/engine/engine.rs b/src/engine/engine/mod.rs similarity index 98% rename from src/engine/engine.rs rename to src/engine/engine/mod.rs index 798ebb3..c510c6c 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine/mod.rs @@ -1,5 +1,8 @@ +pub mod plugin; + use crate::{ - store::{accounts::AccountList, config::Config, plugin::StrategyEngine}, + engine::plugin::StrategyEngine, + store::{accounts::AccountList, config::Config}, terminal::TerminalServer, }; use pulse_wire::prelude::*; diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs new file mode 100644 index 0000000..86545dd --- /dev/null +++ b/src/engine/engine/plugin.rs @@ -0,0 +1,174 @@ +use pulse_wire::prelude::*; +use std::{ + path::PathBuf, + process::Stdio, + sync::{Arc, Weak}, +}; +use tokio::{ + fs, + process::{Child, Command}, +}; + +use crate::{ + engine::Engine, + store::{plugin::Plugin, pulse_plugin}, +}; + +#[derive(Debug)] +pub struct StrategyEngine { + pub strategy: Arc>, + pub risk: Arc>, + + pub engine: Weak, +} + +impl StrategyEngine { + pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { + let strategy = pulse_plugin(strategy_id)?; + let risk = pulse_plugin(risk_id)?; + + let (strategy, strategy_manifest) = get_manifest_plugin_pair( + &strategy, + &strategy.join("strategy.bash"), + &fs::read(strategy.join("strategy.toml")).await?, + )?; + + let (risk, risk_manifest) = get_manifest_plugin_pair( + &risk, + &risk.join("risk.bash"), + &fs::read(risk.join("risk.toml")).await?, + )?; + + Ok(Self { + strategy: Arc::new(Plugin::new(strategy, strategy_manifest)), + risk: Arc::new(Plugin::new(risk, risk_manifest)), + engine: Weak::new(), + }) + } + + pub fn initialize(mut self, engine: Weak) -> Arc { + self.engine = engine; + + Arc::new(self) + } + + pub async fn run_strategy(&self) -> tokio::io::Result<()> { + let engine = self + .engine + .upgrade() + .expect("Failed to upgrade engine (StrategyEngine)"); + + let strategy = self.strategy.clone(); + let risk = self.risk.clone(); + + loop { + match strategy.recv().await? { + None => {} + + Some(StrategyMessage::Log(mut log)) => { + log.name.insert_str(0, "strategy::"); + engine.terminal_server.log_raw(log).await?; + } + + Some(StrategyMessage::Signal(signal)) => { + risk.send(&RiskEngineMessage::Signal(signal)).await?; + } + + Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} + + Some(StrategyMessage::RequestOHLC { + symbol, + timeframe, + count, + }) => {} + } + } + } + + pub async fn run_risk(&self) -> tokio::io::Result<()> { + let engine = self + .engine + .upgrade() + .expect("Failed to upgrade engine (StrategyEngine)"); + + let risk = self.risk.clone(); + + loop { + match risk.recv().await? { + None => {} + + Some(RiskMessage::Log(mut log)) => { + log.name.insert_str(0, "strategy::"); + engine.terminal_server.log_raw(log).await?; + } + + Some(RiskMessage::Approve(signal)) => {} + Some(RiskMessage::Reject { reason }) => {} + } + } + + Ok(()) + } + + pub async fn spawn(self: &Arc) { + let engine = self.clone(); + + tokio::spawn(async move { engine.run_strategy().await }); + + let engine = self.clone(); + + tokio::spawn(async move { engine.run_risk().await }); + } + + pub async fn reload_strategy(self: &Arc, id: &str) -> tokio::io::Result<()> { + let plugin = pulse_plugin(id)?; + + let (child, manifest) = get_manifest_plugin_pair( + &plugin, + &plugin.join("strategy.bash"), + &fs::read(plugin.join("strategy.toml")).await?, + )?; + + self.strategy.reload(child, manifest).await?; + + let engine = self.clone(); + tokio::spawn(async move { engine.run_strategy().await }); + + Ok(()) + } + + pub async fn reload_risk(self: &Arc, id: &str) -> tokio::io::Result<()> { + let plugin = pulse_plugin(id)?; + + let (child, manifest) = get_manifest_plugin_pair( + &plugin, + &plugin.join("risk.bash"), + &fs::read(plugin.join("risk.toml")).await?, + )?; + + self.risk.reload(child, manifest).await?; + + let engine = self.clone(); + tokio::spawn(async move { engine.run_risk().await }); + + Ok(()) + } +} + +pub fn get_manifest_plugin_pair<'de, M: serde::Deserialize<'de>>( + plugin_dir: &PathBuf, + plugin_path: &PathBuf, + manifest: &'de [u8], +) -> tokio::io::Result<(Child, M)> { + Ok(( + Command::new("bash") + .arg(plugin_path) + .current_dir(plugin_dir) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::inherit()) + .spawn()?, + toml::from_slice(manifest) + .map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))?, + )) +} diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 3a83bd9..6f37cb4 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1,28 +1,14 @@ -use std::{ - marker::PhantomData, - path::PathBuf, - process::Stdio, - sync::{Arc, Weak}, -}; +use std::marker::PhantomData; -use pulse_wire::{ - PulseWire, - plugin::{ - RiskEngineMessage, RiskManifest, RiskMessage, StrategyEngineMessage, StrategyManifest, - StrategyMessage, - }, -}; +use pulse_wire::PulseWire; use serde::Deserialize; use tokio::{ - fs, io::{AsyncReadExt, AsyncWriteExt}, - process::{Child, ChildStdout, Command}, + process::{Child, ChildStdout}, sync::Mutex, }; -use crate::{engine::Engine, store::pulse_plugin}; - #[derive(Debug)] pub struct Plugin Deserialize<'de>> { pub manifest: Mutex, @@ -89,181 +75,3 @@ impl Deserialize<'de>> Plugin { Ok(()) } } - -pub fn get_manifest_plugin_pair<'de, M: Deserialize<'de>>( - plugin_dir: &PathBuf, - plugin_path: &PathBuf, - manifest: &'de [u8], -) -> tokio::io::Result<(Child, M)> { - Ok(( - Command::new("bash") - .arg(plugin_path) - .current_dir(plugin_dir) - .stdin(Stdio::piped()) - .stdout(Stdio::piped()) - .stderr(Stdio::inherit()) - .spawn()?, - toml::from_slice(manifest) - .map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))?, - )) -} - -#[derive(Debug)] -pub struct StrategyPair { - pub strategy: Plugin, - pub risk: Plugin, -} - -impl StrategyPair { - pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result> { - let strategy = pulse_plugin(strategy_id)?; - let risk = pulse_plugin(risk_id)?; - - let (strategy, strategy_manifest) = get_manifest_plugin_pair( - &strategy, - &strategy.join("strategy.bash"), - &fs::read(strategy.join("strategy.toml")).await?, - )?; - - let (risk, risk_manifest) = get_manifest_plugin_pair( - &risk, - &risk.join("risk.bash"), - &fs::read(risk.join("risk.toml")).await?, - )?; - - Ok(Arc::new(Self { - strategy: Plugin::new(strategy, strategy_manifest), - risk: Plugin::new(risk, risk_manifest), - })) - } -} - -#[derive(Debug)] -pub struct StrategyEngine { - pub pair: Mutex>, - pub engine: Weak, -} - -impl StrategyEngine { - pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { - Ok(Self { - pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), - engine: Weak::new(), - }) - } - - pub fn initialize(mut self, engine: Weak) -> Arc { - self.engine = engine; - - Arc::new(self) - } - - pub async fn run_strategy(&self) -> tokio::io::Result<()> { - let engine = self - .engine - .upgrade() - .expect("Failed to upgrade engine (StrategyEngine)"); - - let pair = self.pair.lock().await.clone(); - - loop { - match pair.strategy.recv().await? { - None => {} - - Some(StrategyMessage::Log(mut log)) => { - log.name.insert_str(0, "strategy::"); - engine.terminal_server.log_raw(log).await?; - } - - Some(StrategyMessage::Signal(signal)) => { - pair.risk.send(&RiskEngineMessage::Signal(signal)).await?; - } - - Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} - - Some(StrategyMessage::RequestOHLC { - symbol, - timeframe, - count, - }) => {} - } - } - } - - pub async fn run_risk(&self) -> tokio::io::Result<()> { - let engine = self - .engine - .upgrade() - .expect("Failed to upgrade engine (StrategyEngine)"); - - let pair = self.pair.lock().await.clone(); - - loop { - match pair.risk.recv().await? { - None => {} - - Some(RiskMessage::Log(mut log)) => { - log.name.insert_str(0, "strategy::"); - engine.terminal_server.log_raw(log).await?; - } - - Some(RiskMessage::Approve(signal)) => {} - Some(RiskMessage::Reject { reason }) => {} - } - } - - Ok(()) - } - - pub async fn spawn(self: &Arc) { - let engine = self.clone(); - - tokio::spawn(async move { engine.run_strategy().await }); - - let engine = self.clone(); - - tokio::spawn(async move { engine.run_risk().await }); - } - - pub async fn reload_strategy(self: &Arc, id: &str) -> tokio::io::Result<()> { - { - let pair = self.pair.lock().await; - - let plugin = pulse_plugin(id)?; - - let (child, manifest) = get_manifest_plugin_pair( - &plugin, - &plugin.join("strategy.bash"), - &fs::read(plugin.join("strategy.toml")).await?, - )?; - - pair.strategy.reload(child, manifest).await?; - } - - let engine = self.clone(); - tokio::spawn(async move { engine.run_strategy().await }); - - Ok(()) - } - - pub async fn reload_risk(self: &Arc, id: &str) -> tokio::io::Result<()> { - { - let pair = self.pair.lock().await; - - let plugin = pulse_plugin(id)?; - - let (child, manifest) = get_manifest_plugin_pair( - &plugin, - &plugin.join("risk.bash"), - &fs::read(plugin.join("risk.toml")).await?, - )?; - - pair.risk.reload(child, manifest).await?; - } - - let engine = self.clone(); - tokio::spawn(async move { engine.run_risk().await }); - - Ok(()) - } -} diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 13aec15..f76c9d2 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -61,13 +61,12 @@ impl TerminalServer { async fn initialize_client(self: &Arc, id: &usize) -> tokio::io::Result<()> { let engine = self.get_engine(); - let pair = engine.strategy.pair.lock().await; self.send_to( id, pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { - strategy: pair.strategy.manifest.lock().await.clone(), - risk: pair.risk.manifest.lock().await.clone(), + strategy: engine.strategy.strategy.manifest.lock().await.clone(), + risk: engine.strategy.risk.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, cooldown: TimeFrame::M15, From c8be31c88c308994a75d5afae4e9d31b1b4be144 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 19:36:07 +0200 Subject: [PATCH 18/25] risk instead of stratgy on risk log --- src/engine/engine/plugin.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 86545dd..6fdbf1d 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -98,7 +98,7 @@ impl StrategyEngine { None => {} Some(RiskMessage::Log(mut log)) => { - log.name.insert_str(0, "strategy::"); + log.name.insert_str(0, "risk::"); engine.terminal_server.log_raw(log).await?; } From 7d0075e2f90034e1d70493fdbdf6c7c8b830da3e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 00:17:43 +0200 Subject: [PATCH 19/25] Implementing hyper ws and types --- Cargo.lock | 1 + Cargo.toml | 4 +- pulse-wire/Cargo.toml | 1 + pulse-wire/src/hyper_types.rs | 225 ++++++++++++++++++++++++++++++++++ pulse-wire/src/lib.rs | 15 ++- pulse-wire/src/plugin.rs | 6 +- src/engine/engine/mod.rs | 2 +- src/engine/engine/plugin.rs | 7 +- 8 files changed, 249 insertions(+), 12 deletions(-) create mode 100644 pulse-wire/src/hyper_types.rs diff --git a/Cargo.lock b/Cargo.lock index 3cf15f6..288a837 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3906,6 +3906,7 @@ dependencies = [ name = "pulse-wire" version = "0.1.0-alpha.0" dependencies = [ + "hypersdk", "pulse-macros", "serde", ] diff --git a/Cargo.toml b/Cargo.toml index 8b87ecb..79ece7b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,7 +10,7 @@ pulse-wire = { workspace = true } chrono = "0.4.45" tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time", "process"] } crossterm = { workspace = true } -hypersdk = "0.2.14" +hypersdk = { workspace = true } serde_json = "1" rand = "0.8.7" toml = "1.1.3" @@ -23,7 +23,7 @@ members = ["pulse-macros", "pulse-ui", "pulse-wire"] 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" } - +hypersdk = "0.2.14" tokio = "1.52.3" crossterm = "0.29.0" serde = { version = "1.0.229", features = ["serde_derive"] } diff --git a/pulse-wire/Cargo.toml b/pulse-wire/Cargo.toml index 81b1d0c..e96fe28 100644 --- a/pulse-wire/Cargo.toml +++ b/pulse-wire/Cargo.toml @@ -6,3 +6,4 @@ edition = "2024" [dependencies] pulse-macros = { workspace = true } serde = { workspace = true } +hypersdk = { workspace = true } diff --git a/pulse-wire/src/hyper_types.rs b/pulse-wire/src/hyper_types.rs new file mode 100644 index 0000000..dcc44c0 --- /dev/null +++ b/pulse-wire/src/hyper_types.rs @@ -0,0 +1,225 @@ +use hypersdk::{Address, hypercore::Subscription}; + +use crate::PulseWire; + +impl PulseWire for Subscription { + fn from_com(com: &mut Vec) -> Self { + match u8::from_com(com) { + 0 => Self::Bbo { + coin: PulseWire::from_com(com), + }, + 1 => Self::Trades { + coin: PulseWire::from_com(com), + }, + 2 => Self::L2Book { + coin: PulseWire::from_com(com), + n_sig_figs: PulseWire::from_com(com), + mantissa: PulseWire::from_com(com), + fast: PulseWire::from_com(com), + }, + 3 => Self::Candle { + coin: PulseWire::from_com(com), + interval: PulseWire::from_com(com), + }, + 4 => Self::AllMids { + dex: PulseWire::from_com(com), + }, + 5 => Self::OrderUpdates { + user: PulseWire::from_com(com), + }, + 6 => Self::UserFills { + user: PulseWire::from_com(com), + }, + 7 => Self::UserEvents { + user: PulseWire::from_com(com), + }, + 8 => Self::UserTwapSliceFills { + user: PulseWire::from_com(com), + }, + 9 => Self::UserTwapHistory { + user: PulseWire::from_com(com), + }, + 10 => Self::ActiveAssetCtx { + coin: PulseWire::from_com(com), + }, + 11 => Self::ActiveAssetData { + user: PulseWire::from_com(com), + coin: PulseWire::from_com(com), + }, + 12 => Self::WebData2 { + user: PulseWire::from_com(com), + dex: PulseWire::from_com(com), + }, + 13 => Self::ClearinghouseState { + user: PulseWire::from_com(com), + dex: PulseWire::from_com(com), + }, + 14 => Self::AllDexsClearinghouseState { + user: PulseWire::from_com(com), + }, + 15 => Self::OpenOrders { + user: PulseWire::from_com(com), + dex: PulseWire::from_com(com), + }, + 16 => Self::SpotState { + user: PulseWire::from_com(com), + is_portfolio_margin: PulseWire::from_com(com), + }, + 17 => Self::Notification { + user: PulseWire::from_com(com), + }, + 18 => Self::WebData3 { + user: PulseWire::from_com(com), + }, + 19 => Self::TwapStates { + user: PulseWire::from_com(com), + dex: Option::from_com(com), + }, + 20 => Self::UserFundings { + user: PulseWire::from_com(com), + }, + 21 => Self::UserNonFundingLedgerUpdates { + user: PulseWire::from_com(com), + }, + 22 => Self::AllDexsAssetCtxs, + 23 => Self::FastAssetCtxs, + 24 => Self::OutcomeMetaUpdates, + x => panic!("Invalid Subscription discriminant: {}", x), + } + } + + fn to_com(&self) -> Vec { + let mut com = Vec::new(); + + match self { + Self::Bbo { coin } => { + com.push(0); + com.extend(coin.to_com()); + } + Self::Trades { coin } => { + com.push(1); + com.extend(coin.to_com()); + } + Self::L2Book { + coin, + n_sig_figs, + mantissa, + fast, + } => { + com.push(2); + com.extend(coin.to_com()); + com.extend(n_sig_figs.to_com()); + com.extend(mantissa.to_com()); + com.extend(fast.to_com()); + } + Self::Candle { coin, interval } => { + com.push(3); + com.extend(coin.to_com()); + com.extend(interval.to_com()); + } + Self::AllMids { dex } => { + com.push(4); + com.extend(dex.to_com()); + } + Self::OrderUpdates { user } => { + com.push(5); + com.extend(user.to_com()); + } + Self::UserFills { user } => { + com.push(6); + com.extend(user.to_com()); + } + Self::UserEvents { user } => { + com.push(7); + com.extend(user.to_com()); + } + Self::UserTwapSliceFills { user } => { + com.push(8); + com.extend(user.to_com()); + } + Self::UserTwapHistory { user } => { + com.push(9); + com.extend(user.to_com()); + } + Self::ActiveAssetCtx { coin } => { + com.push(10); + com.extend(coin.to_com()); + } + Self::ActiveAssetData { user, coin } => { + com.push(11); + com.extend(user.to_com()); + com.extend(coin.to_com()); + } + Self::WebData2 { user, dex } => { + com.push(12); + com.extend(user.to_com()); + com.extend(dex.to_com()); + } + Self::ClearinghouseState { user, dex } => { + com.push(13); + com.extend(user.to_com()); + com.extend(dex.to_com()); + } + Self::AllDexsClearinghouseState { user } => { + com.push(14); + com.extend(user.to_com()); + } + Self::OpenOrders { user, dex } => { + com.push(15); + com.extend(user.to_com()); + com.extend(dex.to_com()); + } + Self::SpotState { + user, + is_portfolio_margin, + } => { + com.push(16); + com.extend(user.to_com()); + com.extend(is_portfolio_margin.to_com()); + } + Self::Notification { user } => { + com.push(17); + com.extend(user.to_com()); + } + Self::WebData3 { user } => { + com.push(18); + com.extend(user.to_com()); + } + Self::TwapStates { user, dex } => { + com.push(19); + com.extend(user.to_com()); + com.extend(dex.to_com()); + } + Self::UserFundings { user } => { + com.push(20); + com.extend(user.to_com()); + } + Self::UserNonFundingLedgerUpdates { user } => { + com.push(21); + com.extend(user.to_com()); + } + Self::AllDexsAssetCtxs => { + com.push(22); + } + Self::FastAssetCtxs => { + com.push(23); + } + Self::OutcomeMetaUpdates => { + com.push(24); + } + } + + com + } +} + +impl PulseWire for Address { + fn from_com(com: &mut Vec) -> Self { + let bytes: [u8; 20] = com.drain(..20).collect::>().try_into().unwrap(); + Self::from_slice(&bytes) + } + + fn to_com(&self) -> Vec { + self.as_slice().to_vec() + } +} diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index 0d8b6d8..be1a05a 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -2,6 +2,7 @@ use std::path::PathBuf; pub mod general; +mod hyper_types; pub mod plugin; pub mod terminal; pub mod units; @@ -9,8 +10,8 @@ pub mod units; pub mod prelude { pub use crate::PulseWire; pub use crate::general::*; - pub use crate::server_path; pub use crate::plugin::*; + pub use crate::server_path; pub use crate::terminal::*; pub use crate::units::*; } @@ -55,7 +56,7 @@ impl PulseWire for Vec { impl PulseWire for Option { fn to_com(&self) -> Vec { if let Some(v) = self { - v.to_com() + vec![1].into_iter().chain(v.to_com()).collect() } else { vec![0] } @@ -90,6 +91,16 @@ impl PulseWire for String { } } +impl PulseWire for bool { + fn to_com(&self) -> Vec { + if *self { vec![1] } else { vec![0] } + } + + fn from_com(com: &mut Vec) -> Self { + com[0] > 0 + } +} + macro_rules! int_com { ($t:ty) => { impl $crate::PulseWire for $t { diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs index 07d2419..d4b49a8 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -4,6 +4,7 @@ use crate::{ terminal::MarketItem, units::{Direction, TimeFrame}, }; +use hypersdk::hypercore::Subscription; use pulse_macros::pwp; #[pwp] @@ -37,10 +38,7 @@ pub enum StrategyMessage { count: u32, }, - SubscribeCandle { - symbol: String, - timeframe: TimeFrame, - }, + Subscribe(Subscription), Signal(StrategySignal), diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index c510c6c..aa7e18e 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -9,7 +9,7 @@ use pulse_wire::prelude::*; use std::sync::Arc; use tokio::{sync::Mutex, task::JoinHandle}; -#[derive(Debug, Clone)] +#[derive(Clone)] pub struct Engine { pub terminal_server: Arc, pub strategy: Arc, diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 6fdbf1d..86f163f 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -1,3 +1,4 @@ +use hypersdk::hypercore::{self, WebSocket}; use pulse_wire::prelude::*; use std::{ path::PathBuf, @@ -14,12 +15,11 @@ use crate::{ store::{plugin::Plugin, pulse_plugin}, }; -#[derive(Debug)] pub struct StrategyEngine { pub strategy: Arc>, pub risk: Arc>, - pub engine: Weak, + pub ws: WebSocket, } impl StrategyEngine { @@ -43,6 +43,7 @@ impl StrategyEngine { strategy: Arc::new(Plugin::new(strategy, strategy_manifest)), risk: Arc::new(Plugin::new(risk, risk_manifest)), engine: Weak::new(), + ws: hypercore::mainnet_ws(), }) } @@ -74,7 +75,7 @@ impl StrategyEngine { risk.send(&RiskEngineMessage::Signal(signal)).await?; } - Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} + Some(StrategyMessage::Subscribe(val)) => {} Some(StrategyMessage::RequestOHLC { symbol, From 1a47d0c487c35046cd89f76f14498caf939aa7da Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 01:03:12 +0200 Subject: [PATCH 20/25] Subscriptions --- pulse-wire/src/hyper_types.rs | 90 ++++++++++++++++++++++++++++++++++- pulse-wire/src/plugin.rs | 18 +++---- pulse-wire/src/terminal.rs | 5 +- pulse-wire/src/units.rs | 60 ----------------------- src/engine/engine/plugin.rs | 58 ++++++++++++++++++++-- src/engine/terminal.rs | 2 +- src/terminal/formatting.rs | 4 +- 7 files changed, 159 insertions(+), 78 deletions(-) diff --git a/pulse-wire/src/hyper_types.rs b/pulse-wire/src/hyper_types.rs index dcc44c0..9e19906 100644 --- a/pulse-wire/src/hyper_types.rs +++ b/pulse-wire/src/hyper_types.rs @@ -1,4 +1,7 @@ -use hypersdk::{Address, hypercore::Subscription}; +use hypersdk::{ + Address, Decimal, + hypercore::{Candle, CandleInterval, Subscription}, +}; use crate::PulseWire; @@ -223,3 +226,88 @@ impl PulseWire for Address { self.as_slice().to_vec() } } + +impl PulseWire for CandleInterval { + fn from_com(com: &mut Vec) -> Self { + match u8::from_com(com) { + 0 => Self::OneMinute, + 1 => Self::ThreeMinutes, + 2 => Self::FiveMinutes, + 3 => Self::FifteenMinutes, + 4 => Self::ThirtyMinutes, + 5 => Self::OneHour, + 6 => Self::TwoHours, + 7 => Self::FourHours, + 8 => Self::EightHours, + 9 => Self::TwelveHours, + 10 => Self::OneDay, + 11 => Self::ThreeDays, + 12 => Self::OneWeek, + 13 => Self::OneMonth, + x => panic!("Invalid CandleInterval discriminant: {}", x), + } + } + + fn to_com(&self) -> Vec { + vec![match self { + Self::OneMinute => 0, + Self::ThreeMinutes => 1, + Self::FiveMinutes => 2, + Self::FifteenMinutes => 3, + Self::ThirtyMinutes => 4, + Self::OneHour => 5, + Self::TwoHours => 6, + Self::FourHours => 7, + Self::EightHours => 8, + Self::TwelveHours => 9, + Self::OneDay => 10, + Self::ThreeDays => 11, + Self::OneWeek => 12, + Self::OneMonth => 13, + }] + } +} + +impl PulseWire for Candle { + fn from_com(com: &mut Vec) -> Self { + Self { + open_time: u64::from_com(com), + close_time: u64::from_com(com), + coin: String::from_com(com), + interval: String::from_com(com), + open: Decimal::from_com(com), + high: Decimal::from_com(com), + low: Decimal::from_com(com), + close: Decimal::from_com(com), + volume: Decimal::from_com(com), + num_trades: u64::from_com(com), + } + } + + fn to_com(&self) -> Vec { + let mut com = Vec::new(); + + com.extend(self.open_time.to_com()); + com.extend(self.close_time.to_com()); + com.extend(self.coin.to_com()); + com.extend(self.interval.to_com()); + com.extend(self.open.to_com()); + com.extend(self.high.to_com()); + com.extend(self.low.to_com()); + com.extend(self.close.to_com()); + com.extend(self.volume.to_com()); + com.extend(self.num_trades.to_com()); + + com + } +} + +impl PulseWire for Decimal { + fn from_com(com: &mut Vec) -> Self { + Self::deserialize(com[..16].try_into().unwrap()) + } + + fn to_com(&self) -> Vec { + self.serialize().to_vec() + } +} diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs index d4b49a8..db9139b 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -2,9 +2,9 @@ use crate::{ PulseWire, general::{EventLog, Signal}, terminal::MarketItem, - units::{Direction, TimeFrame}, + units::Direction, }; -use hypersdk::hypercore::Subscription; +use hypersdk::hypercore::{Candle, CandleInterval, Subscription}; use pulse_macros::pwp; #[pwp] @@ -14,8 +14,6 @@ pub struct StrategyManifest { description: String, author: String, version: String, - - timeframes: Vec, } #[pwp] @@ -27,19 +25,23 @@ pub struct RiskManifest { version: String, max_loss: u8, - cooldown: TimeFrame, + cooldown: String, } #[pwp] pub enum StrategyMessage { RequestOHLC { symbol: String, - timeframe: TimeFrame, + interval: CandleInterval, count: u32, }, Subscribe(Subscription), + Unsubscribe(Subscription), + + UnsubscribeAll, + Signal(StrategySignal), Log(EventLog), @@ -58,8 +60,8 @@ pub enum StrategyEngineMessage { OHLC { symbol: String, - timeframe: TimeFrame, - candles: Vec, + interval: CandleInterval, + candles: Vec, }, Start, diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index d59b3fc..006c646 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -2,8 +2,9 @@ use crate::{ PulseWire, general::{EventLog, MarketTrend, Position, Signal}, plugin::{RiskManifest, StrategyManifest}, - units::{Symbol, TimeFrame, USD, Volatility}, + units::{Symbol, USD, Volatility}, }; +use hypersdk::hypercore::CandleInterval; use pulse_macros::pwp; #[pwp] @@ -113,7 +114,7 @@ pub struct Strategy { mode: Mode, state: ItemState, - cooldown: TimeFrame, + cooldown: CandleInterval, } impl std::fmt::Display for AlertLevel { diff --git a/pulse-wire/src/units.rs b/pulse-wire/src/units.rs index 06e84ae..c7db0e6 100644 --- a/pulse-wire/src/units.rs +++ b/pulse-wire/src/units.rs @@ -8,42 +8,6 @@ pub struct Symbol(pub String); #[derive(Debug, Clone, Copy)] pub struct USD(pub f64); -#[pwp] -#[derive(Copy, serde::Deserialize, serde::Serialize)] -pub enum TimeFrame { - #[serde(rename = "1m")] - M1, - #[serde(rename = "3m")] - M3, - #[serde(rename = "5m")] - M5, - #[serde(rename = "15m")] - M15, - #[serde(rename = "30m")] - M30, - - #[serde(rename = "1h")] - H1, - #[serde(rename = "2h")] - H2, - #[serde(rename = "4h")] - H4, - #[serde(rename = "8h")] - H8, - #[serde(rename = "12h")] - H12, - - #[serde(rename = "1d")] - D1, - #[serde(rename = "3d")] - D3, - - #[serde(rename = "1w")] - W1, - #[serde(rename = "1M")] - Month1, -} - #[pwp] pub enum Direction { Buy, @@ -148,30 +112,6 @@ pub fn format_f64(value: f64) -> String { formatted } -impl TimeFrame { - pub fn as_str(&self) -> &'static str { - match self { - Self::M1 => "1m", - Self::M3 => "3m", - Self::M5 => "5m", - Self::M15 => "15m", - Self::M30 => "30m", - - Self::H1 => "1h", - Self::H2 => "2h", - Self::H4 => "4h", - Self::H8 => "8h", - Self::H12 => "12h", - - Self::D1 => "1d", - Self::D3 => "3d", - - Self::W1 => "1w", - Self::Month1 => "1M", - } - } -} - impl std::fmt::Display for Direction { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 86f163f..de38592 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -1,13 +1,16 @@ -use hypersdk::hypercore::{self, WebSocket}; +use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket}; use pulse_wire::prelude::*; use std::{ + collections::HashSet, path::PathBuf, process::Stdio, sync::{Arc, Weak}, + time::{SystemTime, UNIX_EPOCH}, }; use tokio::{ fs, process::{Child, Command}, + sync::Mutex, }; use crate::{ @@ -19,7 +22,9 @@ pub struct StrategyEngine { pub strategy: Arc>, pub risk: Arc>, pub engine: Weak, + pub ws: WebSocket, + pub subscriptions: Mutex>, } impl StrategyEngine { @@ -44,6 +49,7 @@ impl StrategyEngine { risk: Arc::new(Plugin::new(risk, risk_manifest)), engine: Weak::new(), ws: hypercore::mainnet_ws(), + subscriptions: Mutex::new(HashSet::new()), }) } @@ -75,13 +81,57 @@ impl StrategyEngine { risk.send(&RiskEngineMessage::Signal(signal)).await?; } - Some(StrategyMessage::Subscribe(val)) => {} + Some(StrategyMessage::Subscribe(subscription)) => { + self.ws.subscribe(subscription.clone()); + self.subscriptions.lock().await.insert(subscription); + } + + Some(StrategyMessage::Unsubscribe(subscription)) => { + self.subscriptions.lock().await.remove(&subscription); + self.ws.unsubscribe(subscription); + } + + Some(StrategyMessage::UnsubscribeAll) => { + for sub in self.subscriptions.lock().await.drain() { + self.ws.unsubscribe(sub); + } + } Some(StrategyMessage::RequestOHLC { symbol, - timeframe, + interval, count, - }) => {} + }) => { + let client = hypercore::mainnet(); + + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis() as u64; + + let interval_ms = match interval { + CandleInterval::OneMinute => 60_000, + CandleInterval::ThreeMinutes => 3 * 60_000, + CandleInterval::FiveMinutes => 5 * 60_000, + CandleInterval::FifteenMinutes => 15 * 60_000, + CandleInterval::ThirtyMinutes => 30 * 60_000, + CandleInterval::OneHour => 60 * 60_000, + CandleInterval::TwoHours => 2 * 60 * 60_000, + CandleInterval::FourHours => 4 * 60 * 60_000, + CandleInterval::EightHours => 8 * 60 * 60_000, + CandleInterval::TwelveHours => 12 * 60 * 60_000, + CandleInterval::OneDay => 24 * 60 * 60_000, + CandleInterval::ThreeDays => 3 * 24 * 60 * 60_000, + CandleInterval::OneWeek => 7 * 24 * 60 * 60_000, + CandleInterval::OneMonth => 30 * 24 * 60 * 60_000, + }; + + let start_time = now.saturating_sub(interval_ms * count as u64); + + client + .candle_snapshot(symbol, interval, start_time, now) + .await?; + } } } } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index f76c9d2..39eb997 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -69,7 +69,7 @@ impl TerminalServer { risk: engine.strategy.risk.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, - cooldown: TimeFrame::M15, + cooldown: hypersdk::hypercore::CandleInterval::ThirtyMinutes, }), ) .await?; diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index ed76ef9..b46ef22 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -254,7 +254,7 @@ impl Formatted for Strategy { Triple("", "", ""), Triple( "\x1b[2mCooldown\x1b[0m", - "\x1b[2mTimeframes\x1b[0m", + "\x1b[2mStrat Author\x1b[0m", "\x1b[2mMax loss\x1b[0m", ), Triple( @@ -262,7 +262,7 @@ impl Formatted for Strategy { "\x1b[96m{:?}\x1b[0m (\x1b[90m{:?} rec\x1b[0m)", self.cooldown, self.risk.cooldown ), - &format!("\x1b[96m{:?}\x1b[0m", self.strategy.timeframes), + &format!("\x1b[96m{:?}\x1b[0m", self.strategy.author), &format!("\x1b[93m{}%\x1b[0m", self.risk.max_loss), ), ] From b412cf50ec2ef9e08355e94a98faaf34b1ecde3f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 01:22:41 +0200 Subject: [PATCH 21/25] Fix cooldown display --- Cargo.lock | 1 + Cargo.toml | 17 +++++++++++++---- pulse-wire/src/lib.rs | 2 ++ src/engine/engine/plugin.rs | 4 +--- src/engine/main.rs | 2 +- src/terminal/formatting.rs | 2 +- 6 files changed, 19 insertions(+), 9 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 288a837..48d222b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3882,6 +3882,7 @@ dependencies = [ name = "pulse-trader" version = "0.1.0-alpha.0" dependencies = [ + "anyhow", "chrono", "crossterm", "hypersdk", diff --git a/Cargo.toml b/Cargo.toml index 79ece7b..5371dc0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,15 +6,23 @@ edition = "2024" [dependencies] 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", "process"] } crossterm = { workspace = true } hypersdk = { workspace = true } +serde = { workspace = true } +anyhow = { workspace = true } +tokio = { workspace = true, features = [ + "rt-multi-thread", + "macros", + "net", + "fs", + "io-util", + "time", + "process", +] } +chrono = "0.4.45" serde_json = "1" rand = "0.8.7" toml = "1.1.3" -serde = { workspace = true } [workspace] members = ["pulse-macros", "pulse-ui", "pulse-wire"] @@ -27,6 +35,7 @@ hypersdk = "0.2.14" tokio = "1.52.3" crossterm = "0.29.0" serde = { version = "1.0.229", features = ["serde_derive"] } +anyhow = "1.0.104" [[bin]] name = "pulse-trader" diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index be1a05a..0ff2c32 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -6,6 +6,7 @@ mod hyper_types; pub mod plugin; pub mod terminal; pub mod units; +pub use hypersdk; pub mod prelude { pub use crate::PulseWire; @@ -14,6 +15,7 @@ pub mod prelude { pub use crate::server_path; pub use crate::terminal::*; pub use crate::units::*; + pub use hypersdk; } pub fn server_path() -> PathBuf { diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index de38592..51ebbcd 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -59,7 +59,7 @@ impl StrategyEngine { Arc::new(self) } - pub async fn run_strategy(&self) -> tokio::io::Result<()> { + pub async fn run_strategy(&self) -> anyhow::Result<()> { let engine = self .engine .upgrade() @@ -157,8 +157,6 @@ impl StrategyEngine { Some(RiskMessage::Reject { reason }) => {} } } - - Ok(()) } pub async fn spawn(self: &Arc) { diff --git a/src/engine/main.rs b/src/engine/main.rs index c3c6dc6..7555ce1 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -4,7 +4,7 @@ pub mod store; pub mod terminal; #[tokio::main] -async fn main() -> tokio::io::Result<()> { +async fn main() -> anyhow::Result<()> { let engine = engine::Engine::new().await?; let terminal_server = engine.spawn_terminal_server().await; diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index b46ef22..504ff42 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -259,7 +259,7 @@ impl Formatted for Strategy { ), Triple( &format!( - "\x1b[96m{:?}\x1b[0m (\x1b[90m{:?} rec\x1b[0m)", + "\x1b[96m{}\x1b[0m (\x1b[90m{} rec\x1b[0m)", self.cooldown, self.risk.cooldown ), &format!("\x1b[96m{:?}\x1b[0m", self.strategy.author), From 99515c17d4157a31057b8157e942b629795dbb32 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 01:52:52 +0200 Subject: [PATCH 22/25] Cooldown config --- pulse-wire/src/plugin.rs | 2 +- pulse-wire/src/terminal.rs | 4 ++-- src/engine/engine/mod.rs | 2 +- src/engine/store/config.rs | 21 +++++++++++++++------ src/engine/terminal.rs | 2 +- 5 files changed, 20 insertions(+), 11 deletions(-) diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs index db9139b..bd91135 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -25,7 +25,7 @@ pub struct RiskManifest { version: String, max_loss: u8, - cooldown: String, + cooldown: CandleInterval, } #[pwp] diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index 006c646..072001d 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -103,7 +103,7 @@ pub enum Mode { #[pwp] pub enum ItemState { Running, - Off, + Stopped, Error, } @@ -140,7 +140,7 @@ impl std::fmt::Display for ItemState { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { Self::Running => write!(f, "\x1b[92mRUNNING\x1b[0m"), - Self::Off => write!(f, "\x1b[90mOFF\x1b[0m"), + Self::Stopped => write!(f, "\x1b[90mSTOPPED\x1b[0m"), Self::Error => write!(f, "\x1b[91mERROR\x1b[0m"), } } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index aa7e18e..dcd8e65 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -54,7 +54,7 @@ impl Engine { loop { refresh.tick().await; - let watch_list = &self.config.lock().await.watchlist.symbols; + let watch_list = &self.config.lock().await.watchlist; match crate::fetch::fetch_watch_list(&client, watch_list).await { Ok(watch_list) => { diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index afaf487..3d3f507 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -1,13 +1,22 @@ -#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct WatchList { - pub symbols: Vec, -} +use hypersdk::hypercore::CandleInterval; -#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] +#[derive(Debug, serde::Serialize, serde::Deserialize)] pub struct Config { - pub watchlist: WatchList, + pub watchlist: Vec, pub strategy: String, pub risk: String, + pub cooldown: CandleInterval, +} + +impl Default for Config { + fn default() -> Self { + Self { + watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()], + strategy: String::new(), + risk: String::new(), + cooldown: CandleInterval::ThirtyMinutes, + } + } } impl Config { diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 39eb997..6623b6a 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -69,7 +69,7 @@ impl TerminalServer { risk: engine.strategy.risk.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, - cooldown: hypersdk::hypercore::CandleInterval::ThirtyMinutes, + cooldown: engine.config.lock().await.cooldown, }), ) .await?; From 9006bc7440a5df0ee863a814e4d8aeba551315d0 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 02:23:59 +0200 Subject: [PATCH 23/25] More commands --- src/engine/engine/command.rs | 163 +++++++++++++++++++++++++++++++++++ src/engine/engine/mod.rs | 101 +--------------------- 2 files changed, 164 insertions(+), 100 deletions(-) create mode 100644 src/engine/engine/command.rs diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs new file mode 100644 index 0000000..67affcf --- /dev/null +++ b/src/engine/engine/command.rs @@ -0,0 +1,163 @@ +use crate::{engine::Engine, store::config::Config}; + +use pulse_wire::terminal::{ItemState, Mode, Strategy}; +use toml::Value; + +impl Engine { + pub async fn execute_command(&self, command: &str, args: Vec<&str>) -> tokio::io::Result<()> { + match command { + "config" | "cfg" => { + if args.len() == 0 { + return self.invalid_command_usage("config").await; + } + + match args[0] { + "reload" => { + *self.config.lock().await = Config::new().await?; + } + + "set" => { + macro_rules! set_cfg { + ($n:ident, $v:expr) => { + if let Ok(Ok($n)) = Value::try_from(args[2]).map(|v| v.try_into()) { + $v; + } else { + self.terminal_server + .error("config::set", "Unable to parse value") + .await?; + } + }; + } + + if args.len() == 3 { + match args[1] { + "watchlist" | "watch" | "wl" => { + set_cfg!(v, self.config.lock().await.watchlist = v); + } + + "strategy" | "strat" | "str" | "sg" => { + set_cfg!(id, { + let id: String = id; + self.strategy.reload_strategy(id.as_str()).await?; + self.config.lock().await.strategy = id; + }); + } + + "risk" | "rs" => { + set_cfg!(id, { + let id: String = id; + self.strategy.reload_risk(id.as_str()).await?; + self.config.lock().await.risk = id; + }); + } + + "cooldown" | "cool" | "cd" => { + set_cfg!(cooldown, { + self.config.lock().await.cooldown = cooldown; + + self.terminal_server.broadcast( + pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { + strategy: self.strategy.strategy.manifest.lock().await.clone(), + risk: self.strategy.risk.manifest.lock().await.clone(), + mode: Mode::Auto, + state: ItemState::Running, + cooldown, + }), + ) + .await?; + }); + } + + _ => { + self.terminal_server + .error("config::set", "Invalid usage, available options: watchlist, strategy, risk, cooldown") + .await?; + } + } + } else { + self.invalid_command_usage("config").await?; + } + } + + _ => { + self.invalid_command_usage("config").await?; + } + } + } + + "account" | "acc" => { + if args.len() == 0 { + return self.invalid_command_usage("account man").await; + } + + match args[0] { + "list" | "ls" => { + self.terminal_server + .info("account man", "ACCOUNT LIST") + .await?; + + let accounts = self.accounts.lock().await; + + for (name, acc) in &accounts.accounts { + self.terminal_server + .info( + "account man", + &if name == &accounts.active { + format!( + "{} (active) -> {}", + name, + acc.get_truncated_address() + ) + } else { + format!("{} -> {}", name, acc.get_truncated_address()) + }, + ) + .await?; + } + } + + "use" | "set" => { + if args.len() < 2 { + return self.invalid_command_usage("account man").await; + } + + let new_active = args[1]; + + let mut accounts = self.accounts.lock().await; + + if !accounts.accounts.contains_key(new_active) { + return self + .terminal_server + .error("account man", &format!("Account not found ({new_active})")) + .await; + } + + accounts.active = new_active.to_string(); + + self.terminal_server + .info( + "account man", + &format!("Account set to {new_active} successfully!"), + ) + .await?; + } + + _ => { + return self.invalid_command_usage("account man").await; + } + } + } + + _ => { + self.terminal_server + .error( + "Command executor", + &format!("Command '{}' not found", command), + ) + .await?; + } + } + + Ok(()) + } +} diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index dcd8e65..6afa49f 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -1,4 +1,5 @@ pub mod plugin; +pub mod command; use crate::{ engine::plugin::StrategyEngine, @@ -136,104 +137,4 @@ impl Engine { pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> { self.terminal_server.error(name, "Invalid usage").await } - - pub async fn execute_command(&self, command: &str, args: Vec<&str>) -> tokio::io::Result<()> { - match command { - "config" | "cfg" => { - if args.len() != 1 { - return self.invalid_command_usage("config").await; - } - - match args[0] { - "reload" => { - *self.config.lock().await = Config::new().await?; - } - - _ => { - self.terminal_server - .error("config", "Invalid usage") - .await?; - } - } - } - - "account" | "acc" => { - if args.len() == 0 { - return self.invalid_command_usage("account man").await; - } - - match args[0] { - "reload" => { - *self.accounts.lock().await = AccountList::new().await?; - } - - "list" | "ls" => { - self.terminal_server - .info("account man", "ACCOUNT LIST") - .await?; - - let accounts = self.accounts.lock().await; - - for (name, acc) in &accounts.accounts { - self.terminal_server - .info( - "account man", - &if name == &accounts.active { - format!( - "{} (active) -> {}", - name, - acc.get_truncated_address() - ) - } else { - format!("{} -> {}", name, acc.get_truncated_address()) - }, - ) - .await?; - } - } - - "use" | "set" => { - if args.len() < 2 { - return self.invalid_command_usage("account man").await; - } - - let new_active = args[1]; - - let mut accounts = self.accounts.lock().await; - - if !accounts.accounts.contains_key(new_active) { - return self - .terminal_server - .error("account man", &format!("Account not found ({new_active})")) - .await; - } - - accounts.active = new_active.to_string(); - - self.terminal_server - .info( - "account man", - &format!("Account set to {new_active} successfully!"), - ) - .await?; - } - - _ => { - return self.invalid_command_usage("account man").await; - } - } - } - - _ => { - self.terminal_server - .error( - "Command executor", - &format!("Command '{}' not found", command), - ) - .await?; - } - } - - Ok(()) - } } From 61dc17987d98965a1f70209df1c0ea1b19656138 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 02:38:19 +0200 Subject: [PATCH 24/25] Smarter risk and strategy, and logging --- src/engine/engine/command.rs | 44 ++++++++++++++++++++++++++++++++++++ src/engine/store/config.rs | 8 +++++++ 2 files changed, 52 insertions(+) diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 67affcf..836f104 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -16,6 +16,10 @@ impl Engine { *self.config.lock().await = Config::new().await?; } + "save" => { + self.config.lock().await.save().await?; + } + "set" => { macro_rules! set_cfg { ($n:ident, $v:expr) => { @@ -33,22 +37,62 @@ impl Engine { match args[1] { "watchlist" | "watch" | "wl" => { set_cfg!(v, self.config.lock().await.watchlist = v); + + self.terminal_server + .info("config::set", "watchlist set successfully, use `config save` to persist changes") + .await?; } "strategy" | "strat" | "str" | "sg" => { set_cfg!(id, { let id: String = id; + + if !crate::store::pulse_plugin(&id)? + .join("strategy.toml") + .exists() + { + return self + .terminal_server + .error( + "config::set::strategy", + &format!("Non existent strategy `{id}`"), + ) + .await; + } + self.strategy.reload_strategy(id.as_str()).await?; self.config.lock().await.strategy = id; }); + + self.terminal_server + .info("config::set", "strategy set successfully, use `config save` to persist changes") + .await?; } "risk" | "rs" => { set_cfg!(id, { let id: String = id; + + if !crate::store::pulse_plugin(&id)? + .join("risk.toml") + .exists() + { + return self + .terminal_server + .error( + "config::set::risk", + &format!("Non existent risk `{id}`"), + ) + .await; + } + self.strategy.reload_risk(id.as_str()).await?; self.config.lock().await.risk = id; }); + + self.terminal_server + .info("config::set", "risk set successfully, use `config save` to persist changes") + .await?; } "cooldown" | "cool" | "cd" => { diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index 3d3f507..c16d916 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -37,6 +37,14 @@ impl Config { Self::from_str(&output) } + pub async fn save(&self) -> tokio::io::Result<()> { + let path = crate::store::pulse_config_file()?; + + tokio::fs::write(path, self.to_string()?).await?; + + Ok(()) + } + 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())) From dd48045dd621ee906aeff844fc146614b38fb3bd Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 02:41:46 +0200 Subject: [PATCH 25/25] renamed to configuration --- src/terminal/formatting.rs | 2 +- src/terminal/main.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 504ff42..b807ead 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -231,7 +231,7 @@ impl Formatted for Strategy { fn get_formatted(&self) -> Vec { vec![ Triple( - "\x1b[2mName\x1b[0m", + "\x1b[2mStrategy\x1b[0m", "\x1b[2mRisk\x1b[0m", "\x1b[2mStrat Ver\x1b[0m", ), diff --git a/src/terminal/main.rs b/src/terminal/main.rs index e1ccf47..979c8a6 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -144,7 +144,7 @@ impl App for PulseTradeApp { ( LayoutItem::Widget(Size::Flex(1)), Box::new( - advanced_option_draw(&self.scroll, 3, "STRATEGY", &self.strategy).await, + advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.strategy).await, ), ), (