From f9c67ce923a9b8c867349af1f18f4b4bb9cf5bd8 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 31 Jul 2026 17:02:38 +0200 Subject: [PATCH 01/26] Config presets --- pulse-sdk/src/terminal.rs | 1 + src/engine/engine/command.rs | 118 +++++++++++++++++++++++++--------- src/engine/engine/mod.rs | 9 +-- src/engine/engine/terminal.rs | 7 +- src/engine/store/config.rs | 45 +++++++++++-- src/terminal/formatting.rs | 4 +- 6 files changed, 141 insertions(+), 43 deletions(-) diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index c82562b..168cc73 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -63,5 +63,6 @@ pub enum InspectItem { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct EngineConfig { + pub preset: String, pub strategy: StrategyManifest, } diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 0a815e1..06724f6 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -1,4 +1,4 @@ -use crate::{engine::Engine, store::config::Config}; +use crate::{engine::Engine, store::config::ConfigManager}; use pulse_sdk::prelude::*; use toml::Value; @@ -21,7 +21,7 @@ impl Engine { match args[0] { "reload" => { - *self.config.lock().await = Config::new().await?; + *self.config.lock().await = ConfigManager::new().await?; self.terminal_server .info("engine::config", "successfully reloaded") .await?; @@ -50,17 +50,11 @@ impl Engine { if args.len() == 3 { match args[1] { "watchlist" | "watch" | "wl" => { - set_cfg!(v, self.config.lock().await.watchlist = v); + set_cfg!(v, self.config.lock().await.get_mut()?.watchlist = v); self.terminal_server .info("config::set", "watchlist set successfully, use `config save` to persist changes") .await?; - - self.terminal_server - .broadcast(TerminalServerMessage::ConfigUpdated( - self.get_config_status().await, - )) - .await?; } "strategy" | "strat" | "sg" => { @@ -81,30 +75,12 @@ impl Engine { } self.strategy_engine.reload(id.as_str()).await?; - self.config.lock().await.strategy = id; + self.config.lock().await.get_mut()?.strategy = id; }); self.terminal_server .info("config::set", "strategy set successfully, use `config save` to persist changes") .await?; - - self.terminal_server - .broadcast(TerminalServerMessage::ConfigUpdated( - self.get_config_status().await, - )) - .await?; - } - - "cooldown" | "cool" | "cd" => { - set_cfg!(cooldown, { - self.config.lock().await.cooldown = cooldown; - }); - - self.terminal_server - .broadcast(TerminalServerMessage::ConfigUpdated( - self.get_config_status().await, - )) - .await?; } _ => { @@ -113,6 +89,12 @@ impl Engine { .await?; } } + + self.terminal_server + .broadcast(TerminalServerMessage::ConfigUpdated( + self.get_config_status().await, + )) + .await?; } else { self.invalid_command_usage("engine::config").await?; } @@ -169,7 +151,7 @@ impl Engine { .terminal_server .error( "engine::account", - &format!("Account not found ({new_active})"), + &format!("account not found ({new_active})"), ) .await; } @@ -179,7 +161,7 @@ impl Engine { self.terminal_server .info( "engine::account", - &format!("Account set to {new_active} successfully!"), + &format!("account set to {new_active} successfully!"), ) .await?; } @@ -190,6 +172,80 @@ impl Engine { } } + "preset" => { + if args.len() == 0 { + return self.invalid_command_usage("engine::preset").await; + } + + match args[0] { + "list" | "ls" => { + self.terminal_server + .info("engine::preset", "PRESET LIST") + .await?; + + let config = self.config.lock().await; + + for (name, cfg) in &config.configs { + let description = if let Some(desc) = &cfg.description { + format!("- {desc}") + } else { + String::new() + }; + + self.terminal_server + .info( + "engine::preset", + &if name == &config.preset { + format!("{name} (active) {description}") + } else { + format!("{name} {description}") + }, + ) + .await?; + } + } + + "use" | "set" => { + if args.len() < 2 { + return self.invalid_command_usage("engine::preset").await; + } + + let new_active = args[1]; + + let mut config = self.config.lock().await; + + if !config.configs.contains_key(new_active) { + return self + .terminal_server + .error( + "engine::preset", + &format!("preset not found ({new_active})"), + ) + .await; + } + + config.preset = new_active.to_string(); + + self.terminal_server + .info( + "engine::preset", + &format!("preset set to {new_active} successfully!"), + ) + .await?; + + self.terminal_server + .broadcast(TerminalServerMessage::ConfigUpdated( + self.get_config_status().await, + )) + .await?; + } + + _ => { + return self.invalid_command_usage("engine::preset").await; + } + } + } + "strategy" | "strat" | "sg" => { let Some(strategy_command) = args.drain(0..=0).next() else { return self.invalid_command_usage("strategy").await; @@ -211,7 +267,7 @@ impl Engine { _ => { self.terminal_server - .error("engine::cmd", &format!("Command '{}' not found", command)) + .error("engine::cmd", &format!("command '{}' not found", command)) .await?; } } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 492dc7f..8e3aa98 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -5,7 +5,7 @@ pub mod terminal; use crate::{ engine::{strategy::StrategyEngine, terminal::TerminalServer}, - store::{accounts::AccountList, config::Config}, + store::{accounts::AccountList, config::ConfigManager}, }; use hypersdk::hypercore::ws::ConnectionStream; use pulse_sdk::prelude::*; @@ -26,7 +26,7 @@ pub struct Engine { pub ws_stream: Arc>, // data - pub config: Arc>, + pub config: Arc>, pub accounts: Arc>, // live data / status @@ -37,11 +37,11 @@ pub struct Engine { impl Engine { pub async fn new() -> tokio::io::Result> { - let config = Config::new().await?; + let config = ConfigManager::new().await?; let (ws_handle, ws_stream) = hypersdk::hypercore::mainnet_ws().split(); - let strategy = StrategyEngine::new(&config.strategy, ws_handle).await?; + let strategy = StrategyEngine::new(&config.get_ref()?.strategy, ws_handle).await?; let accounts = Arc::new(Mutex::new(AccountList::new().await?)); let config = Arc::new(Mutex::new(config)); @@ -91,6 +91,7 @@ impl Engine { pub async fn get_config_status(&self) -> EngineConfig { EngineConfig { + preset: self.config.lock().await.preset.clone(), strategy: self.strategy_engine.strategy.lock().await.manifest.clone(), } } diff --git a/src/engine/engine/terminal.rs b/src/engine/engine/terminal.rs index f2e6211..39003ad 100644 --- a/src/engine/engine/terminal.rs +++ b/src/engine/engine/terminal.rs @@ -194,8 +194,11 @@ impl TerminalServer { loop { refresh.tick().await; - match crate::fetch::fetch_watch_list(&client, &engine.config.lock().await.watchlist) - .await + match crate::fetch::fetch_watch_list( + &client, + &engine.config.lock().await.get_ref()?.watchlist, + ) + .await { Ok(watch_list) => { *engine.watch_list.lock().await = watch_list.clone(); diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index ba42aa2..7dfd733 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -1,23 +1,60 @@ -use hypersdk::hypercore::CandleInterval; +use std::collections::HashMap; #[derive(Debug, serde::Serialize, serde::Deserialize)] pub struct Config { + pub description: Option, pub watchlist: Vec, pub strategy: String, - pub cooldown: CandleInterval, +} + +#[derive(Debug, serde::Serialize, serde::Deserialize)] +pub struct ConfigManager { + pub preset: String, + + #[serde(flatten)] + pub configs: HashMap, } impl Default for Config { fn default() -> Self { Self { + description: None, watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()], strategy: String::new(), - cooldown: CandleInterval::ThirtyMinutes, } } } -impl Config { +impl Default for ConfigManager { + fn default() -> Self { + Self { + preset: "Main".to_string(), + configs: HashMap::from([("Main".to_string(), Config::default())]), + } + } +} + +impl ConfigManager { + pub fn get_ref(&self) -> tokio::io::Result<&Config> { + self.configs.get(&self.preset).ok_or_else(|| { + tokio::io::Error::new( + std::io::ErrorKind::NotFound, + format!("Failed to get preset '{}'", self.preset), + ) + }) + } + + pub fn get_mut(&mut self) -> tokio::io::Result<&mut Config> { + self.configs.get_mut(&self.preset).ok_or_else(|| { + tokio::io::Error::new( + std::io::ErrorKind::NotFound, + format!("Failed to get preset '{}'", self.preset), + ) + }) + } +} + +impl ConfigManager { pub async fn new() -> tokio::io::Result { let path = crate::store::pulse_config_file()?; diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 64744a6..b8fb1a3 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -217,11 +217,11 @@ impl Formatted for EngineConfig { fn get_formatted(&self) -> Vec { vec![ Triple( + "\x1b[2mPreset\x1b[0m", "\x1b[2mStrategy\x1b[0m", "\x1b[2m..\x1b[0m", - "\x1b[2m..\x1b[0m", ), - Triple(&self.strategy.name, "..", ".."), + Triple(&self.preset, &self.strategy.name, ".."), ] .get_formatted() } From 7e4e264eb00707395109ed006821726ba2aad98b Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 31 Jul 2026 20:13:31 +0200 Subject: [PATCH 02/26] Risk config --- pulse-sdk/src/general.rs | 71 +++++++++++++++++++++++++++++++++++++- src/engine/store/config.rs | 21 +++++++++++ 2 files changed, 91 insertions(+), 1 deletion(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 03582af..56e6253 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -1,4 +1,4 @@ -use hypersdk::hypercore::Side; +use hypersdk::{dec, hypercore::Side}; use rust_decimal::Decimal; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -30,6 +30,12 @@ pub enum Mode { Manual, } +#[derive(Debug, Clone)] +pub enum Allocation { + Fixed(Decimal), + Percent(Decimal), +} + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Signal { pub symbol: String, @@ -101,3 +107,66 @@ impl std::fmt::Display for ItemState { } } } + +impl<'de> serde::Deserialize<'de> for Allocation { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + struct AllocationVisitor; + + impl<'de> serde::de::Visitor<'de> for AllocationVisitor { + type Value = Allocation; + + fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result { + formatter.write_str("a number or percentage string") + } + + fn visit_f64(self, value: f64) -> Result + where + E: serde::de::Error, + { + Ok(Allocation::Fixed(Decimal::from_f64_retain(value).unwrap())) + } + + fn visit_str(self, value: &str) -> Result + where + E: serde::de::Error, + { + if let Some(percent) = value.strip_suffix('%') { + let value = percent.parse::().map_err(E::custom)?; + + Ok(Allocation::Percent(value)) + } else { + Err(E::custom("invalid allocation format")) + } + } + } + + deserializer.deserialize_any(AllocationVisitor) + } +} + +impl serde::Serialize for Allocation { + fn serialize(&self, serializer: S) -> Result + where + S: serde::Serializer, + { + match self { + Allocation::Fixed(value) => serde::Serialize::serialize(value, serializer), + Allocation::Percent(value) => { + let s = format!("{}%", value); + serializer.serialize_str(&s) + } + } + } +} + +impl Allocation { + pub fn get(&self, full_alloc: Decimal) -> Decimal { + match self { + Self::Fixed(f) => *f, + Self::Percent(p) => (*p / dec!(100.0)) * full_alloc, + } + } +} diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index 7dfd733..7e321fe 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -1,10 +1,20 @@ use std::collections::HashMap; +use pulse_sdk::general::Allocation; + +#[derive(Debug, serde::Serialize, serde::Deserialize)] +pub struct RiskConfig { + pub risk_per_trade: Allocation, + pub max_open_positions: u32, + pub max_daily_loss: Allocation, +} + #[derive(Debug, serde::Serialize, serde::Deserialize)] pub struct Config { pub description: Option, pub watchlist: Vec, pub strategy: String, + pub risk: RiskConfig, } #[derive(Debug, serde::Serialize, serde::Deserialize)] @@ -15,12 +25,23 @@ pub struct ConfigManager { pub configs: HashMap, } +impl Default for RiskConfig { + fn default() -> Self { + Self { + risk_per_trade: Allocation::Percent(10.into()), + max_daily_loss: Allocation::Percent(20.into()), + max_open_positions: 1, + } + } +} + impl Default for Config { fn default() -> Self { Self { description: None, watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()], strategy: String::new(), + risk: Default::default(), } } } From a9ccd84633545503e795aa45c68a67db604fefd6 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Fri, 31 Jul 2026 23:27:56 +0200 Subject: [PATCH 03/26] Risk engine --- src/engine/engine/execution.rs | 40 +++++++++++++++++++++---- src/engine/engine/mod.rs | 5 +++- src/engine/engine/risk.rs | 55 ++++++++++++++++++++++++++++++++++ src/engine/engine/strategy.rs | 18 +++++++++-- 4 files changed, 109 insertions(+), 9 deletions(-) create mode 100644 src/engine/engine/risk.rs diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index a7370c0..88a4c4d 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -1,7 +1,9 @@ -use hypersdk::hypercore::{self, BatchOrder, OrderRequest, OrderTypePlacement, Side, TimeInForce}; +use hypersdk::hypercore::{ + self, BatchOrder, OrderRequest, OrderResponseStatus, OrderTypePlacement, Side, TimeInForce, +}; use pulse_sdk::prelude::*; -use crate::engine::Engine; +use crate::engine::{Engine, risk::OrderIds}; impl Engine { pub async fn execute_signal(&self, signal: &Signal) -> tokio::io::Result<()> { @@ -30,6 +32,8 @@ impl Engine { return Ok(()); }; + let order_ids = OrderIds::new(); + let order = BatchOrder { orders: vec![ OrderRequest { @@ -41,7 +45,7 @@ impl Engine { order_type: OrderTypePlacement::Limit { tif: TimeInForce::Gtc, }, - cloid: Default::default(), + cloid: order_ids.entry, }, OrderRequest { asset: asset_id, @@ -54,7 +58,7 @@ impl Engine { trigger_px: signal.take_profit, tpsl: hypercore::TpSl::Tp, }, - cloid: Default::default(), + cloid: order_ids.take_profit, }, OrderRequest { asset: asset_id, @@ -67,7 +71,7 @@ impl Engine { trigger_px: signal.stop_loss, tpsl: hypercore::TpSl::Sl, }, - cloid: Default::default(), + cloid: order_ids.stop_loss, }, ], grouping: hypercore::OrderGrouping::Na, @@ -80,7 +84,31 @@ impl Engine { .place(&acc.private_key.0, order, nonce, None, None) .await { - Ok(_) => {} + Ok(o) if o.iter().any(|o| matches!(o, OrderResponseStatus::Error(_))) => { + self.risk_engine.order_placed(order_ids).await; + } + + Ok(e) => { + self.terminal_server + .error( + "self::order", + &format!( + "order rejected: {}", + e.into_iter() + .filter_map(|res| { + if let OrderResponseStatus::Error(e) = res { + Some(e) + } else { + None + } + }) + .collect::>() + .join(", ") + ), + ) + .await?; + } + Err(e) => { self.terminal_server .error("self::order", &e.to_string()) diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 8e3aa98..9989ac9 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -1,10 +1,11 @@ pub mod command; pub mod execution; +pub mod risk; pub mod strategy; pub mod terminal; use crate::{ - engine::{strategy::StrategyEngine, terminal::TerminalServer}, + engine::{risk::RiskEngine, strategy::StrategyEngine, terminal::TerminalServer}, store::{accounts::AccountList, config::ConfigManager}, }; use hypersdk::hypercore::ws::ConnectionStream; @@ -23,6 +24,7 @@ pub struct Engine { // engine pub terminal_server: Arc, pub strategy_engine: Arc, + pub risk_engine: Arc, pub ws_stream: Arc>, // data @@ -50,6 +52,7 @@ impl Engine { // engine terminal_server: TerminalServer::new(engine.clone()), strategy_engine: strategy.initialize(engine.clone()), + risk_engine: RiskEngine::new(engine.clone()), ws_stream: Arc::new(Mutex::new(ws_stream)), // data diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs new file mode 100644 index 0000000..a16e85d --- /dev/null +++ b/src/engine/engine/risk.rs @@ -0,0 +1,55 @@ +use std::sync::{Arc, Weak}; + +use hypersdk::hypercore::{self, Cloid, WebSocket}; +use pulse_sdk::general::Signal; +use tokio::sync::Mutex; + +use crate::engine::Engine; + +#[derive(Debug, Clone)] +pub struct OrderIds { + pub entry: Cloid, + pub take_profit: Cloid, + pub stop_loss: Cloid, +} + +impl OrderIds { + pub fn new() -> Self { + Self { + entry: Cloid::random(), + take_profit: Cloid::random(), + stop_loss: Cloid::random(), + } + } +} + +pub struct RiskEngine { + // Orders made by the engine + pub orders: Mutex>, + pub ws: Mutex, + pub engine: Weak, +} + +impl RiskEngine { + pub fn new(engine: Weak) -> Arc { + Arc::new(Self { + orders: Mutex::new(Vec::new()), + ws: Mutex::new(hypercore::mainnet_ws()), + engine, + }) + } + + pub async fn validate_signal(&self, _signal: &mut Signal) -> bool { + true + } + + pub async fn order_placed(&self, order: OrderIds) { + self.orders.lock().await.push(order); + } + + pub fn get_engine(&self) -> Arc { + self.engine + .upgrade() + .expect("Failed to upgrade engine(Weak) to Arc") + } +} diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index b79e26e..e358aee 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -117,7 +117,18 @@ impl StrategyEngine { engine.terminal_server.log_raw(log).await?; } - Some(StrategyMessage::Signal(signal)) => { + Some(StrategyMessage::Signal(mut signal)) => { + if !engine.risk_engine.validate_signal(&mut signal).await { + engine.signals.lock().await.push(Err(signal)); + + engine + .terminal_server + .error("engine::risk", "Signal rejected") + .await?; + + continue; + } + match engine.execute_signal(&signal).await { Ok(_) => engine.signals.lock().await.push(Ok(signal)), Err(e) => { @@ -125,7 +136,10 @@ impl StrategyEngine { engine .terminal_server - .error("signal", &format!("Failed to execute signal: {e}")) + .error( + "engine::strategy", + &format!("Failed to execute signal: {e}"), + ) .await? } } From f48e239d055123a099d62944f23c384940a2661f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 06:56:15 +0200 Subject: [PATCH 04/26] Risk stuff and trying to test different ws archs --- src/engine/engine/mod.rs | 58 +++++++++++++++++------------ src/engine/engine/risk.rs | 77 +++++++++++++++++++++++++++++++++++++-- src/engine/main.rs | 2 +- 3 files changed, 108 insertions(+), 29 deletions(-) diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 9989ac9..a11422a 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -25,7 +25,6 @@ pub struct Engine { pub terminal_server: Arc, pub strategy_engine: Arc, pub risk_engine: Arc, - pub ws_stream: Arc>, // data pub config: Arc>, @@ -37,39 +36,50 @@ pub struct Engine { pub signals: Arc>>, } +pub struct EngineSides { + pub strategy_stream: ConnectionStream, + pub risk_stream: ConnectionStream, +} + impl Engine { - pub async fn new() -> tokio::io::Result> { + pub async fn new() -> tokio::io::Result<(Arc, EngineSides)> { let config = ConfigManager::new().await?; - let (ws_handle, ws_stream) = hypersdk::hypercore::mainnet_ws().split(); + let (strategy_handle, strategy_stream) = hypersdk::hypercore::mainnet_ws().split(); + let (risk_handle, risk_stream) = hypersdk::hypercore::mainnet_ws().split(); - let strategy = StrategyEngine::new(&config.get_ref()?.strategy, ws_handle).await?; + let strategy = StrategyEngine::new(&config.get_ref()?.strategy, strategy_handle).await?; let accounts = Arc::new(Mutex::new(AccountList::new().await?)); let config = Arc::new(Mutex::new(config)); - Ok(Arc::new_cyclic(|engine| Self { - // engine - terminal_server: TerminalServer::new(engine.clone()), - strategy_engine: strategy.initialize(engine.clone()), - risk_engine: RiskEngine::new(engine.clone()), - ws_stream: Arc::new(Mutex::new(ws_stream)), + Ok(( + Arc::new_cyclic(|engine| Self { + // engine + terminal_server: TerminalServer::new(engine.clone()), + strategy_engine: strategy.initialize(engine.clone()), + risk_engine: RiskEngine::new(engine.clone(), risk_handle), - // data - config, - accounts, + // data + config, + accounts, - // live data / status - status: Arc::new(Mutex::new(EngineStatus { - strategy_mode: Mode::Auto, - strategy_state: ItemState::Stopped, - })), - watch_list: Arc::new(Mutex::new(WatchList { - name_to_index: HashMap::new(), - items: Vec::new(), - })), - signals: Arc::new(Mutex::new(Vec::new())), - })) + // live data / status + status: Arc::new(Mutex::new(EngineStatus { + strategy_mode: Mode::Auto, + strategy_state: ItemState::Stopped, + })), + watch_list: Arc::new(Mutex::new(WatchList { + name_to_index: HashMap::new(), + items: Vec::new(), + })), + signals: Arc::new(Mutex::new(Vec::new())), + }), + EngineSides { + strategy_stream, + risk_stream, + }, + )) } /// Starts the main strategy server diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index a16e85d..58eb06d 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -1,6 +1,13 @@ use std::sync::{Arc, Weak}; -use hypersdk::hypercore::{self, Cloid, WebSocket}; +use futures::StreamExt; +use hypersdk::{ + Decimal, + hypercore::{ + self, Cloid, + ws::{ConnectionHandle, ConnectionStream}, + }, +}; use pulse_sdk::general::Signal; use tokio::sync::Mutex; @@ -23,22 +30,84 @@ impl OrderIds { } } +pub struct RiskState { + pub starting_equity: Decimal, + pub realized_pnl_today: Decimal, + pub open_positions: usize, +} + pub struct RiskEngine { // Orders made by the engine pub orders: Mutex>, - pub ws: Mutex, + pub handle: Mutex, + pub state: Mutex, pub engine: Weak, } impl RiskEngine { - pub fn new(engine: Weak) -> Arc { + pub fn new(engine: Weak, handle: ConnectionHandle) -> Arc { Arc::new(Self { orders: Mutex::new(Vec::new()), - ws: Mutex::new(hypercore::mainnet_ws()), + handle: Mutex::new(handle), + state: Mutex::new(RiskState { + starting_equity: 0.into(), + realized_pnl_today: 0.into(), + open_positions: 0, + }), engine, }) } + pub async fn day_tick(&self) -> anyhow::Result<()> { + let engine = self.get_engine(); + + let client = hypercore::mainnet(); + + if let Some(acc) = engine.accounts.lock().await.get_active() { + self.state.lock().await.starting_equity = client + .user_vault_equities(acc.address) + .await? + .into_iter() + .map(|v| v.equity) + .sum(); + } + + Ok(()) + } + + pub async fn run(&self, mut stream: ConnectionStream) -> anyhow::Result<()> { + while let Some(e) = stream.next().await { + match e { + hypercore::ws::Event::Message(hypercore::Incoming::OrderUpdates(order_updates)) => { + for order in order_updates { + if let Some(cloid) = order.order.cloid { + let orders = self.orders.lock().await.clone(); + + let mut rm = Vec::new(); + + for (i, order_ids) in orders.iter().enumerate() { + if order_ids.stop_loss == cloid { + if order.status.is_filled() { + rm.push(i); + } + + break; + } + } + + for r in rm { + self.orders.lock().await.remove(r); + } + } + } + } + _ => {} + } + } + + Ok(()) + } + pub async fn validate_signal(&self, _signal: &mut Signal) -> bool { true } diff --git a/src/engine/main.rs b/src/engine/main.rs index c9d3244..8531472 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -4,7 +4,7 @@ pub mod store; #[tokio::main] async fn main() -> anyhow::Result<()> { - let engine = engine::Engine::new().await?; + let (engine, sides) = engine::Engine::new().await?; let server = engine.terminal_server.spawn_server().await; let broadcaster = engine.terminal_server.spawn_broadcaster().await; From fdb0525ca2a6bab0b162b72c3101def3018f218b Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 19:03:34 +0200 Subject: [PATCH 05/26] Strategy and Risk websocket clients --- src/engine/engine/risk.rs | 36 ++++++++++++++++++++++--- src/engine/engine/strategy.rs | 49 +++++++++++++++++++++++++++++++---- src/engine/main.rs | 17 ++++++++++++ 3 files changed, 93 insertions(+), 9 deletions(-) diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index 58eb06d..428c034 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -5,7 +5,7 @@ use hypersdk::{ Decimal, hypercore::{ self, Cloid, - ws::{ConnectionHandle, ConnectionStream}, + ws::{ConnectionHandle, ConnectionStream, Event}, }, }; use pulse_sdk::general::Signal; @@ -75,10 +75,15 @@ impl RiskEngine { Ok(()) } - pub async fn run(&self, mut stream: ConnectionStream) -> anyhow::Result<()> { + pub async fn run_event_stream( + self: Arc, + mut stream: ConnectionStream, + ) -> anyhow::Result<()> { + let engine = self.get_engine(); + while let Some(e) = stream.next().await { match e { - hypercore::ws::Event::Message(hypercore::Incoming::OrderUpdates(order_updates)) => { + Event::Message(hypercore::Incoming::OrderUpdates(order_updates)) => { for order in order_updates { if let Some(cloid) = order.order.cloid { let orders = self.orders.lock().await.clone(); @@ -101,7 +106,30 @@ impl RiskEngine { } } } - _ => {} + + Event::Connected => { + engine + .terminal_server + .info("engine::risk", "HyperLiquid WebSocket connected") + .await?; + } + + Event::Disconnected => { + engine + .terminal_server + .warn("engine::risk", "HyperLiquid WebSocket disconnected") + .await?; + } + + _ => { + engine + .terminal_server + .warn( + "engine::risk", + &format!("Unexpected event from HyperLiquid WebSocket: {e:?}"), + ) + .await?; + } } } diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index e358aee..56621fc 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -1,5 +1,9 @@ use anyhow::Context; -use hypersdk::hypercore::{self, CandleInterval, Subscription, ws::ConnectionHandle}; +use futures::StreamExt; +use hypersdk::hypercore::{ + self, CandleInterval, Subscription, + ws::{ConnectionHandle, ConnectionStream, Event}, +}; use pulse_sdk::prelude::*; use std::{ collections::HashSet, @@ -75,11 +79,40 @@ impl StrategyEngine { Ok(()) } + pub async fn run_event_stream( + self: Arc, + mut stream: ConnectionStream, + ) -> anyhow::Result<()> { + let engine = self.get_engine(); + + while let Some(e) = stream.next().await { + match e { + Event::Message(incoming) => { + self.send(&StrategyEngineMessage::Incoming(incoming)) + .await?; + } + + Event::Connected => { + engine + .terminal_server + .info("engine::strategy", "HyperLiquid WebSocket connected") + .await?; + } + + Event::Disconnected => { + engine + .terminal_server + .warn("engine::strategy", "HyperLiquid WebSocket disconnected") + .await?; + } + } + } + + Ok(()) + } + pub async fn run(self: &Arc) -> anyhow::Result<()> { - let engine = self - .engine - .upgrade() - .expect("Failed to upgrade engine (StrategyEngine)"); + let engine = self.get_engine(); let mut stdout = { let mut child = self.strategy.lock().await; @@ -211,4 +244,10 @@ impl StrategyEngine { } } } + + pub fn get_engine(&self) -> Arc { + self.engine + .upgrade() + .expect("Failed to upgrade engine (StrategyEngine)") + } } diff --git a/src/engine/main.rs b/src/engine/main.rs index 8531472..e8be395 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -9,10 +9,27 @@ async fn main() -> anyhow::Result<()> { let server = engine.terminal_server.spawn_server().await; let broadcaster = engine.terminal_server.spawn_broadcaster().await; + let risk = tokio::spawn( + engine + .risk_engine + .clone() + .run_event_stream(sides.risk_stream), + ); + + let strategy = tokio::spawn( + engine + .strategy_engine + .clone() + .run_event_stream(sides.strategy_stream), + ); + engine.run().await?; server.await??; broadcaster.await??; + risk.await??; + strategy.await??; + Ok(()) } From f64d12d8c78823dac4b173da85e434e48b142a7b Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 19:16:38 +0200 Subject: [PATCH 06/26] Stop loss order management --- src/engine/engine/execution.rs | 4 ++-- src/engine/engine/risk.rs | 27 +++++++++++++++++++++++---- 2 files changed, 25 insertions(+), 6 deletions(-) diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index 88a4c4d..5eded2c 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -3,7 +3,7 @@ use hypersdk::hypercore::{ }; use pulse_sdk::prelude::*; -use crate::engine::{Engine, risk::OrderIds}; +use crate::engine::Engine; impl Engine { pub async fn execute_signal(&self, signal: &Signal) -> tokio::io::Result<()> { @@ -32,7 +32,7 @@ impl Engine { return Ok(()); }; - let order_ids = OrderIds::new(); + let order_ids = self.risk_engine.create_order().await?; let order = BatchOrder { orders: vec![ diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index 428c034..923a9a0 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -18,21 +18,23 @@ pub struct OrderIds { pub entry: Cloid, pub take_profit: Cloid, pub stop_loss: Cloid, + pub risk_equity: Decimal, } impl OrderIds { - pub fn new() -> Self { + pub fn new(risk_equity: Decimal) -> Self { Self { entry: Cloid::random(), take_profit: Cloid::random(), stop_loss: Cloid::random(), + risk_equity, } } } pub struct RiskState { pub starting_equity: Decimal, - pub realized_pnl_today: Decimal, + pub pnl: Decimal, pub open_positions: usize, } @@ -51,7 +53,7 @@ impl RiskEngine { handle: Mutex::new(handle), state: Mutex::new(RiskState { starting_equity: 0.into(), - realized_pnl_today: 0.into(), + pnl: 0.into(), open_positions: 0, }), engine, @@ -101,7 +103,9 @@ impl RiskEngine { } for r in rm { - self.orders.lock().await.remove(r); + let order = self.orders.lock().await.remove(r); + + self.state.lock().await.pnl -= order.risk_equity; } } } @@ -140,6 +144,21 @@ impl RiskEngine { true } + pub async fn create_order(&self) -> tokio::io::Result { + let state = self.state.lock().await; + + Ok(OrderIds::new( + self.get_engine() + .config + .lock() + .await + .get_ref()? + .risk + .risk_per_trade + .get(state.starting_equity), + )) + } + pub async fn order_placed(&self, order: OrderIds) { self.orders.lock().await.push(order); } From 32b93ce2a598a79a33ec50f76e8242fe3d59508c Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 19:45:12 +0200 Subject: [PATCH 07/26] Validate signal --- pulse-sdk/src/general.rs | 2 +- src/engine/engine/risk.rs | 29 ++++++++++++++++++++--------- src/engine/engine/strategy.rs | 2 +- src/engine/store/config.rs | 2 +- 4 files changed, 23 insertions(+), 12 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 56e6253..f7491aa 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -30,7 +30,7 @@ pub enum Mode { Manual, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, Copy)] pub enum Allocation { Fixed(Decimal), Percent(Decimal), diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index 923a9a0..e49f295 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -11,7 +11,7 @@ use hypersdk::{ use pulse_sdk::general::Signal; use tokio::sync::Mutex; -use crate::engine::Engine; +use crate::{engine::Engine, store::config::RiskConfig}; #[derive(Debug, Clone)] pub struct OrderIds { @@ -74,6 +74,8 @@ impl RiskEngine { .sum(); } + self.state.lock().await.pnl = 0.into(); + Ok(()) } @@ -140,20 +142,25 @@ impl RiskEngine { Ok(()) } - pub async fn validate_signal(&self, _signal: &mut Signal) -> bool { - true + pub async fn validate_signal(&self, _signal: &mut Signal) -> tokio::io::Result { + let state = self.state.lock().await; + + let under_max_losses = state.pnl + < self + .get_risk_config() + .await? + .max_daily_loss + .get(state.starting_equity); + + Ok(under_max_losses) } pub async fn create_order(&self) -> tokio::io::Result { let state = self.state.lock().await; Ok(OrderIds::new( - self.get_engine() - .config - .lock() - .await - .get_ref()? - .risk + self.get_risk_config() + .await? .risk_per_trade .get(state.starting_equity), )) @@ -163,6 +170,10 @@ impl RiskEngine { self.orders.lock().await.push(order); } + pub async fn get_risk_config(&self) -> tokio::io::Result { + Ok(self.get_engine().config.lock().await.get_ref()?.risk) + } + pub fn get_engine(&self) -> Arc { self.engine .upgrade() diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index 56621fc..88e02e1 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -151,7 +151,7 @@ impl StrategyEngine { } Some(StrategyMessage::Signal(mut signal)) => { - if !engine.risk_engine.validate_signal(&mut signal).await { + if !engine.risk_engine.validate_signal(&mut signal).await? { engine.signals.lock().await.push(Err(signal)); engine diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index 7e321fe..a39bd9a 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -2,7 +2,7 @@ use std::collections::HashMap; use pulse_sdk::general::Allocation; -#[derive(Debug, serde::Serialize, serde::Deserialize)] +#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] pub struct RiskConfig { pub risk_per_trade: Allocation, pub max_open_positions: u32, From a56b0004cf6bdc3d3a6488c7f648007bda5a05da Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 19:58:10 +0200 Subject: [PATCH 08/26] Updated EngineConfig --- pulse-sdk/src/general.rs | 9 +++++++++ pulse-sdk/src/terminal.rs | 6 +++++- src/engine/engine/command.rs | 4 ++-- src/engine/engine/mod.rs | 15 +++++++++++---- src/engine/engine/terminal.rs | 2 +- src/terminal/formatting.rs | 8 ++++++-- 6 files changed, 34 insertions(+), 10 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index f7491aa..70674d9 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -170,3 +170,12 @@ impl Allocation { } } } + +impl std::fmt::Display for Allocation { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Fixed(v) => write!(f, "${v}"), + Self::Percent(p) => write!(f, "{p}%"), + } + } +} diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 168cc73..d985d49 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,5 +1,5 @@ use crate::{ - general::{EngineStatus, EventLog, Position, Signal}, + general::{Allocation, EngineStatus, EventLog, Position, Signal}, strategy::StrategyManifest, }; use hypersdk::Decimal; @@ -65,4 +65,8 @@ pub enum InspectItem { pub struct EngineConfig { pub preset: String, pub strategy: StrategyManifest, + + pub risk_per_trade: Allocation, + pub max_open_positions: u32, + pub max_daily_loss: Allocation, } diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 06724f6..8fa4f9d 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -92,7 +92,7 @@ impl Engine { self.terminal_server .broadcast(TerminalServerMessage::ConfigUpdated( - self.get_config_status().await, + self.get_config_status().await?, )) .await?; } else { @@ -235,7 +235,7 @@ impl Engine { self.terminal_server .broadcast(TerminalServerMessage::ConfigUpdated( - self.get_config_status().await, + self.get_config_status().await?, )) .await?; } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index a11422a..7e62fd7 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -102,11 +102,18 @@ impl Engine { } } - pub async fn get_config_status(&self) -> EngineConfig { - EngineConfig { - preset: self.config.lock().await.preset.clone(), + pub async fn get_config_status(&self) -> tokio::io::Result { + let manager = self.config.lock().await; + let config = manager.get_ref()?; + + Ok(EngineConfig { + preset: manager.preset.clone(), strategy: self.strategy_engine.strategy.lock().await.manifest.clone(), - } + + risk_per_trade: config.risk.risk_per_trade, + max_daily_loss: config.risk.max_daily_loss, + max_open_positions: config.risk.max_open_positions, + }) } pub async fn update_status(&self) -> tokio::io::Result<()> { diff --git a/src/engine/engine/terminal.rs b/src/engine/engine/terminal.rs index 39003ad..40a4c4d 100644 --- a/src/engine/engine/terminal.rs +++ b/src/engine/engine/terminal.rs @@ -91,7 +91,7 @@ impl TerminalServer { self.send_to( id, - TerminalServerMessage::ConfigUpdated(engine.get_config_status().await), + TerminalServerMessage::ConfigUpdated(engine.get_config_status().await?), ) .await?; diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index b8fb1a3..33c02c1 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -219,9 +219,13 @@ impl Formatted for EngineConfig { Triple( "\x1b[2mPreset\x1b[0m", "\x1b[2mStrategy\x1b[0m", - "\x1b[2m..\x1b[0m", + "\x1b[2mRisk\x1b[0m", + ), + Triple( + &self.preset, + &self.strategy.name, + &self.risk_per_trade.to_string(), ), - Triple(&self.preset, &self.strategy.name, ".."), ] .get_formatted() } From bfac0c340f02dd56d4861ce9fab9c8e50f807bc9 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 20:25:50 +0200 Subject: [PATCH 09/26] Added more formatting options --- pulse-sdk/src/terminal.rs | 7 ++++--- src/engine/engine/mod.rs | 5 +++-- src/terminal/formatting.rs | 18 ++++++++++++++++-- src/terminal/terminal.rs | 22 +++++++++++++++++++++- 4 files changed, 44 insertions(+), 8 deletions(-) diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index d985d49..db231be 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,5 +1,5 @@ use crate::{ - general::{Allocation, EngineStatus, EventLog, Position, Signal}, + general::{EngineStatus, EventLog, Position, Signal}, strategy::StrategyManifest, }; use hypersdk::Decimal; @@ -65,8 +65,9 @@ pub enum InspectItem { pub struct EngineConfig { pub preset: String, pub strategy: StrategyManifest, + pub description: Option, - pub risk_per_trade: Allocation, + pub risk_per_trade: String, pub max_open_positions: u32, - pub max_daily_loss: Allocation, + pub max_daily_loss: String, } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 7e62fd7..b65245f 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -109,9 +109,10 @@ impl Engine { Ok(EngineConfig { preset: manager.preset.clone(), strategy: self.strategy_engine.strategy.lock().await.manifest.clone(), + description: config.description.clone(), - risk_per_trade: config.risk.risk_per_trade, - max_daily_loss: config.risk.max_daily_loss, + risk_per_trade: config.risk.risk_per_trade.to_string(), + max_daily_loss: config.risk.max_daily_loss.to_string(), max_open_positions: config.risk.max_open_positions, }) } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 33c02c1..623e386 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -219,12 +219,26 @@ impl Formatted for EngineConfig { Triple( "\x1b[2mPreset\x1b[0m", "\x1b[2mStrategy\x1b[0m", - "\x1b[2mRisk\x1b[0m", + "\x1b[2mDescription\x1b[0m", ), Triple( &self.preset, &self.strategy.name, - &self.risk_per_trade.to_string(), + self.description + .as_ref() + .map(String::as_str) + .unwrap_or("none"), + ), + Triple("", "", ""), + Triple( + "\x1b[2mRisk/Trade\x1b[0m", + "\x1b[2mMax Positions\x1b[0m", + "\x1b[2mMax daily loss\x1b[0m", + ), + Triple( + &self.risk_per_trade, + &self.max_open_positions.to_string(), + &self.max_daily_loss, ), ] .get_formatted() diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index 9acf585..fe8aea3 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -66,7 +66,27 @@ impl TerminalClient { ) .await { - panic!("{v}") + crossterm::terminal::disable_raw_mode().unwrap(); + + crossterm::execute!(std::io::stdout(), crossterm::cursor::Show).unwrap(); + + crossterm::execute!( + std::io::stdout(), + crossterm::terminal::Clear(crossterm::terminal::ClearType::Purge) + ) + .unwrap(); + + crossterm::execute!( + std::io::stdout(), + crossterm::terminal::Clear(crossterm::terminal::ClearType::All) + ) + .unwrap(); + + crossterm::execute!(std::io::stdout(), crossterm::cursor::MoveTo(0, 0)).unwrap(); + + println!("Error using app: {v}"); + + std::process::exit(1); } }); From e2cf0a00d60148e762ff64bb3e18700944ff172e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 21:30:16 +0200 Subject: [PATCH 10/26] Daily scheduler --- pulse-sdk/src/general.rs | 1 + src/engine/engine/mod.rs | 38 ++++++++++++++++++++++++++++++++++++-- src/engine/engine/risk.rs | 2 ++ src/engine/main.rs | 12 ++++++++---- 4 files changed, 47 insertions(+), 6 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 70674d9..1b8dd78 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -63,6 +63,7 @@ pub struct Position { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct EngineStatus { + pub equity: Decimal, pub strategy_mode: Mode, pub strategy_state: ItemState, } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index b65245f..80a9f91 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -10,8 +10,8 @@ use crate::{ }; use hypersdk::hypercore::ws::ConnectionStream; use pulse_sdk::prelude::*; -use std::{collections::HashMap, sync::Arc}; -use tokio::sync::Mutex; +use std::{collections::HashMap, sync::Arc, time::Duration}; +use tokio::{sync::Mutex, time::Instant}; #[derive(Debug, Clone)] pub struct WatchList { @@ -66,6 +66,7 @@ impl Engine { // live data / status status: Arc::new(Mutex::new(EngineStatus { + equity: 0.into(), strategy_mode: Mode::Auto, strategy_state: ItemState::Stopped, })), @@ -102,6 +103,39 @@ impl Engine { } } + pub async fn run_daily_scheduler(self: Arc) -> anyhow::Result<()> { + loop { + self.risk_engine.day_tick().await?; + + let now_local = chrono::Local::now(); + + let midnight = chrono::NaiveTime::from_hms_opt(0, 0, 0).unwrap(); + + let tomorrow_local = now_local + .date_naive() + .succ_opt() + .unwrap() + .and_time(midnight) + .and_local_timezone(chrono::Local) + .unwrap(); + + let duration_until_midnight = tomorrow_local.signed_duration_since(now_local); + + let std_duration = Duration::from_secs(duration_until_midnight.num_seconds() as u64); + + let deadline = Instant::now() + std_duration; + + self.terminal_server + .info( + "engine::schedule", + &format!("next daily reset: {tomorrow_local}"), + ) + .await?; + + tokio::time::sleep_until(deadline).await; + } + } + pub async fn get_config_status(&self) -> tokio::io::Result { let manager = self.config.lock().await; let config = manager.get_ref()?; diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index e49f295..e3b56d8 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -61,6 +61,8 @@ impl RiskEngine { } pub async fn day_tick(&self) -> anyhow::Result<()> { + println!("TICK"); + let engine = self.get_engine(); let client = hypercore::mainnet(); diff --git a/src/engine/main.rs b/src/engine/main.rs index e8be395..fa87e1b 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -23,13 +23,17 @@ async fn main() -> anyhow::Result<()> { .run_event_stream(sides.strategy_stream), ); + let daily_scheduler = tokio::spawn(engine.clone().run_daily_scheduler()); + engine.run().await?; - server.await??; - broadcaster.await??; + server.abort(); + broadcaster.abort(); - risk.await??; - strategy.await??; + risk.abort(); + strategy.abort(); + + daily_scheduler.abort(); Ok(()) } From 08b80d04de39dcc1f771cde24326ca6ac8d647e3 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 21:33:58 +0200 Subject: [PATCH 11/26] Trading equity --- src/engine/engine/risk.rs | 8 +++++--- src/terminal/formatting.rs | 4 ++-- 2 files changed, 7 insertions(+), 5 deletions(-) diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index e3b56d8..ba45dcc 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -61,12 +61,12 @@ impl RiskEngine { } pub async fn day_tick(&self) -> anyhow::Result<()> { - println!("TICK"); - let engine = self.get_engine(); let client = hypercore::mainnet(); + let mut state = self.state.lock().await; + if let Some(acc) = engine.accounts.lock().await.get_active() { self.state.lock().await.starting_equity = client .user_vault_equities(acc.address) @@ -74,9 +74,11 @@ impl RiskEngine { .into_iter() .map(|v| v.equity) .sum(); + + engine.status.lock().await.equity = state.starting_equity; } - self.state.lock().await.pnl = 0.into(); + state.pnl = 0.into(); Ok(()) } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 623e386..3ca77ca 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -199,14 +199,14 @@ impl Formatted for EngineStatus { fn get_formatted(&self) -> Vec { vec![ Triple( + "\x1b[2mTrading Equity\x1b[0m", "\x1b[2mState\x1b[0m", "\x1b[2mMode\x1b[0m", - "\x1b[2m..\x1b[0m", ), Triple( + &format_usd(self.equity.as_f64()), &self.strategy_state.to_string(), &self.strategy_mode.to_string(), - "..", ), ] .get_formatted() From 8f4c984486be461ac81577e81e44897d31657181 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 22:23:31 +0200 Subject: [PATCH 12/26] Perps equity --- pulse-sdk/src/general.rs | 2 +- src/engine/engine/mod.rs | 2 +- src/engine/engine/risk.rs | 9 ++------- src/engine/engine/terminal.rs | 5 +++++ src/terminal/formatting.rs | 8 ++++---- 5 files changed, 13 insertions(+), 13 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 1b8dd78..fd70443 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -63,7 +63,7 @@ pub struct Position { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct EngineStatus { - pub equity: Decimal, + pub perps_equity: Decimal, pub strategy_mode: Mode, pub strategy_state: ItemState, } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 80a9f91..c062433 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -66,7 +66,7 @@ impl Engine { // live data / status status: Arc::new(Mutex::new(EngineStatus { - equity: 0.into(), + perps_equity: 0.into(), strategy_mode: Mode::Auto, strategy_state: ItemState::Stopped, })), diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index ba45dcc..cee9b65 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -68,14 +68,9 @@ impl RiskEngine { let mut state = self.state.lock().await; if let Some(acc) = engine.accounts.lock().await.get_active() { - self.state.lock().await.starting_equity = client - .user_vault_equities(acc.address) - .await? - .into_iter() - .map(|v| v.equity) - .sum(); + let clearing_house = client.clearinghouse_state(acc.address, None).await?; - engine.status.lock().await.equity = state.starting_equity; + state.starting_equity = clearing_house.margin_summary.account_value; } state.pnl = 0.into(); diff --git a/src/engine/engine/terminal.rs b/src/engine/engine/terminal.rs index 40a4c4d..20d700e 100644 --- a/src/engine/engine/terminal.rs +++ b/src/engine/engine/terminal.rs @@ -227,6 +227,11 @@ impl TerminalServer { if let Some(acc) = engine.accounts.lock().await.get_active() { match client.clearinghouse_state(acc.address, None).await { Ok(state) => { + engine.status.lock().await.perps_equity = + state.margin_summary.account_value; + + engine.update_status().await?; + self.broadcast(TerminalServerMessage::PositionsUpdated( state .asset_positions diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 3ca77ca..17bd12e 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -199,12 +199,12 @@ impl Formatted for EngineStatus { fn get_formatted(&self) -> Vec { vec![ Triple( - "\x1b[2mTrading Equity\x1b[0m", - "\x1b[2mState\x1b[0m", - "\x1b[2mMode\x1b[0m", + "\x1b[2mPerps Equity\x1b[0m", + "\x1b[2mStrategy State\x1b[0m", + "\x1b[2mStrategy Mode\x1b[0m", ), Triple( - &format_usd(self.equity.as_f64()), + &format_usd(self.perps_equity.as_f64()), &self.strategy_state.to_string(), &self.strategy_mode.to_string(), ), From cab27466f26942f7dba152ea83d125f983148548 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 23:01:24 +0200 Subject: [PATCH 13/26] Status and positions linked --- pulse-sdk/src/general.rs | 1 + pulse-sdk/src/terminal.rs | 3 --- src/engine/engine/mod.rs | 1 + src/engine/engine/terminal.rs | 31 +++++++++++++++---------------- src/terminal/main.rs | 27 ++++++++++++++++++++++++--- src/terminal/terminal.rs | 7 ------- 6 files changed, 41 insertions(+), 29 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index fd70443..565523e 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -66,6 +66,7 @@ pub struct EngineStatus { pub perps_equity: Decimal, pub strategy_mode: Mode, pub strategy_state: ItemState, + pub positions: Vec, } impl std::fmt::Display for MarketTrend { diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index db231be..a046ace 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -11,9 +11,6 @@ pub enum TerminalServerMessage { // WatchList WatchListUpdated(Vec), - // Positions - PositionsUpdated(Vec), - // Configuration ConfigUpdated(EngineConfig), diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index c062433..141682d 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -69,6 +69,7 @@ impl Engine { perps_equity: 0.into(), strategy_mode: Mode::Auto, strategy_state: ItemState::Stopped, + positions: Vec::new(), })), watch_list: Arc::new(Mutex::new(WatchList { name_to_index: HashMap::new(), diff --git a/src/engine/engine/terminal.rs b/src/engine/engine/terminal.rs index 20d700e..20ab85b 100644 --- a/src/engine/engine/terminal.rs +++ b/src/engine/engine/terminal.rs @@ -227,24 +227,23 @@ impl TerminalServer { if let Some(acc) = engine.accounts.lock().await.get_active() { match client.clearinghouse_state(acc.address, None).await { Ok(state) => { - engine.status.lock().await.perps_equity = - state.margin_summary.account_value; + let mut status = engine.status.lock().await; + + status.perps_equity = state.margin_summary.account_value; + status.positions = state + .asset_positions + .into_iter() + .map(|position| Position { + symbol: position.position.coin, + size: position.position.szi, + entry_price: position.position.entry_px.unwrap_or_default(), + pnl: position.position.unrealized_pnl, + }) + .collect(); + + drop(status); engine.update_status().await?; - - self.broadcast(TerminalServerMessage::PositionsUpdated( - state - .asset_positions - .into_iter() - .map(|position| Position { - symbol: position.position.coin, - size: position.position.szi, - entry_price: position.position.entry_px.unwrap_or_default(), - pnl: position.position.unrealized_pnl, - }) - .collect(), - )) - .await?; } Err(e) => { self.error("orders", &format!("Unable to get open orders: {e}")) diff --git a/src/terminal/main.rs b/src/terminal/main.rs index dc8cb13..05aa89a 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -33,7 +33,6 @@ pub struct PulseTradeApp { status: State>, watch_list: State>, - active_positions: State>, logs: State>, signals: State>, inspect: State, @@ -128,7 +127,10 @@ impl App for PulseTradeApp { ( LayoutItem::Widget(Size::Flex(1)), Box::new( - advanced_draw(&self.scroll, 1, "POSITIONS", &self.active_positions).await, + advanced_draw_map(&self.scroll, 1, "POSITIONS", &self.status, |status| { + &status.positions + }) + .await, ), ), ( @@ -177,7 +179,6 @@ async fn main() -> tokio::io::Result<()> { command: ctx.use_state(InputState::new()), 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()), logs: ctx.use_state(Vec::new()), inspect: ctx.use_state(InspectTarget::None), @@ -190,6 +191,26 @@ async fn main() -> tokio::io::Result<()> { Ok(()) } +pub async fn advanced_draw_map( + scroll: &State>, + index: usize, + title: &'static str, + state: &State>, + m: impl for<'a> FnOnce(&'a T) -> &'a B, +) -> ScrollText { + let scroll = scroll.lock().await; + + scroll.scroll( + index, + format!(" {}{title}\x1b[0m", scroll.get_selected(index)), + if let Some(state) = &*state.lock().await { + apply_padding(m(state).get_formatted()).join("\n") + } else { + " Loading..".to_string() + }, + ) +} + pub async fn advanced_draw( scroll: &State>, index: usize, diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index fe8aea3..d6b7ec1 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -46,7 +46,6 @@ impl TerminalClient { let 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.config.clone(); @@ -57,7 +56,6 @@ impl TerminalClient { if let Err(v) = Self::run_client( reader, watch_list, - active_positions, logs, signals, market_overview, @@ -99,7 +97,6 @@ impl TerminalClient { mut reader: OwnedReadHalf, watch_list: State>, - active_positions: State>, logs: State>, signals: State>, config: State>, @@ -127,10 +124,6 @@ impl TerminalClient { *watch_list.lock().await = v; } - TerminalServerMessage::PositionsUpdated(v) => { - *active_positions.lock().await = v; - } - TerminalServerMessage::ConfigUpdated(v) => { *config.lock().await = Some(v); } From 51998aad684f9c2eafc3bd30b1b6456fe48f13f9 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 1 Aug 2026 23:54:28 +0200 Subject: [PATCH 14/26] Removed unused --- pulse-sdk/src/terminal.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index a046ace..35e923c 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,5 +1,5 @@ use crate::{ - general::{EngineStatus, EventLog, Position, Signal}, + general::{EngineStatus, EventLog, Signal}, strategy::StrategyManifest, }; use hypersdk::Decimal; From 05cf081339b2fe114599a7531ef47e7d338c4063 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 00:19:47 +0200 Subject: [PATCH 15/26] Better signal formatting --- src/engine/engine/execution.rs | 158 +++++++++++++++++---------------- src/engine/engine/risk.rs | 42 +++++---- src/terminal/formatting.rs | 35 ++++---- 3 files changed, 117 insertions(+), 118 deletions(-) diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index 5eded2c..e8f3b51 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -32,89 +32,91 @@ impl Engine { return Ok(()); }; - let order_ids = self.risk_engine.create_order().await?; + let order_ids = self.risk_engine.create_order(signal).await?; - let order = BatchOrder { - orders: vec![ - OrderRequest { - asset: asset_id, - is_buy: matches!(signal.side, Side::Bid), - limit_px: signal.entry_price, - sz: 0.into(), - reduce_only: false, - order_type: OrderTypePlacement::Limit { - tif: TimeInForce::Gtc, - }, - cloid: order_ids.entry, - }, - OrderRequest { - asset: asset_id, - is_buy: matches!(signal.side, Side::Bid), - limit_px: signal.entry_price, - sz: 0.into(), - reduce_only: true, - order_type: OrderTypePlacement::Trigger { - is_market: true, - trigger_px: signal.take_profit, - tpsl: hypercore::TpSl::Tp, - }, - cloid: order_ids.take_profit, - }, - OrderRequest { - asset: asset_id, - is_buy: matches!(signal.side, Side::Bid), - limit_px: signal.entry_price, - sz: 0.into(), - reduce_only: true, - order_type: OrderTypePlacement::Trigger { - is_market: true, - trigger_px: signal.stop_loss, - tpsl: hypercore::TpSl::Sl, - }, - cloid: order_ids.stop_loss, - }, - ], - grouping: hypercore::OrderGrouping::Na, - builder: None, - }; + println!("{order_ids:?}"); - let nonce = chrono::Utc::now().timestamp_millis() as u64; + // let order = BatchOrder { + // orders: vec![ + // OrderRequest { + // asset: asset_id, + // is_buy: matches!(signal.side, Side::Bid), + // limit_px: signal.entry_price, + // sz: 0.into(), + // reduce_only: false, + // order_type: OrderTypePlacement::Limit { + // tif: TimeInForce::Gtc, + // }, + // cloid: order_ids.entry, + // }, + // OrderRequest { + // asset: asset_id, + // is_buy: matches!(signal.side, Side::Bid), + // limit_px: signal.entry_price, + // sz: 0.into(), + // reduce_only: true, + // order_type: OrderTypePlacement::Trigger { + // is_market: true, + // trigger_px: signal.take_profit, + // tpsl: hypercore::TpSl::Tp, + // }, + // cloid: order_ids.take_profit, + // }, + // OrderRequest { + // asset: asset_id, + // is_buy: matches!(signal.side, Side::Bid), + // limit_px: signal.entry_price, + // sz: 0.into(), + // reduce_only: true, + // order_type: OrderTypePlacement::Trigger { + // is_market: true, + // trigger_px: signal.stop_loss, + // tpsl: hypercore::TpSl::Sl, + // }, + // cloid: order_ids.stop_loss, + // }, + // ], + // grouping: hypercore::OrderGrouping::Na, + // builder: None, + // }; - match client - .place(&acc.private_key.0, order, nonce, None, None) - .await - { - Ok(o) if o.iter().any(|o| matches!(o, OrderResponseStatus::Error(_))) => { - self.risk_engine.order_placed(order_ids).await; - } + // let nonce = chrono::Utc::now().timestamp_millis() as u64; - Ok(e) => { - self.terminal_server - .error( - "self::order", - &format!( - "order rejected: {}", - e.into_iter() - .filter_map(|res| { - if let OrderResponseStatus::Error(e) = res { - Some(e) - } else { - None - } - }) - .collect::>() - .join(", ") - ), - ) - .await?; - } + // match client + // .place(&acc.private_key.0, order, nonce, None, None) + // .await + // { + // Ok(o) if o.iter().any(|o| matches!(o, OrderResponseStatus::Error(_))) => { + // self.risk_engine.order_placed(order_ids).await; + // } - Err(e) => { - self.terminal_server - .error("self::order", &e.to_string()) - .await?; - } - } + // Ok(e) => { + // self.terminal_server + // .error( + // "self::order", + // &format!( + // "order rejected: {}", + // e.into_iter() + // .filter_map(|res| { + // if let OrderResponseStatus::Error(e) = res { + // Some(e) + // } else { + // None + // } + // }) + // .collect::>() + // .join(", ") + // ), + // ) + // .await?; + // } + + // Err(e) => { + // self.terminal_server + // .error("self::order", &e.to_string()) + // .await?; + // } + // } } else { self.terminal_server .error("Engine::order", "Unable to get active account") diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index cee9b65..1993247 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -14,22 +14,13 @@ use tokio::sync::Mutex; use crate::{engine::Engine, store::config::RiskConfig}; #[derive(Debug, Clone)] -pub struct OrderIds { +pub struct EngineOrder { + pub risk_equity: Decimal, + pub size: Decimal, + pub entry: Cloid, pub take_profit: Cloid, pub stop_loss: Cloid, - pub risk_equity: Decimal, -} - -impl OrderIds { - pub fn new(risk_equity: Decimal) -> Self { - Self { - entry: Cloid::random(), - take_profit: Cloid::random(), - stop_loss: Cloid::random(), - risk_equity, - } - } } pub struct RiskState { @@ -40,7 +31,7 @@ pub struct RiskState { pub struct RiskEngine { // Orders made by the engine - pub orders: Mutex>, + pub orders: Mutex>, pub handle: Mutex, pub state: Mutex, pub engine: Weak, @@ -154,18 +145,25 @@ impl RiskEngine { Ok(under_max_losses) } - pub async fn create_order(&self) -> tokio::io::Result { + pub async fn create_order(&self, signal: &Signal) -> tokio::io::Result { let state = self.state.lock().await; - Ok(OrderIds::new( - self.get_risk_config() - .await? - .risk_per_trade - .get(state.starting_equity), - )) + let risk = self + .get_risk_config() + .await? + .risk_per_trade + .get(state.starting_equity); + + Ok(EngineOrder { + entry: Cloid::random(), + take_profit: Cloid::random(), + stop_loss: Cloid::random(), + risk_equity: risk, + size: risk / (signal.entry_price - signal.stop_loss).abs(), + }) } - pub async fn order_placed(&self, order: OrderIds) { + pub async fn order_placed(&self, order: EngineOrder) { self.orders.lock().await.push(order); } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 17bd12e..13313c1 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -54,31 +54,22 @@ impl Formatted for EventLog { impl Formatted for SignalStatus { fn get_formatted(&self) -> Vec { match self { - Ok(signal) => { + Ok(signal) | Err(signal) => { vec![ - format!("\x1b[33mOK\x1b[0m"), - if matches!(signal.side, Side::Ask) { + if matches!(self, Ok(_)) { + format!("\x1b[34mOK\x1b[0m") + } else { + format!("\x1b[31mREJ\x1b[0m") + }, + if matches!(signal.side, Side::Bid) { format!("\x1b[32mBUY\x1b[0m") } else { format!("\x1b[31mSELL\x1b[0m") }, format_symbol(&signal.symbol), format_usd(signal.entry_price.as_f64()), - format!("\x1b[33mAPR {}\x1b[0m", signal.confidence), - ] - } - - Err(signal) => { - vec![ - format!("\x1b[31mERR\x1b[0m"), - if matches!(signal.side, Side::Ask) { - format!("\x1b[32mBUY\x1b[0m") - } else { - format!("\x1b[31mSELL\x1b[0m") - }, - format_symbol(&signal.symbol), - format_usd(signal.entry_price.as_f64()), - format!("\x1b[33mAPR {}\x1b[0m", signal.confidence), + format_usd(signal.take_profit.as_f64()), + format_usd_reverse(signal.stop_loss.as_f64()), ] } } @@ -274,6 +265,14 @@ pub fn format_usd(value: f64) -> String { } } +pub fn format_usd_reverse(value: f64) -> String { + if value.is_sign_positive() { + format!("\x1b[31m${}\x1b[0m", format_f64(value)) + } else { + format!("\x1b[32m${}\x1b[0m", format_f64(value)) + } +} + pub fn format_f64(value: f64) -> String { let abs = value.abs(); From f9f218ad802445320202d4f429ad3f203736b039 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 00:51:12 +0200 Subject: [PATCH 16/26] Signal execution --- src/engine/engine/execution.rs | 192 +++++++++++++++------------------ src/engine/engine/strategy.rs | 5 +- 2 files changed, 88 insertions(+), 109 deletions(-) diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index e8f3b51..db94da3 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -6,121 +6,103 @@ use pulse_sdk::prelude::*; use crate::engine::Engine; impl Engine { - pub async fn execute_signal(&self, signal: &Signal) -> tokio::io::Result<()> { + pub async fn execute_signal(&self, signal: &Signal) -> anyhow::Result<()> { let client = hypercore::mainnet(); let accounts = self.accounts.lock().await; - if let Some(acc) = accounts.get_active() { - let Some(asset_id) = self - .watch_list - .lock() - .await - .name_to_index - .get(&signal.symbol) - .cloned() - else { - self.terminal_server - .error( - "self::order", - &format!( - "Invalid Symbol: {:?}, Unable to get asset id", - signal.symbol - ), - ) - .await?; + let acc = accounts + .get_active() + .ok_or_else(|| anyhow::anyhow!("Unable to get active account"))?; - return Ok(()); - }; + let Some(asset_id) = self + .watch_list + .lock() + .await + .name_to_index + .get(&signal.symbol) + .cloned() + else { + return Err(anyhow::anyhow!( + "Invalid Symbol: {:?}, Unable to get asset id", + signal.symbol + )); + }; - let order_ids = self.risk_engine.create_order(signal).await?; + let engine_order = self.risk_engine.create_order(signal).await?; - println!("{order_ids:?}"); + let order = BatchOrder { + orders: vec![ + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.side, Side::Bid), + limit_px: signal.entry_price, + sz: engine_order.size, + reduce_only: false, + order_type: OrderTypePlacement::Limit { + tif: TimeInForce::Gtc, + }, + cloid: engine_order.entry, + }, + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.side, Side::Bid), + limit_px: signal.entry_price, + sz: engine_order.size, + reduce_only: true, + order_type: OrderTypePlacement::Trigger { + is_market: true, + trigger_px: signal.take_profit, + tpsl: hypercore::TpSl::Tp, + }, + cloid: engine_order.take_profit, + }, + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.side, Side::Bid), + limit_px: signal.entry_price, + sz: engine_order.size, + reduce_only: true, + order_type: OrderTypePlacement::Trigger { + is_market: true, + trigger_px: signal.stop_loss, + tpsl: hypercore::TpSl::Sl, + }, + cloid: engine_order.stop_loss, + }, + ], + grouping: hypercore::OrderGrouping::Na, + builder: None, + }; - // let order = BatchOrder { - // orders: vec![ - // OrderRequest { - // asset: asset_id, - // is_buy: matches!(signal.side, Side::Bid), - // limit_px: signal.entry_price, - // sz: 0.into(), - // reduce_only: false, - // order_type: OrderTypePlacement::Limit { - // tif: TimeInForce::Gtc, - // }, - // cloid: order_ids.entry, - // }, - // OrderRequest { - // asset: asset_id, - // is_buy: matches!(signal.side, Side::Bid), - // limit_px: signal.entry_price, - // sz: 0.into(), - // reduce_only: true, - // order_type: OrderTypePlacement::Trigger { - // is_market: true, - // trigger_px: signal.take_profit, - // tpsl: hypercore::TpSl::Tp, - // }, - // cloid: order_ids.take_profit, - // }, - // OrderRequest { - // asset: asset_id, - // is_buy: matches!(signal.side, Side::Bid), - // limit_px: signal.entry_price, - // sz: 0.into(), - // reduce_only: true, - // order_type: OrderTypePlacement::Trigger { - // is_market: true, - // trigger_px: signal.stop_loss, - // tpsl: hypercore::TpSl::Sl, - // }, - // cloid: order_ids.stop_loss, - // }, - // ], - // grouping: hypercore::OrderGrouping::Na, - // builder: None, - // }; + let nonce = chrono::Utc::now().timestamp_millis() as u64; - // let nonce = chrono::Utc::now().timestamp_millis() as u64; + match client + .place(&acc.private_key.0, order, nonce, None, None) + .await + { + Ok(o) if o.iter().any(|o| matches!(o, OrderResponseStatus::Error(_))) => { + self.risk_engine.order_placed(engine_order).await; + } - // match client - // .place(&acc.private_key.0, order, nonce, None, None) - // .await - // { - // Ok(o) if o.iter().any(|o| matches!(o, OrderResponseStatus::Error(_))) => { - // self.risk_engine.order_placed(order_ids).await; - // } + Ok(e) => { + return Err(anyhow::anyhow!( + "order rejected: {}", + e.into_iter() + .filter_map(|res| { + if let OrderResponseStatus::Error(e) = res { + Some(e) + } else { + None + } + }) + .collect::>() + .join(", ") + )); + } - // Ok(e) => { - // self.terminal_server - // .error( - // "self::order", - // &format!( - // "order rejected: {}", - // e.into_iter() - // .filter_map(|res| { - // if let OrderResponseStatus::Error(e) = res { - // Some(e) - // } else { - // None - // } - // }) - // .collect::>() - // .join(", ") - // ), - // ) - // .await?; - // } - - // Err(e) => { - // self.terminal_server - // .error("self::order", &e.to_string()) - // .await?; - // } - // } - } else { - self.terminal_server - .error("Engine::order", "Unable to get active account") - .await?; + e => { + e?; + } } Ok(()) diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index 88e02e1..552149a 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -169,10 +169,7 @@ impl StrategyEngine { engine .terminal_server - .error( - "engine::strategy", - &format!("Failed to execute signal: {e}"), - ) + .error("engine::order", &format!("Failed to execute signal: {e}")) .await? } } From 3d35cbf0a5dc9ccbb0c497bf1f8b92c137b68685 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 01:01:16 +0200 Subject: [PATCH 17/26] Usage cleanup and client initial signals --- src/engine/engine/risk.rs | 4 ++-- src/engine/engine/strategy.rs | 3 ++- src/engine/engine/terminal.rs | 20 +++++++++++++------- src/engine/store/config.rs | 4 ++-- src/engine/store/strategy.rs | 7 ++----- src/terminal/formatting.rs | 3 ++- src/terminal/main.rs | 4 ++-- src/terminal/terminal.rs | 1 + 8 files changed, 26 insertions(+), 20 deletions(-) diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index 1993247..7c42595 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -1,4 +1,4 @@ -use std::sync::{Arc, Weak}; +use pulse_sdk::prelude::*; use futures::StreamExt; use hypersdk::{ @@ -8,7 +8,7 @@ use hypersdk::{ ws::{ConnectionHandle, ConnectionStream, Event}, }, }; -use pulse_sdk::general::Signal; +use std::sync::{Arc, Weak}; use tokio::sync::Mutex; use crate::{engine::Engine, store::config::RiskConfig}; diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index 552149a..3d09a10 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -1,10 +1,11 @@ +use pulse_sdk::prelude::*; + use anyhow::Context; use futures::StreamExt; use hypersdk::hypercore::{ self, CandleInterval, Subscription, ws::{ConnectionHandle, ConnectionStream, Event}, }; -use pulse_sdk::prelude::*; use std::{ collections::HashSet, sync::{Arc, Weak}, diff --git a/src/engine/engine/terminal.rs b/src/engine/engine/terminal.rs index 20ab85b..781ac9b 100644 --- a/src/engine/engine/terminal.rs +++ b/src/engine/engine/terminal.rs @@ -1,5 +1,6 @@ -use crate::engine::Engine; use pulse_sdk::{map_postcard_err, prelude::*}; + +use crate::engine::Engine; use std::{ collections::HashMap, sync::{Arc, Weak}, @@ -37,7 +38,7 @@ impl TerminalServer { } pub async fn run_server(self: Arc) -> tokio::io::Result<()> { - let path = pulse_sdk::server_path(); + let path = server_path(); if path.exists() { tokio::fs::remove_file(&path).await?; @@ -103,7 +104,13 @@ impl TerminalServer { self.send_to( id, - pulse_sdk::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), + TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), + ) + .await?; + + self.send_to( + id, + TerminalServerMessage::SignalsUpdated(engine.signals.lock().await.clone()), ) .await?; @@ -151,7 +158,7 @@ impl TerminalServer { pub async fn broadcast( self: &Arc, - message: pulse_sdk::terminal::TerminalServerMessage, + message: TerminalServerMessage, ) -> tokio::io::Result<()> { let msg = map_postcard_err(postcard::to_allocvec(&message))?; @@ -263,7 +270,7 @@ impl TerminalServer { pub async fn send_to( self: &Arc, id: &usize, - message: pulse_sdk::terminal::TerminalServerMessage, + message: TerminalServerMessage, ) -> tokio::io::Result<()> { Self::send_to_client( self.clients.lock().await.get_mut(id).ok_or_else(|| { @@ -319,8 +326,7 @@ impl TerminalServer { self.logs.lock().await.push(log.clone()); - self.broadcast(pulse_sdk::terminal::TerminalServerMessage::AddLog(log)) - .await + self.broadcast(TerminalServerMessage::AddLog(log)).await } pub async fn info(self: &Arc, name: &str, message: &str) -> tokio::io::Result<()> { diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index a39bd9a..c9c7be7 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -1,7 +1,7 @@ -use std::collections::HashMap; - use pulse_sdk::general::Allocation; +use std::collections::HashMap; + #[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] pub struct RiskConfig { pub risk_per_trade: Allocation, diff --git a/src/engine/store/strategy.rs b/src/engine/store/strategy.rs index 858fd9a..bfcb8cd 100644 --- a/src/engine/store/strategy.rs +++ b/src/engine/store/strategy.rs @@ -1,9 +1,6 @@ -use std::{path::PathBuf, process::Stdio}; +use pulse_sdk::{map_postcard_err, prelude::*}; -use pulse_sdk::{ - map_postcard_err, - strategy::{StrategyEngineMessage, StrategyManifest, StrategyMessage}, -}; +use std::{path::PathBuf, process::Stdio}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, process::{Child, ChildStdout, Command}, diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 13313c1..93c1c87 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -1,6 +1,7 @@ -use hypersdk::hypercore::Side; use pulse_sdk::prelude::*; +use hypersdk::hypercore::Side; + pub trait Formatted { fn get_formatted(&self) -> Vec; } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 05aa89a..5b80956 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -2,6 +2,8 @@ pub mod command; pub mod formatting; pub mod terminal; +use pulse_sdk::prelude::*; + use std::any::Any; use chrono::{Local, Utc}; @@ -21,8 +23,6 @@ use pulse_ui::{ use crate::formatting::{Formatted, apply_padding}; -use pulse_sdk::prelude::*; - pub struct PulseTradeApp { sock: Option, diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index d6b7ec1..ed503b7 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -1,4 +1,5 @@ use pulse_sdk::{map_postcard_err, prelude::*}; + use pulse_ui::state::State; use tokio::{ From 11aff3bcb25d89548a555f75fcf7c6263a8b3559 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 01:45:25 +0200 Subject: [PATCH 18/26] Made confidence less confusing --- pulse-sdk/src/general.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 565523e..670ff9c 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -40,7 +40,8 @@ pub enum Allocation { pub struct Signal { pub symbol: String, pub side: Side, - pub confidence: f32, + /// 0-100 + pub confidence: u8, pub entry_price: Decimal, pub take_profit: Decimal, pub stop_loss: Decimal, From e61aa231ea90fa486e3fd12766e715acb7b4ce96 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 02:36:23 +0200 Subject: [PATCH 19/26] better commands --- src/engine/engine/command.rs | 158 ++++++++++++++++++++--------------- 1 file changed, 89 insertions(+), 69 deletions(-) diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 8fa4f9d..3d18823 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -1,7 +1,6 @@ use crate::{engine::Engine, store::config::ConfigManager}; use pulse_sdk::prelude::*; -use toml::Value; impl Engine { pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> { @@ -13,7 +12,43 @@ impl Engine { command: &str, mut args: Vec<&str>, ) -> tokio::io::Result<()> { + macro_rules! set_cfg { + ($n:ident) => { + let Ok(Ok($n)) = toml::value::Value::try_from(args[2]).map(|v| v.try_into()) else { + return self + .terminal_server + .error("config::set", "Unable to parse value") + .await; + }; + }; + } + match command { + "help" | "?" => { + self.terminal_server + .info("engine::help", "AVAILABLE COMMANDS") + .await?; + + let commands = [ + ("help", "Show this help message"), + ("config reload", "Reload configuration from disk"), + ("config save", "Save current configuration"), + ("account list", "List available accounts"), + ("account use ", "Switch active account"), + ("preset list", "List available presets"), + ("preset use ", "Switch active preset"), + ("strategy start", "Start strategy runtime"), + ("strategy set ", "Change active strategy"), + ("strategy ", "Send command to strategy runtime"), + ]; + + for (command, description) in commands { + self.terminal_server + .info("engine::help", &format!("{:<30} {}", command, description)) + .await?; + } + } + "config" | "cfg" => { if args.len() == 0 { return self.invalid_command_usage("config").await; @@ -21,10 +56,23 @@ impl Engine { match args[0] { "reload" => { - *self.config.lock().await = ConfigManager::new().await?; + let mut config = self.config.lock().await; + + let new_config = ConfigManager::new().await?; + let strategy_changed = + new_config.get_ref()?.strategy != config.get_ref()?.strategy; + + *config = new_config; + self.terminal_server .info("engine::config", "successfully reloaded") .await?; + + if strategy_changed { + self.strategy_engine + .reload(&config.get_ref()?.strategy) + .await?; + } } "save" => { @@ -34,72 +82,6 @@ impl Engine { .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.get_mut()?.watchlist = v); - - self.terminal_server - .info("config::set", "watchlist set successfully, use `config save` to persist changes") - .await?; - } - - "strategy" | "strat" | "sg" => { - set_cfg!(id, { - let id: String = id; - - if !crate::store::pulse_strategy(&id)? - .join("strategy.toml") - .exists() - { - return self - .terminal_server - .error( - "config::set[strategy]", - &format!("Non existent strategy `{id}`"), - ) - .await; - } - - self.strategy_engine.reload(id.as_str()).await?; - self.config.lock().await.get_mut()?.strategy = id; - }); - - self.terminal_server - .info("config::set", "strategy set successfully, use `config save` to persist changes") - .await?; - } - - _ => { - self.terminal_server - .error("config::set", "Invalid usage, available options: watchlist, strategy, cooldown") - .await?; - } - } - - self.terminal_server - .broadcast(TerminalServerMessage::ConfigUpdated( - self.get_config_status().await?, - )) - .await?; - } else { - self.invalid_command_usage("engine::config").await?; - } - } - _ => { self.invalid_command_usage("engine::config").await?; } @@ -254,6 +236,35 @@ impl Engine { match strategy_command { "start" => {} + "set" => { + set_cfg!(id); + + let id: String = id; + + if !crate::store::pulse_strategy(&id)? + .join("strategy.toml") + .exists() + { + return self + .terminal_server + .error( + "config::set[strategy]", + &format!("Non existent strategy `{id}`"), + ) + .await; + } + + self.strategy_engine.reload(id.as_str()).await?; + self.config.lock().await.get_mut()?.strategy = id; + + self.terminal_server + .info( + "config::set", + "strategy set successfully, use `config save` to persist changes", + ) + .await?; + } + _ => { self.strategy_engine .send(&StrategyEngineMessage::Command { @@ -267,11 +278,20 @@ impl Engine { _ => { self.terminal_server - .error("engine::cmd", &format!("command '{}' not found", command)) + .error( + "engine::cmd", + &format!("command '{}' not found, use help or ?", command), + ) .await?; } } + self.terminal_server + .broadcast(TerminalServerMessage::ConfigUpdated( + self.get_config_status().await?, + )) + .await?; + Ok(()) } } From 1a1b68ff4f20d472bbd7c702ca2865617b50f3a2 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 03:13:01 +0200 Subject: [PATCH 20/26] Fix order handling --- src/engine/engine/execution.rs | 13 +++++++------ src/engine/engine/strategy.rs | 14 ++++++++------ 2 files changed, 15 insertions(+), 12 deletions(-) diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index db94da3..2f2fbac 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -80,14 +80,11 @@ impl Engine { .place(&acc.private_key.0, order, nonce, None, None) .await { + // Invalid order Ok(o) if o.iter().any(|o| matches!(o, OrderResponseStatus::Error(_))) => { - self.risk_engine.order_placed(engine_order).await; - } - - Ok(e) => { return Err(anyhow::anyhow!( - "order rejected: {}", - e.into_iter() + "order rejected by HyperLiquid: {}", + o.into_iter() .filter_map(|res| { if let OrderResponseStatus::Error(e) = res { Some(e) @@ -100,6 +97,10 @@ impl Engine { )); } + Ok(_) => { + self.risk_engine.order_placed(engine_order).await; + } + e => { e?; } diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index 3d09a10..27a0355 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -163,17 +163,19 @@ impl StrategyEngine { continue; } - match engine.execute_signal(&signal).await { - Ok(_) => engine.signals.lock().await.push(Ok(signal)), + let result = match engine.execute_signal(&signal).await { + Ok(_) => Ok(signal), Err(e) => { - engine.signals.lock().await.push(Err(signal)); - engine .terminal_server .error("engine::order", &format!("Failed to execute signal: {e}")) - .await? + .await?; + + Err(signal) } - } + }; + + engine.signals.lock().await.push(result); engine .terminal_server From f2b31bbdbcb0e8680bece30f4aedc68ea5a7f67e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 03:28:04 +0200 Subject: [PATCH 21/26] Order size --- pulse-sdk/src/terminal.rs | 2 +- src/engine/engine/execution.rs | 15 ++++++++++++--- src/engine/engine/strategy.rs | 10 ++++++---- src/terminal/formatting.rs | 3 ++- 4 files changed, 21 insertions(+), 9 deletions(-) diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 35e923c..8a0151f 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -4,7 +4,7 @@ use crate::{ }; use hypersdk::Decimal; -pub type SignalStatus = Result; +pub type SignalStatus = Result<(Signal, Option), (Signal, Option)>; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum TerminalServerMessage { diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index 2f2fbac..01cd54d 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -1,12 +1,19 @@ -use hypersdk::hypercore::{ - self, BatchOrder, OrderRequest, OrderResponseStatus, OrderTypePlacement, Side, TimeInForce, +use hypersdk::{ + Decimal, + hypercore::{ + self, BatchOrder, OrderRequest, OrderResponseStatus, OrderTypePlacement, Side, TimeInForce, + }, }; use pulse_sdk::prelude::*; use crate::engine::Engine; impl Engine { - pub async fn execute_signal(&self, signal: &Signal) -> anyhow::Result<()> { + pub async fn execute_signal( + &self, + size: &mut Option, + signal: &Signal, + ) -> anyhow::Result<()> { let client = hypercore::mainnet(); let accounts = self.accounts.lock().await; @@ -30,6 +37,8 @@ impl Engine { let engine_order = self.risk_engine.create_order(signal).await?; + *size = Some(engine_order.size); + let order = BatchOrder { orders: vec![ OrderRequest { diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index 27a0355..d407b60 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -153,7 +153,7 @@ impl StrategyEngine { Some(StrategyMessage::Signal(mut signal)) => { if !engine.risk_engine.validate_signal(&mut signal).await? { - engine.signals.lock().await.push(Err(signal)); + engine.signals.lock().await.push(Err((signal, None))); engine .terminal_server @@ -163,15 +163,17 @@ impl StrategyEngine { continue; } - let result = match engine.execute_signal(&signal).await { - Ok(_) => Ok(signal), + let mut size = None; + + let result = match engine.execute_signal(&mut size, &signal).await { + Ok(_) => Ok((signal, size)), Err(e) => { engine .terminal_server .error("engine::order", &format!("Failed to execute signal: {e}")) .await?; - Err(signal) + Err((signal, size)) } }; diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 93c1c87..3ade1b2 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -55,7 +55,7 @@ impl Formatted for EventLog { impl Formatted for SignalStatus { fn get_formatted(&self) -> Vec { match self { - Ok(signal) | Err(signal) => { + Ok((signal, size)) | Err((signal, size)) => { vec![ if matches!(self, Ok(_)) { format!("\x1b[34mOK\x1b[0m") @@ -67,6 +67,7 @@ impl Formatted for SignalStatus { } else { format!("\x1b[31mSELL\x1b[0m") }, + size.map(|v| format_f64(v.as_f64())).unwrap_or_default(), format_symbol(&signal.symbol), format_usd(signal.entry_price.as_f64()), format_usd(signal.take_profit.as_f64()), From f2d6d9a8b9c7a3f9e9d2acafcfe9aa3fd817dbc8 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 04:23:53 +0200 Subject: [PATCH 22/26] Fix order serialization --- src/engine/engine/execution.rs | 7 ++++--- src/engine/engine/risk.rs | 2 +- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index 01cd54d..4be479d 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -1,7 +1,8 @@ use hypersdk::{ Decimal, hypercore::{ - self, BatchOrder, OrderRequest, OrderResponseStatus, OrderTypePlacement, Side, TimeInForce, + self, BatchOrder, NonceHandler, OrderRequest, OrderResponseStatus, OrderTypePlacement, + Side, TimeInForce, }, }; use pulse_sdk::prelude::*; @@ -83,10 +84,10 @@ impl Engine { builder: None, }; - let nonce = chrono::Utc::now().timestamp_millis() as u64; + let nonce = NonceHandler::default(); match client - .place(&acc.private_key.0, order, nonce, None, None) + .place(&acc.private_key.0, order, nonce.next(), None, None) .await { // Invalid order diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index 7c42595..fd3ce95 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -159,7 +159,7 @@ impl RiskEngine { take_profit: Cloid::random(), stop_loss: Cloid::random(), risk_equity: risk, - size: risk / (signal.entry_price - signal.stop_loss).abs(), + size: (risk / (signal.entry_price - signal.stop_loss).abs()).round_dp(3), }) } From c18dceb29167e4b7193d9c03d4b1ad170ae7317c Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 05:32:11 +0200 Subject: [PATCH 23/26] Fix risk --- pulse-sdk/src/general.rs | 15 ++++++++++++++- pulse-sdk/src/terminal.rs | 4 ++-- src/engine/engine/execution.rs | 28 +++++++++++++--------------- src/engine/engine/risk.rs | 34 ++++++++++++++++++++-------------- src/engine/store/config.rs | 2 ++ src/terminal/formatting.rs | 14 +++++++++----- 6 files changed, 60 insertions(+), 37 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 670ff9c..0de0dfe 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -1,4 +1,7 @@ -use hypersdk::{dec, hypercore::Side}; +use hypersdk::{ + dec, + hypercore::{Cloid, Side}, +}; use rust_decimal::Decimal; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -44,7 +47,17 @@ pub struct Signal { pub confidence: u8, pub entry_price: Decimal, pub take_profit: Decimal, +} + +#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] +pub struct EngineOrder { + pub risk_equity: Decimal, + pub size: Decimal, pub stop_loss: Decimal, + + pub entry: Cloid, + pub take_profit: Cloid, + pub stop_loss_id: Cloid, } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 8a0151f..20cbe6d 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,10 +1,10 @@ use crate::{ - general::{EngineStatus, EventLog, Signal}, + general::{EngineOrder, EngineStatus, EventLog, Signal}, strategy::StrategyManifest, }; use hypersdk::Decimal; -pub type SignalStatus = Result<(Signal, Option), (Signal, Option)>; +pub type SignalStatus = Result<(Signal, Option), (Signal, Option)>; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum TerminalServerMessage { diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index 4be479d..69e4569 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -1,9 +1,6 @@ -use hypersdk::{ - Decimal, - hypercore::{ - self, BatchOrder, NonceHandler, OrderRequest, OrderResponseStatus, OrderTypePlacement, - Side, TimeInForce, - }, +use hypersdk::hypercore::{ + self, BatchOrder, NonceHandler, OrderRequest, OrderResponseStatus, OrderTypePlacement, Side, + TimeInForce, }; use pulse_sdk::prelude::*; @@ -12,7 +9,7 @@ use crate::engine::Engine; impl Engine { pub async fn execute_signal( &self, - size: &mut Option, + size: &mut Option, signal: &Signal, ) -> anyhow::Result<()> { let client = hypercore::mainnet(); @@ -38,7 +35,7 @@ impl Engine { let engine_order = self.risk_engine.create_order(signal).await?; - *size = Some(engine_order.size); + *size = Some(engine_order); let order = BatchOrder { orders: vec![ @@ -55,7 +52,7 @@ impl Engine { }, OrderRequest { asset: asset_id, - is_buy: matches!(signal.side, Side::Bid), + is_buy: !matches!(signal.side, Side::Bid), limit_px: signal.entry_price, sz: engine_order.size, reduce_only: true, @@ -68,19 +65,19 @@ impl Engine { }, OrderRequest { asset: asset_id, - is_buy: matches!(signal.side, Side::Bid), + is_buy: !matches!(signal.side, Side::Bid), limit_px: signal.entry_price, sz: engine_order.size, reduce_only: true, order_type: OrderTypePlacement::Trigger { is_market: true, - trigger_px: signal.stop_loss, + trigger_px: engine_order.stop_loss, tpsl: hypercore::TpSl::Sl, }, - cloid: engine_order.stop_loss, + cloid: engine_order.stop_loss_id, }, ], - grouping: hypercore::OrderGrouping::Na, + grouping: hypercore::OrderGrouping::NormalTpsl, builder: None, }; @@ -95,9 +92,10 @@ impl Engine { return Err(anyhow::anyhow!( "order rejected by HyperLiquid: {}", o.into_iter() - .filter_map(|res| { + .enumerate() + .filter_map(|(i, res)| { if let OrderResponseStatus::Error(e) = res { - Some(e) + Some(format!("{i}={e}")) } else { None } diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index fd3ce95..d217d16 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -4,7 +4,7 @@ use futures::StreamExt; use hypersdk::{ Decimal, hypercore::{ - self, Cloid, + self, Cloid, Side, ws::{ConnectionHandle, ConnectionStream, Event}, }, }; @@ -13,16 +13,6 @@ use tokio::sync::Mutex; use crate::{engine::Engine, store::config::RiskConfig}; -#[derive(Debug, Clone)] -pub struct EngineOrder { - pub risk_equity: Decimal, - pub size: Decimal, - - pub entry: Cloid, - pub take_profit: Cloid, - pub stop_loss: Cloid, -} - pub struct RiskState { pub starting_equity: Decimal, pub pnl: Decimal, @@ -85,7 +75,7 @@ impl RiskEngine { let mut rm = Vec::new(); for (i, order_ids) in orders.iter().enumerate() { - if order_ids.stop_loss == cloid { + if order_ids.stop_loss_id == cloid { if order.status.is_filled() { rm.push(i); } @@ -154,12 +144,28 @@ impl RiskEngine { .risk_per_trade .get(state.starting_equity); + let size = self + .get_risk_config() + .await? + .size_per_trade + .get(state.starting_equity); + + let sl_distance = risk / (size / signal.entry_price); + Ok(EngineOrder { entry: Cloid::random(), take_profit: Cloid::random(), - stop_loss: Cloid::random(), + stop_loss_id: Cloid::random(), + + size: size.round_dp(3), risk_equity: risk, - size: (risk / (signal.entry_price - signal.stop_loss).abs()).round_dp(3), + + stop_loss: if matches!(signal.side, Side::Ask) { + signal.entry_price + sl_distance + } else { + signal.entry_price - sl_distance + } + .round_dp(3), }) } diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index c9c7be7..cc9d21b 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -4,6 +4,7 @@ use std::collections::HashMap; #[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] pub struct RiskConfig { + pub size_per_trade: Allocation, pub risk_per_trade: Allocation, pub max_open_positions: u32, pub max_daily_loss: Allocation, @@ -28,6 +29,7 @@ pub struct ConfigManager { impl Default for RiskConfig { fn default() -> Self { Self { + size_per_trade: Allocation::Percent(20.into()), risk_per_trade: Allocation::Percent(10.into()), max_daily_loss: Allocation::Percent(20.into()), max_open_positions: 1, diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 3ade1b2..de475d5 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -55,8 +55,8 @@ impl Formatted for EventLog { impl Formatted for SignalStatus { fn get_formatted(&self) -> Vec { match self { - Ok((signal, size)) | Err((signal, size)) => { - vec![ + Ok((signal, order)) | Err((signal, order)) => { + let mut base = vec![ if matches!(self, Ok(_)) { format!("\x1b[34mOK\x1b[0m") } else { @@ -67,12 +67,16 @@ impl Formatted for SignalStatus { } else { format!("\x1b[31mSELL\x1b[0m") }, - size.map(|v| format_f64(v.as_f64())).unwrap_or_default(), format_symbol(&signal.symbol), format_usd(signal.entry_price.as_f64()), format_usd(signal.take_profit.as_f64()), - format_usd_reverse(signal.stop_loss.as_f64()), - ] + ]; + + if let Some(order) = order { + base.extend(vec![format_usd_reverse(order.stop_loss.as_f64())]); + } + + base } } } From 1cba37296d970095c1504d43e67d5dbdf98df5dd Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 06:05:15 +0200 Subject: [PATCH 24/26] Working order --- pulse-sdk/src/general.rs | 2 +- pulse-sdk/src/lib.rs | 19 ++++++++++++++++++- src/engine/engine/execution.rs | 21 +++++++++++++-------- src/engine/engine/risk.rs | 9 ++++----- 4 files changed, 36 insertions(+), 15 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 0de0dfe..ec178ec 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -51,7 +51,7 @@ pub struct Signal { #[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] pub struct EngineOrder { - pub risk_equity: Decimal, + pub risk: Decimal, pub size: Decimal, pub stop_loss: Decimal, diff --git a/pulse-sdk/src/lib.rs b/pulse-sdk/src/lib.rs index af0744e..340dfb9 100644 --- a/pulse-sdk/src/lib.rs +++ b/pulse-sdk/src/lib.rs @@ -3,6 +3,8 @@ pub mod strategy; pub mod terminal; pub use hypersdk; +use hypersdk::Decimal; +use rust_decimal::RoundingStrategy; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use crate::{general::LogKind, strategy::StrategyEngineMessage}; @@ -10,9 +12,9 @@ use crate::{general::LogKind, strategy::StrategyEngineMessage}; pub mod prelude { pub use crate::Strategy; pub use crate::general::*; - pub use crate::server_path; pub use crate::strategy::*; pub use crate::terminal::*; + pub use crate::{round_price, server_path}; pub use hypersdk; pub use postcard; @@ -26,6 +28,21 @@ pub fn map_postcard_err(res: postcard::Result) -> tokio::io::Result { res.map_err(|e| tokio::io::Error::new(std::io::ErrorKind::Other, e)) } +pub fn round_price(price: Decimal) -> Decimal { + let digits = price.trunc().to_string().len() as u32; + let decimal_places = 5_i32 - digits as i32; + + if decimal_places < 0 { + let factor = Decimal::from(10_u64.pow((-decimal_places) as u32)); + (price / factor).round() * factor + } else { + price.round_dp_with_strategy( + decimal_places as u32, + RoundingStrategy::MidpointAwayFromZero, + ) + } +} + pub async fn send_raw(data: &[u8]) -> tokio::io::Result<()> { let mut stdout = tokio::io::stdout(); diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs index 69e4569..96279ec 100644 --- a/src/engine/engine/execution.rs +++ b/src/engine/engine/execution.rs @@ -37,13 +37,16 @@ impl Engine { *size = Some(engine_order); + let limit_px = round_price(signal.entry_price); + let sz = round_price(engine_order.size / signal.entry_price); + let order = BatchOrder { orders: vec![ OrderRequest { asset: asset_id, is_buy: matches!(signal.side, Side::Bid), - limit_px: signal.entry_price, - sz: engine_order.size, + limit_px, + sz, reduce_only: false, order_type: OrderTypePlacement::Limit { tif: TimeInForce::Gtc, @@ -53,12 +56,12 @@ impl Engine { OrderRequest { asset: asset_id, is_buy: !matches!(signal.side, Side::Bid), - limit_px: signal.entry_price, - sz: engine_order.size, + limit_px, + sz, reduce_only: true, order_type: OrderTypePlacement::Trigger { is_market: true, - trigger_px: signal.take_profit, + trigger_px: round_price(signal.take_profit), tpsl: hypercore::TpSl::Tp, }, cloid: engine_order.take_profit, @@ -66,12 +69,12 @@ impl Engine { OrderRequest { asset: asset_id, is_buy: !matches!(signal.side, Side::Bid), - limit_px: signal.entry_price, - sz: engine_order.size, + limit_px, + sz, reduce_only: true, order_type: OrderTypePlacement::Trigger { is_market: true, - trigger_px: engine_order.stop_loss, + trigger_px: round_price(engine_order.stop_loss), tpsl: hypercore::TpSl::Sl, }, cloid: engine_order.stop_loss_id, @@ -81,6 +84,8 @@ impl Engine { builder: None, }; + println!("{order:?}"); + let nonce = NonceHandler::default(); match client diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index d217d16..07e4ff8 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -87,7 +87,7 @@ impl RiskEngine { for r in rm { let order = self.orders.lock().await.remove(r); - self.state.lock().await.pnl -= order.risk_equity; + self.state.lock().await.pnl -= order.risk; } } } @@ -157,15 +157,14 @@ impl RiskEngine { take_profit: Cloid::random(), stop_loss_id: Cloid::random(), - size: size.round_dp(3), - risk_equity: risk, + size, + risk, stop_loss: if matches!(signal.side, Side::Ask) { signal.entry_price + sl_distance } else { signal.entry_price - sl_distance - } - .round_dp(3), + }, }) } From c7fb587cbb394b0a0e63a889b4d3fab9375257a8 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 06:15:52 +0200 Subject: [PATCH 25/26] Removed max positions --- pulse-sdk/src/terminal.rs | 1 - src/engine/engine/mod.rs | 1 - src/engine/engine/risk.rs | 2 -- src/engine/store/config.rs | 2 -- src/terminal/formatting.rs | 8 ++------ 5 files changed, 2 insertions(+), 12 deletions(-) diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 20cbe6d..0063acb 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -65,6 +65,5 @@ pub struct EngineConfig { pub description: Option, pub risk_per_trade: String, - pub max_open_positions: u32, pub max_daily_loss: String, } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 141682d..fe2fd7f 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -148,7 +148,6 @@ impl Engine { risk_per_trade: config.risk.risk_per_trade.to_string(), max_daily_loss: config.risk.max_daily_loss.to_string(), - max_open_positions: config.risk.max_open_positions, }) } diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index 07e4ff8..28569fc 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -16,7 +16,6 @@ use crate::{engine::Engine, store::config::RiskConfig}; pub struct RiskState { pub starting_equity: Decimal, pub pnl: Decimal, - pub open_positions: usize, } pub struct RiskEngine { @@ -35,7 +34,6 @@ impl RiskEngine { state: Mutex::new(RiskState { starting_equity: 0.into(), pnl: 0.into(), - open_positions: 0, }), engine, }) diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index cc9d21b..d3c6f03 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -6,7 +6,6 @@ use std::collections::HashMap; pub struct RiskConfig { pub size_per_trade: Allocation, pub risk_per_trade: Allocation, - pub max_open_positions: u32, pub max_daily_loss: Allocation, } @@ -32,7 +31,6 @@ impl Default for RiskConfig { size_per_trade: Allocation::Percent(20.into()), risk_per_trade: Allocation::Percent(10.into()), max_daily_loss: Allocation::Percent(20.into()), - max_open_positions: 1, } } } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index de475d5..4b3db06 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -229,14 +229,10 @@ impl Formatted for EngineConfig { Triple("", "", ""), Triple( "\x1b[2mRisk/Trade\x1b[0m", - "\x1b[2mMax Positions\x1b[0m", + "\x1b[2m..\x1b[0m", "\x1b[2mMax daily loss\x1b[0m", ), - Triple( - &self.risk_per_trade, - &self.max_open_positions.to_string(), - &self.max_daily_loss, - ), + Triple(&self.risk_per_trade, "..", &self.max_daily_loss), ] .get_formatted() } From 6dfd142d3c5847ca4bec0c41604b0d2a4c7e439f Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 2 Aug 2026 07:07:13 +0200 Subject: [PATCH 26/26] Restart command --- src/engine/engine/command.rs | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 3d18823..07bf9e2 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -38,6 +38,7 @@ impl Engine { ("preset list", "List available presets"), ("preset use ", "Switch active preset"), ("strategy start", "Start strategy runtime"), + ("strategy restart", "Restart strategy runtime"), ("strategy set ", "Change active strategy"), ("strategy ", "Send command to strategy runtime"), ]; @@ -236,6 +237,12 @@ impl Engine { match strategy_command { "start" => {} + "restart" => { + self.strategy_engine + .reload(&self.config.lock().await.get_ref()?.strategy) + .await?; + } + "set" => { set_cfg!(id);