From 9c29695728fb1b3157e68b4936b49b0d5bc28c48 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 01:35:36 +0200 Subject: [PATCH 01/12] Strategy and Risk sdk --- pulse-sdk/Cargo.toml | 2 +- pulse-sdk/src/lib.rs | 137 ++++++++++++++++++++++++++++++++++-- pulse-sdk/src/plugin.rs | 21 +++--- src/engine/engine/plugin.rs | 3 +- 4 files changed, 146 insertions(+), 17 deletions(-) diff --git a/pulse-sdk/Cargo.toml b/pulse-sdk/Cargo.toml index bb06f35..f8b9607 100644 --- a/pulse-sdk/Cargo.toml +++ b/pulse-sdk/Cargo.toml @@ -7,4 +7,4 @@ edition = "2024" serde = { workspace = true } hypersdk = { workspace = true } postcard = { workspace = true } -tokio = { workspace = true } +tokio = { workspace = true, features = ["io-std"]} diff --git a/pulse-sdk/src/lib.rs b/pulse-sdk/src/lib.rs index dd2211e..189435a 100644 --- a/pulse-sdk/src/lib.rs +++ b/pulse-sdk/src/lib.rs @@ -1,25 +1,152 @@ -#[cfg(target_os = "macos")] -use std::path::PathBuf; - pub mod general; pub mod plugin; pub mod terminal; pub mod units; pub use hypersdk; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; + +use crate::plugin::{RiskEngineMessage, StrategyEngineMessage}; + pub mod prelude { pub use crate::general::*; pub use crate::plugin::*; pub use crate::server_path; pub use crate::terminal::*; pub use crate::units::*; + pub use hypersdk; + pub use postcard; } -pub fn server_path() -> PathBuf { - PathBuf::from("/tmp/pulse-engine.sock") +pub fn server_path() -> std::path::PathBuf { + std::path::PathBuf::from("/tmp/pulse-engine.sock") } 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 async fn send_raw(data: &[u8]) -> tokio::io::Result<()> { + let mut stdout = tokio::io::stdout(); + + stdout.write_all(&data.len().to_le_bytes()).await?; + stdout.write_all(data).await?; + stdout.flush().await?; + + Ok(()) +} + +macro_rules! engine_methods { + ($t:ty) => { + async fn start(&self) -> tokio::io::Result<()> { + let mut stdin = tokio::io::stdin(); + + loop { + let mut len_buf = [0u8; size_of::()]; + let size = stdin.read_exact(&mut len_buf).await?; + + let len = usize::from_le_bytes(len_buf); + + if size == 0 || len == 0 { + break Ok(()); + } + + let mut buffer = vec![0u8; len]; + + stdin.read_exact(&mut buffer).await?; + + self.on_raw(&buffer).await?; + } + } + + async fn send(&self, msg: &$t) -> tokio::io::Result<()> { + $crate::send_raw(&$crate::map_postcard_err( + $crate::prelude::postcard::to_allocvec(msg), + )?) + .await + } + }; +} + +#[allow(async_fn_in_trait)] +pub trait Strategy { + engine_methods!(prelude::StrategyMessage); + + async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> { + match map_postcard_err(postcard::from_bytes(&data))? { + StrategyEngineMessage::Initialize => self.initialize().await, + StrategyEngineMessage::Command { command, args } => self.command(command, args).await, + StrategyEngineMessage::WatchList(watchlist) => self.watchlist(watchlist).await, + StrategyEngineMessage::Incoming(incoming) => self.incoming(incoming).await, + StrategyEngineMessage::Candlestick { + symbol, + interval, + candles, + } => self.candlestick(symbol, interval, candles).await, + } + } + + async fn initialize(&self) -> tokio::io::Result<()> { + unimplemented!("Strategy::initialize") + } + + async fn command(&self, _command: String, _args: Vec) -> tokio::io::Result<()> { + unimplemented!("Strategy::command") + } + + async fn watchlist(&self, _watchlist: Vec) -> tokio::io::Result<()> { + unimplemented!("Strategy::watchlist") + } + + async fn incoming(&self, _incoming: hypersdk::hypercore::Incoming) -> tokio::io::Result<()> { + unimplemented!("Strategy::event") + } + + async fn candlestick( + &self, + _symbol: String, + _interval: hypersdk::hypercore::CandleInterval, + _candles: Vec, + ) -> tokio::io::Result<()> { + unimplemented!("Strategy::candlestick") + } +} + +#[allow(async_fn_in_trait)] +pub trait Risk { + engine_methods!(prelude::RiskMessage); + + async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> { + match map_postcard_err(postcard::from_bytes(&data))? { + RiskEngineMessage::Initialize => self.initialize().await, + RiskEngineMessage::Command { command, args } => self.command(command, args).await, + RiskEngineMessage::WatchList(watchlist) => self.watchlist(watchlist).await, + RiskEngineMessage::Signal(signal) => { + self.send(&prelude::RiskMessage::Signal( + self.signal_request(signal).await?, + )) + .await + } + } + } + + async fn initialize(&self) -> tokio::io::Result<()> { + unimplemented!("Risk::initialize") + } + + async fn command(&self, _command: String, _args: Vec) -> tokio::io::Result<()> { + unimplemented!("Risk::command") + } + + async fn watchlist(&self, _watchlist: Vec) -> tokio::io::Result<()> { + unimplemented!("Risk::watchlist") + } + + async fn signal_request( + &self, + _signal: prelude::StrategySignal, + ) -> tokio::io::Result { + unimplemented!("Risk::signal_request") + } +} diff --git a/pulse-sdk/src/plugin.rs b/pulse-sdk/src/plugin.rs index 821a226..fc54bf7 100644 --- a/pulse-sdk/src/plugin.rs +++ b/pulse-sdk/src/plugin.rs @@ -3,7 +3,7 @@ use crate::{ terminal::MarketItem, units::Direction, }; -use hypersdk::hypercore::{Candle, CandleInterval, Subscription}; +use hypersdk::hypercore::{Candle, CandleInterval, Incoming, Subscription}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct StrategyManifest { @@ -51,9 +51,7 @@ pub enum RiskMessage { GetWatchList, - Approve(Signal), - - Reject { reason: String }, + Signal(RiskSignal), } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -67,11 +65,7 @@ pub enum StrategyEngineMessage { args: Vec, }, - CandleUpdate { - symbol: String, - interval: CandleInterval, - candle: Candle, - }, + Incoming(Incoming), Candlestick { symbol: String, @@ -98,3 +92,12 @@ pub struct StrategySignal { pub confidence: f32, pub price: Option, } + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub enum RiskSignal { + Approve(Signal), + Reject { + rejection_confidence: f32, + reason: String, + }, +} diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 2c30d1b..5c69307 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -178,8 +178,7 @@ impl StrategyEngine { engine.terminal_server.log_raw(log).await?; } - Some(RiskMessage::Approve(signal)) => {} - Some(RiskMessage::Reject { reason }) => {} + Some(RiskMessage::Signal(signal)) => {} } } } From 2696618ca00782bb496eca84a0eebbfe9e50bd9a Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 01:48:28 +0200 Subject: [PATCH 02/12] More flexible signal handling --- pulse-sdk/src/lib.rs | 12 ++---------- pulse-sdk/src/plugin.rs | 1 + 2 files changed, 3 insertions(+), 10 deletions(-) diff --git a/pulse-sdk/src/lib.rs b/pulse-sdk/src/lib.rs index 189435a..5147e98 100644 --- a/pulse-sdk/src/lib.rs +++ b/pulse-sdk/src/lib.rs @@ -122,12 +122,7 @@ pub trait Risk { RiskEngineMessage::Initialize => self.initialize().await, RiskEngineMessage::Command { command, args } => self.command(command, args).await, RiskEngineMessage::WatchList(watchlist) => self.watchlist(watchlist).await, - RiskEngineMessage::Signal(signal) => { - self.send(&prelude::RiskMessage::Signal( - self.signal_request(signal).await?, - )) - .await - } + RiskEngineMessage::Signal(signal) => self.signal_request(signal).await, } } @@ -143,10 +138,7 @@ pub trait Risk { unimplemented!("Risk::watchlist") } - async fn signal_request( - &self, - _signal: prelude::StrategySignal, - ) -> tokio::io::Result { + async fn signal_request(&self, _signal: prelude::StrategySignal) -> tokio::io::Result<()> { unimplemented!("Risk::signal_request") } } diff --git a/pulse-sdk/src/plugin.rs b/pulse-sdk/src/plugin.rs index fc54bf7..e69c51c 100644 --- a/pulse-sdk/src/plugin.rs +++ b/pulse-sdk/src/plugin.rs @@ -97,6 +97,7 @@ pub struct StrategySignal { pub enum RiskSignal { Approve(Signal), Reject { + signal: StrategySignal, rejection_confidence: f32, reason: String, }, From e0230425ec922ce0457759ba73cba00d4e5541ac Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 02:28:54 +0200 Subject: [PATCH 03/12] Order execution and using decimal --- pulse-sdk/src/general.rs | 4 ++- pulse-sdk/src/terminal.rs | 4 +-- pulse-sdk/src/units.rs | 10 ++++--- src/engine/engine/mod.rs | 5 ++-- src/engine/engine/plugin.rs | 54 +++++++++++++++++++++++++++++++++++-- src/engine/fetch.rs | 10 ++++--- src/terminal/formatting.rs | 2 +- 7 files changed, 72 insertions(+), 17 deletions(-) diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index a2fffad..93e312f 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -1,3 +1,5 @@ +use hypersdk::Decimal; + use crate::units::{Direction, Symbol, USD}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -20,7 +22,7 @@ pub struct Signal { pub symbol: String, pub kind: Direction, pub confidence: f32, - pub size: f64, + pub size: Decimal, pub price: USD, pub take_profit: USD, pub stop_loss: USD, diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index a07e9f1..625af89 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -3,7 +3,7 @@ use crate::{ plugin::{RiskManifest, StrategyManifest}, units::{Symbol, USD, Volatility}, }; -use hypersdk::hypercore::CandleInterval; +use hypersdk::{Decimal, hypercore::CandleInterval}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum TerminalServerMessage { @@ -39,7 +39,7 @@ pub enum TerminalClientMessage { pub struct MarketItem { pub symbol: Symbol, pub price: USD, - pub trend: f64, + pub trend: Decimal, pub volume_24h: USD, } diff --git a/pulse-sdk/src/units.rs b/pulse-sdk/src/units.rs index e2758e9..9dd091f 100644 --- a/pulse-sdk/src/units.rs +++ b/pulse-sdk/src/units.rs @@ -1,8 +1,10 @@ +use hypersdk::Decimal; + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Symbol(pub String); #[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] -pub struct USD(pub f64); +pub struct USD(pub Decimal); #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum Direction { @@ -25,10 +27,10 @@ impl std::fmt::Display for Symbol { impl std::fmt::Display for USD { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - if self.0 > 0.0 { - write!(f, "\x1b[32m${}\x1b[0m", format_f64(self.0)) + if self.0.is_sign_positive() { + write!(f, "\x1b[32m${}\x1b[0m", format_f64(self.0.as_f64())) } else { - write!(f, "\x1b[31m${}\x1b[0m", format_f64(self.0)) + write!(f, "\x1b[31m${}\x1b[0m", format_f64(self.0.as_f64())) } } } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index bb1d46f..b8e4681 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -102,9 +102,8 @@ impl Engine { entry_price: USD(position .position .entry_px - .map(|px| px.as_f64()) - .unwrap_or(0.0)), - profit: USD(position.position.unrealized_pnl.as_f64()), + .unwrap_or_default()), + profit: USD(position.position.unrealized_pnl), }) .collect(), ), diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 5c69307..ada7363 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -1,4 +1,7 @@ -use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket}; +use hypersdk::hypercore::{ + self, BatchOrder, CandleInterval, OrderRequest, OrderTypePlacement, Subscription, TimeInForce, + WebSocket, +}; use pulse_sdk::prelude::*; use std::{ collections::HashSet, @@ -178,7 +181,54 @@ impl StrategyEngine { engine.terminal_server.log_raw(log).await?; } - Some(RiskMessage::Signal(signal)) => {} + Some(RiskMessage::Signal(RiskSignal::Approve(signal))) => { + let client = hypercore::mainnet(); + let accounts = engine.accounts.lock().await; + + if let Some(acc) = accounts.get_active() { + let order = BatchOrder { + orders: vec![OrderRequest { + asset: 0, + is_buy: matches!(signal.kind, Direction::Buy), + limit_px: signal.price.0, + sz: signal.size, + reduce_only: false, + order_type: OrderTypePlacement::Limit { + tif: TimeInForce::Gtc, + }, + cloid: Default::default(), + }], + grouping: hypercore::OrderGrouping::Na, + builder: None, + }; + + let nonce = chrono::Utc::now().timestamp_millis() as u64; + + match client + .place(&acc.private_key.0, order, nonce, None, None) + .await + { + Ok(_) => {} + Err(e) => { + engine + .terminal_server + .error("Engine::order", &e.to_string()) + .await?; + } + } + } else { + engine + .terminal_server + .error("Engine::order", "Unable to get active account") + .await?; + } + } + + Some(RiskMessage::Signal(RiskSignal::Reject { + signal, + rejection_confidence, + reason, + })) => {} } } } diff --git a/src/engine/fetch.rs b/src/engine/fetch.rs index d8be404..394618d 100644 --- a/src/engine/fetch.rs +++ b/src/engine/fetch.rs @@ -1,13 +1,14 @@ +use hypersdk::Decimal; use pulse_sdk::prelude::*; use serde_json::Value; use std::collections::HashMap; -fn number(value: &Value, field: &str) -> Result { +fn number(value: &Value, field: &str) -> Result { let raw = value[field] .as_str() .ok_or_else(|| format!("asset context field {field} must be a string"))?; - raw.parse::() + raw.parse::() .map_err(|error| format!("could not parse asset context field {field} ({raw}): {error}")) } @@ -56,7 +57,7 @@ pub async fn fetch_watch_list( let previous_day_price = number(context, "prevDayPx")?; let volume_24h = number(context, "dayNtlVlm")?; - if previous_day_price <= 0.0 { + if previous_day_price.is_zero() || previous_day_price.is_sign_negative() { return Err(format!("{symbol} has an invalid previous-day price")); } @@ -66,7 +67,8 @@ pub async fn fetch_watch_list( symbol: Symbol(symbol.to_owned()), price: USD(price), volume_24h: USD(volume_24h), - trend: ((price / previous_day_price) - 1.0) * 100.0, + trend: ((price / previous_day_price) - >::from(1)) + * >::from(100), }, ); } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 50211d6..58d4439 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -85,7 +85,7 @@ impl Formatted for MarketItem { } else { "\x1b[31m▼\x1b[0m" }, - format_f64(self.trend.abs()) + format_f64(self.trend.as_f64().abs()) ), ] } From a7df23bc429b2d9764e65fc31a90dc8fdd17c66b Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 03:03:30 +0200 Subject: [PATCH 04/12] Proper asset id --- src/engine/engine/mod.rs | 17 +++++++++--- src/engine/engine/plugin.rs | 28 +++++++++++++++++--- src/engine/fetch.rs | 53 ++++++++++++++++++++++--------------- 3 files changed, 70 insertions(+), 28 deletions(-) diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index b8e4681..48382f4 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -7,16 +7,22 @@ use crate::{ terminal::TerminalServer, }; use pulse_sdk::prelude::*; -use std::sync::Arc; +use std::{collections::HashMap, sync::Arc}; use tokio::{sync::Mutex, task::JoinHandle}; +#[derive(Clone)] +pub struct WatchList { + pub name_to_index: HashMap, + pub items: Vec, +} + #[derive(Clone)] pub struct Engine { pub terminal_server: Arc, pub strategy: Arc, pub config: Arc>, pub accounts: Arc>, - pub watch_list: Arc>>, + pub watch_list: Arc>, } impl Engine { @@ -33,7 +39,10 @@ impl Engine { strategy: strategy.initialize(engine.clone()), config, accounts, - watch_list: Arc::new(Mutex::new(Vec::new())), + watch_list: Arc::new(Mutex::new(WatchList { + name_to_index: HashMap::new(), + items: Vec::new(), + })), })) } @@ -61,7 +70,7 @@ impl Engine { .terminal_server .broadcast( pulse_sdk::terminal::TerminalServerMessage::WatchListUpdated( - watch_list, + watch_list.items, ), ) .await diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index ada7363..e99f983 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -79,7 +79,7 @@ impl StrategyEngine { Some(StrategyMessage::GetWatchList) => { self.strategy .send(&StrategyEngineMessage::WatchList( - engine.watch_list.lock().await.clone(), + engine.watch_list.lock().await.clone().items, )) .await?; } @@ -171,7 +171,7 @@ impl StrategyEngine { Some(RiskMessage::GetWatchList) => { self.risk .send(&&RiskEngineMessage::WatchList( - engine.watch_list.lock().await.clone(), + engine.watch_list.lock().await.clone().items, )) .await?; } @@ -186,9 +186,31 @@ impl StrategyEngine { let accounts = engine.accounts.lock().await; if let Some(acc) = accounts.get_active() { + let Some(asset_id) = engine + .watch_list + .lock() + .await + .name_to_index + .get(&signal.symbol) + .cloned() + else { + engine + .terminal_server + .error( + "Engine::order", + &format!( + "Invalid Symbol: {:?}, Unable to get asset id", + signal.symbol + ), + ) + .await?; + + continue; + }; + let order = BatchOrder { orders: vec![OrderRequest { - asset: 0, + asset: asset_id, is_buy: matches!(signal.kind, Direction::Buy), limit_px: signal.price.0, sz: signal.size, diff --git a/src/engine/fetch.rs b/src/engine/fetch.rs index 394618d..09bdf31 100644 --- a/src/engine/fetch.rs +++ b/src/engine/fetch.rs @@ -3,6 +3,8 @@ use pulse_sdk::prelude::*; use serde_json::Value; use std::collections::HashMap; +use crate::engine::WatchList; + fn number(value: &Value, field: &str) -> Result { let raw = value[field] .as_str() @@ -15,7 +17,7 @@ fn number(value: &Value, field: &str) -> Result { pub async fn fetch_watch_list( client: &hypersdk::hypercore::HttpClient, symbols: &[String], -) -> Result, String> { +) -> Result { let response = client .meta_and_asset_ctxs(None) .await @@ -35,6 +37,7 @@ pub async fn fetch_watch_list( let universe = response[0]["universe"] .as_array() .ok_or("metaAndAssetCtxs response is missing meta.universe")?; + let contexts = response[1] .as_array() .ok_or("metaAndAssetCtxs response contexts must be an array")?; @@ -47,12 +50,27 @@ pub async fn fetch_watch_list( )); } - let mut by_symbol = HashMap::with_capacity(universe.len()); + // Build symbol -> perp asset index + let mut name_to_index = HashMap::with_capacity(universe.len()); - for (meta, context) in universe.iter().zip(contexts) { + for (index, meta) in universe.iter().enumerate() { let symbol = meta["name"] .as_str() .ok_or("instrument metadata is missing a name")?; + + name_to_index.insert(symbol.to_owned(), index); + } + + // Now only process requested symbols + let mut items = Vec::with_capacity(symbols.len()); + + for symbol in symbols { + let index = *name_to_index + .get(symbol) + .ok_or_else(|| format!("{symbol} is not in the Hyperliquid perpetual universe"))?; + + let context = &contexts[index]; + let price = number(context, "markPx")?; let previous_day_price = number(context, "prevDayPx")?; let volume_24h = number(context, "dayNtlVlm")?; @@ -61,24 +79,17 @@ pub async fn fetch_watch_list( return Err(format!("{symbol} has an invalid previous-day price")); } - by_symbol.insert( - symbol, - MarketItem { - symbol: Symbol(symbol.to_owned()), - price: USD(price), - volume_24h: USD(volume_24h), - trend: ((price / previous_day_price) - >::from(1)) - * >::from(100), - }, - ); + items.push(MarketItem { + symbol: Symbol(symbol.clone()), + price: USD(price), + volume_24h: USD(volume_24h), + trend: ((price / previous_day_price) - >::from(1)) + * >::from(100), + }); } - symbols - .iter() - .map(|symbol| { - by_symbol - .remove(symbol.as_str()) - .ok_or_else(|| format!("{symbol} is not in the Hyperliquid perpetual universe")) - }) - .collect() + Ok(WatchList { + items, + name_to_index, + }) } From b8ea0c46148cdcdf9d83cbc73f19841fa2b4596d Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 03:18:08 +0200 Subject: [PATCH 05/12] Order making --- src/engine/engine/plugin.rs | 48 +++++++++++++++++++++++++++++-------- 1 file changed, 38 insertions(+), 10 deletions(-) diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index e99f983..dd8574f 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -209,17 +209,45 @@ impl StrategyEngine { }; let order = BatchOrder { - orders: vec![OrderRequest { - asset: asset_id, - is_buy: matches!(signal.kind, Direction::Buy), - limit_px: signal.price.0, - sz: signal.size, - reduce_only: false, - order_type: OrderTypePlacement::Limit { - tif: TimeInForce::Gtc, + orders: vec![ + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.kind, Direction::Buy), + limit_px: signal.price.0, + sz: signal.size, + reduce_only: false, + order_type: OrderTypePlacement::Limit { + tif: TimeInForce::Gtc, + }, + cloid: Default::default(), }, - cloid: Default::default(), - }], + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.kind, Direction::Buy), + limit_px: signal.price.0, + sz: signal.size, + reduce_only: true, + order_type: OrderTypePlacement::Trigger { + is_market: true, + trigger_px: signal.take_profit.0, + tpsl: hypercore::TpSl::Tp, + }, + cloid: Default::default(), + }, + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.kind, Direction::Buy), + limit_px: signal.price.0, + sz: signal.size, + reduce_only: true, + order_type: OrderTypePlacement::Trigger { + is_market: true, + trigger_px: signal.stop_loss.0, + tpsl: hypercore::TpSl::Sl, + }, + cloid: Default::default(), + }, + ], grouping: hypercore::OrderGrouping::Na, builder: None, }; From 5f1b5c9aada1d3d44ce308eb9e151940790c4fcc Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 04:00:58 +0200 Subject: [PATCH 06/12] better Signals --- pulse-sdk/src/terminal.rs | 8 ++- src/engine/engine/execution.rs | 98 +++++++++++++++++++++++++++ src/engine/engine/mod.rs | 3 + src/engine/engine/plugin.rs | 120 +++++++-------------------------- src/terminal/formatting.rs | 74 +++++++++++++++++--- src/terminal/main.rs | 2 +- src/terminal/terminal.rs | 2 +- 7 files changed, 197 insertions(+), 110 deletions(-) create mode 100644 src/engine/engine/execution.rs diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 625af89..1998607 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,10 +1,12 @@ use crate::{ - general::{EventLog, MarketTrend, Position, Signal}, - plugin::{RiskManifest, StrategyManifest}, + general::{EventLog, MarketTrend, Position}, + plugin::{RiskManifest, RiskSignal, StrategyManifest}, units::{Symbol, USD, Volatility}, }; use hypersdk::{Decimal, hypercore::CandleInterval}; +pub type SignalStatus = Result; + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum TerminalServerMessage { // WatchList @@ -17,7 +19,7 @@ pub enum TerminalServerMessage { StrategyUpdated(Strategy), // Signals - SignalsUpdated(Vec), + SignalsUpdated(Vec), // Inspector Inspect(InspectTarget), diff --git a/src/engine/engine/execution.rs b/src/engine/engine/execution.rs new file mode 100644 index 0000000..9d56b4e --- /dev/null +++ b/src/engine/engine/execution.rs @@ -0,0 +1,98 @@ +use hypersdk::hypercore::{self, BatchOrder, OrderRequest, OrderTypePlacement, TimeInForce}; +use pulse_sdk::prelude::*; + +use crate::engine::Engine; + +impl Engine { + pub async fn execute_signal(&self, signal: &Signal) -> tokio::io::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?; + + return Ok(()); + }; + + let order = BatchOrder { + orders: vec![ + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.kind, Direction::Buy), + limit_px: signal.price.0, + sz: signal.size, + reduce_only: false, + order_type: OrderTypePlacement::Limit { + tif: TimeInForce::Gtc, + }, + cloid: Default::default(), + }, + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.kind, Direction::Buy), + limit_px: signal.price.0, + sz: signal.size, + reduce_only: true, + order_type: OrderTypePlacement::Trigger { + is_market: true, + trigger_px: signal.take_profit.0, + tpsl: hypercore::TpSl::Tp, + }, + cloid: Default::default(), + }, + OrderRequest { + asset: asset_id, + is_buy: matches!(signal.kind, Direction::Buy), + limit_px: signal.price.0, + sz: signal.size, + reduce_only: true, + order_type: OrderTypePlacement::Trigger { + is_market: true, + trigger_px: signal.stop_loss.0, + tpsl: hypercore::TpSl::Sl, + }, + cloid: Default::default(), + }, + ], + grouping: hypercore::OrderGrouping::Na, + builder: None, + }; + + let nonce = chrono::Utc::now().timestamp_millis() as u64; + + match client + .place(&acc.private_key.0, order, nonce, None, None) + .await + { + Ok(_) => {} + 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?; + } + + Ok(()) + } +} diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 48382f4..f728c80 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -1,5 +1,6 @@ pub mod command; pub mod plugin; +pub mod execution; use crate::{ engine::plugin::StrategyEngine, @@ -23,6 +24,7 @@ pub struct Engine { pub config: Arc>, pub accounts: Arc>, pub watch_list: Arc>, + pub signals: Arc>>, } impl Engine { @@ -43,6 +45,7 @@ impl Engine { name_to_index: HashMap::new(), items: Vec::new(), })), + signals: Arc::new(Mutex::new(Vec::new())), })) } diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index dd8574f..1e5dd28 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -1,7 +1,4 @@ -use hypersdk::hypercore::{ - self, BatchOrder, CandleInterval, OrderRequest, OrderTypePlacement, Subscription, TimeInForce, - WebSocket, -}; +use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket}; use pulse_sdk::prelude::*; use std::{ collections::HashSet, @@ -181,104 +178,37 @@ impl StrategyEngine { engine.terminal_server.log_raw(log).await?; } - Some(RiskMessage::Signal(RiskSignal::Approve(signal))) => { - let client = hypercore::mainnet(); - let accounts = engine.accounts.lock().await; + Some(RiskMessage::Signal(ref sig @ RiskSignal::Approve(ref signal))) => { + match engine.execute_signal(signal).await { + Ok(_) => engine.signals.lock().await.push(Ok(sig.clone())), + Err(e) => { + engine.signals.lock().await.push(Err(sig.clone())); - if let Some(acc) = accounts.get_active() { - let Some(asset_id) = engine - .watch_list - .lock() - .await - .name_to_index - .get(&signal.symbol) - .cloned() - else { engine .terminal_server - .error( - "Engine::order", - &format!( - "Invalid Symbol: {:?}, Unable to get asset id", - signal.symbol - ), - ) - .await?; - - continue; - }; - - let order = BatchOrder { - orders: vec![ - OrderRequest { - asset: asset_id, - is_buy: matches!(signal.kind, Direction::Buy), - limit_px: signal.price.0, - sz: signal.size, - reduce_only: false, - order_type: OrderTypePlacement::Limit { - tif: TimeInForce::Gtc, - }, - cloid: Default::default(), - }, - OrderRequest { - asset: asset_id, - is_buy: matches!(signal.kind, Direction::Buy), - limit_px: signal.price.0, - sz: signal.size, - reduce_only: true, - order_type: OrderTypePlacement::Trigger { - is_market: true, - trigger_px: signal.take_profit.0, - tpsl: hypercore::TpSl::Tp, - }, - cloid: Default::default(), - }, - OrderRequest { - asset: asset_id, - is_buy: matches!(signal.kind, Direction::Buy), - limit_px: signal.price.0, - sz: signal.size, - reduce_only: true, - order_type: OrderTypePlacement::Trigger { - is_market: true, - trigger_px: signal.stop_loss.0, - tpsl: hypercore::TpSl::Sl, - }, - cloid: Default::default(), - }, - ], - grouping: hypercore::OrderGrouping::Na, - builder: None, - }; - - let nonce = chrono::Utc::now().timestamp_millis() as u64; - - match client - .place(&acc.private_key.0, order, nonce, None, None) - .await - { - Ok(_) => {} - Err(e) => { - engine - .terminal_server - .error("Engine::order", &e.to_string()) - .await?; - } + .error("signal", &format!("Failed to execute signal: {e}")) + .await? } - } else { - engine - .terminal_server - .error("Engine::order", "Unable to get active account") - .await?; } + + engine + .terminal_server + .broadcast(TerminalServerMessage::SignalsUpdated( + engine.signals.lock().await.clone(), + )) + .await?; } - Some(RiskMessage::Signal(RiskSignal::Reject { - signal, - rejection_confidence, - reason, - })) => {} + Some(RiskMessage::Signal(reject)) => { + engine.signals.lock().await.push(Ok(reject)); + + engine + .terminal_server + .broadcast(TerminalServerMessage::SignalsUpdated( + engine.signals.lock().await.clone(), + )) + .await?; + } } } } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 58d4439..2a2744b 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -58,17 +58,71 @@ impl Formatted for EventLog { } } -impl Formatted for Signal { +impl Formatted for SignalStatus { fn get_formatted(&self) -> Vec { - vec![ - if matches!(self.kind, Direction::Buy) { - format!("\x1b[32m{}\x1b[0m", self.kind) - } else { - format!("\x1b[31m{}\x1b[0m", self.kind) - }, - format!("\x1b[35m{}\x1b[0m", self.symbol), - self.price.to_string(), - ] + match self { + Ok(RiskSignal::Approve(signal)) => { + vec![ + format!("\x1b[33mOK\x1b[0m"), + if matches!(signal.kind, Direction::Buy) { + format!("\x1b[32mBUY\x1b[0m") + } else { + format!("\x1b[31mSELL\x1b[0m") + }, + format!("\x1b[35m{}\x1b[0m", signal.symbol), + signal.price.to_string(), + format!("\x1b[33mAPR {}\x1b[0m", signal.confidence), + ] + } + Ok(RiskSignal::Reject { + signal, + rejection_confidence, + reason, + }) => { + vec![ + format!("\x1b[33mOK\x1b[0m"), + if matches!(signal.side, Direction::Buy) { + format!("\x1b[32mBUY\x1b[0m") + } else { + format!("\x1b[31mSELL\x1b[0m") + }, + format!("\x1b[35m{}\x1b[0m", signal.symbol), + format!("\x1b[31mREJ {}\x1b[0m", rejection_confidence), + reason.to_owned(), + ] + } + + Err(RiskSignal::Approve(signal)) => { + vec![ + format!("\x1b[31mERR\x1b[0m"), + if matches!(signal.kind, Direction::Buy) { + format!("\x1b[32mBUY\x1b[0m") + } else { + format!("\x1b[31mSELL\x1b[0m") + }, + format!("\x1b[35m{}\x1b[0m", signal.symbol), + signal.price.to_string(), + format!("\x1b[33mAPR {}\x1b[0m", signal.confidence), + ] + } + Err(RiskSignal::Reject { + signal, + rejection_confidence, + reason, + }) => { + vec![ + format!("\x1b[31mERR\x1b[0m"), + if matches!(signal.side, Direction::Buy) { + format!("\x1b[32mBUY\x1b[0m") + } else { + format!("\x1b[31mSELL\x1b[0m") + }, + format!("\x1b[35m{}\x1b[0m", signal.symbol), + format!("\x1b[31mREJ {}\x1b[0m", rejection_confidence), + reason.to_owned(), + ] + } + } } } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 25ae2f5..e751356 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -35,7 +35,7 @@ pub struct PulseTradeApp { watch_list: State>, active_positions: State>, logs: State>, - signals: State>, + signals: State>, inspect: State, } diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index 3657ce2..b032f9d 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -75,7 +75,7 @@ impl TerminalClient { watch_list: State>, active_positions: State>, logs: State>, - signals: State>, + signals: State>, market_overview: State>, status: State>, inspect: State, From 90d3c59ff486a0ba28dc7bc707ce7a674ae11009 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 04:30:26 +0200 Subject: [PATCH 07/12] Fix inconsistencies --- src/engine/engine/mod.rs | 42 +++++++++++++++++----------------------- src/engine/terminal.rs | 4 ++-- src/terminal/terminal.rs | 4 ++-- 3 files changed, 22 insertions(+), 28 deletions(-) diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index f728c80..c7d8c4c 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -1,6 +1,6 @@ pub mod command; -pub mod plugin; pub mod execution; +pub mod plugin; use crate::{ engine::plugin::StrategyEngine, @@ -11,7 +11,7 @@ use pulse_sdk::prelude::*; use std::{collections::HashMap, sync::Arc}; use tokio::{sync::Mutex, task::JoinHandle}; -#[derive(Clone)] +#[derive(Debug, Clone)] pub struct WatchList { pub name_to_index: HashMap, pub items: Vec, @@ -71,11 +71,7 @@ impl Engine { if let Err(error) = self .terminal_server - .broadcast( - pulse_sdk::terminal::TerminalServerMessage::WatchListUpdated( - watch_list.items, - ), - ) + .broadcast(TerminalServerMessage::WatchListUpdated(watch_list.items)) .await { self.terminal_server @@ -103,23 +99,21 @@ impl Engine { match client.clearinghouse_state(acc.address, None).await { Ok(state) => { self.terminal_server - .broadcast( - pulse_sdk::terminal::TerminalServerMessage::PositionsUpdated( - state - .asset_positions - .into_iter() - .map(|position| Position { - symbol: Symbol(position.position.coin), - size: position.position.szi.as_f64(), - entry_price: USD(position - .position - .entry_px - .unwrap_or_default()), - profit: USD(position.position.unrealized_pnl), - }) - .collect(), - ), - ) + .broadcast(TerminalServerMessage::PositionsUpdated( + state + .asset_positions + .into_iter() + .map(|position| Position { + symbol: Symbol(position.position.coin), + size: position.position.szi.as_f64(), + entry_price: USD(position + .position + .entry_px + .unwrap_or_default()), + profit: USD(position.position.unrealized_pnl), + }) + .collect(), + )) .await?; } Err(e) => { diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index ba14436..ec2e9f1 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -164,8 +164,8 @@ impl TerminalServer { } pub async fn send_to_client(client: &mut OwnedWriteHalf, msg: &[u8]) -> tokio::io::Result<()> { - client.write(&msg.len().to_le_bytes()).await?; - client.write(msg).await?; + client.write_all(&msg.len().to_le_bytes()).await?; + client.write_all(msg).await?; client.flush().await?; Ok(()) diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index b032f9d..6224913 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -31,8 +31,8 @@ impl TerminalClient { message: pulse_sdk::terminal::TerminalClientMessage, ) -> tokio::io::Result<()> { let msg = map_postcard_err(postcard::to_allocvec(&message))?; - self.writer.write(&msg.len().to_le_bytes()).await?; - self.writer.write(&msg).await?; + self.writer.write_all(&msg.len().to_le_bytes()).await?; + self.writer.write_all(&msg).await?; self.writer.flush().await?; Ok(()) From f3406c7b97775d06bd9a72508fe70294e465d845 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 06:57:25 +0200 Subject: [PATCH 08/12] Fix bug --- Cargo.lock | 2 + Cargo.toml | 2 + pulse-sdk/Cargo.toml | 3 +- pulse-sdk/src/general.rs | 2 +- pulse-sdk/src/units.rs | 2 +- src/engine/fetch.rs | 2 +- src/terminal/terminal.rs | 100 +++++++++++++++++++++------------------ 7 files changed, 62 insertions(+), 51 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index eac0541..2242481 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3947,6 +3947,7 @@ version = "0.1.0-alpha.0" dependencies = [ "hypersdk", "postcard", + "rust_decimal", "serde", "tokio", ] @@ -3963,6 +3964,7 @@ dependencies = [ "pulse-sdk", "pulse-ui", "rand 0.8.7", + "rust_decimal", "serde", "serde_json", "tokio", diff --git a/Cargo.toml b/Cargo.toml index 9badfe3..c1d3ab1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,6 +11,7 @@ hypersdk = { workspace = true } serde = { workspace = true } anyhow = { workspace = true } postcard = { workspace = true } +rust_decimal = { workspace = true } tokio = { workspace = true, features = [ "rt-multi-thread", "macros", @@ -29,6 +30,7 @@ toml = "1.1.3" members = ["pulse-ui", "pulse-sdk"] [workspace.dependencies] +rust_decimal = { version = "1.39", features = ["serde-str"] } postcard = { version = "1.1.3", features = ["alloc"] } pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" } pulse-sdk = { path = "pulse-sdk", version = "0.1.0-alpha.0" } diff --git a/pulse-sdk/Cargo.toml b/pulse-sdk/Cargo.toml index f8b9607..c545d4d 100644 --- a/pulse-sdk/Cargo.toml +++ b/pulse-sdk/Cargo.toml @@ -7,4 +7,5 @@ edition = "2024" serde = { workspace = true } hypersdk = { workspace = true } postcard = { workspace = true } -tokio = { workspace = true, features = ["io-std"]} +tokio = { workspace = true, features = ["io-std"] } +rust_decimal = { workspace = true } diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 93e312f..f36fc66 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -1,4 +1,4 @@ -use hypersdk::Decimal; +use rust_decimal::Decimal; use crate::units::{Direction, Symbol, USD}; diff --git a/pulse-sdk/src/units.rs b/pulse-sdk/src/units.rs index 9dd091f..d598482 100644 --- a/pulse-sdk/src/units.rs +++ b/pulse-sdk/src/units.rs @@ -1,4 +1,4 @@ -use hypersdk::Decimal; +use rust_decimal::Decimal; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Symbol(pub String); diff --git a/src/engine/fetch.rs b/src/engine/fetch.rs index 09bdf31..bbd6038 100644 --- a/src/engine/fetch.rs +++ b/src/engine/fetch.rs @@ -1,4 +1,4 @@ -use hypersdk::Decimal; +use rust_decimal::Decimal; use pulse_sdk::prelude::*; use serde_json::Value; use std::collections::HashMap; diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index 6224913..608d1b8 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -53,16 +53,22 @@ impl TerminalClient { let status = app.status.clone(); let inspect = app.inspect.clone(); - tokio::spawn(Self::run_client( - reader, - watch_list, - active_positions, - logs, - signals, - market_overview, - status, - inspect, - )); + tokio::spawn(async { + if let Err(v) = Self::run_client( + reader, + watch_list, + active_positions, + logs, + signals, + market_overview, + status, + inspect, + ) + .await + { + panic!("{v}") + } + }); app.sock = Some(self); @@ -80,55 +86,55 @@ impl TerminalClient { status: State>, inspect: State, ) -> tokio::io::Result<()> { - let mut len_buf = [0u8; size_of::()]; - reader - .read_exact(&mut len_buf) - .await - .expect("Failed to get header length"); + loop { + let mut len_buf = [0u8; size_of::()]; + reader + .read_exact(&mut len_buf) + .await + .expect("Failed to get header length"); - let len = usize::from_le_bytes(len_buf); + let len = usize::from_le_bytes(len_buf); - let mut buffer = vec![0u8; len]; + let mut buffer = vec![0u8; len]; - reader - .read_exact(&mut buffer) - .await - .expect("Failed to read socket"); + reader + .read_exact(&mut buffer) + .await + .expect("Failed to read socket"); - match map_postcard_err(postcard::from_bytes(&buffer))? { - TerminalServerMessage::WatchListUpdated(v) => { - *watch_list.lock().await = v; - } + match map_postcard_err(postcard::from_bytes(&buffer))? { + TerminalServerMessage::WatchListUpdated(v) => { + *watch_list.lock().await = v; + } - TerminalServerMessage::PositionsUpdated(v) => { - *active_positions.lock().await = v; - } + TerminalServerMessage::PositionsUpdated(v) => { + *active_positions.lock().await = v; + } - TerminalServerMessage::StrategyUpdated(v) => { - *market_overview.lock().await = Some(v); - } + TerminalServerMessage::StrategyUpdated(v) => { + *market_overview.lock().await = Some(v); + } - TerminalServerMessage::SignalsUpdated(v) => { - *signals.lock().await = v; - } + TerminalServerMessage::SignalsUpdated(v) => { + *signals.lock().await = v; + } - TerminalServerMessage::Inspect(v) => { - *inspect.lock().await = v; - } + TerminalServerMessage::Inspect(v) => { + *inspect.lock().await = v; + } - TerminalServerMessage::StatusUpdated(v) => { - *status.lock().await = Some(v); - } + TerminalServerMessage::StatusUpdated(v) => { + *status.lock().await = Some(v); + } - TerminalServerMessage::SetLogs(v) => { - *logs.lock().await = v; - } + TerminalServerMessage::SetLogs(v) => { + *logs.lock().await = v; + } - TerminalServerMessage::AddLog(v) => { - logs.lock().await.push(v); + TerminalServerMessage::AddLog(v) => { + logs.lock().await.push(v); + } } } - - Ok(()) } } From d0f683b4ca3789ca37f0a7b452a68a1140b76b58 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 30 Jul 2026 02:50:22 +0200 Subject: [PATCH 09/12] Removed risk --- pulse-sdk/src/lib.rs | 32 +--------- pulse-sdk/src/plugin.rs | 52 +--------------- pulse-sdk/src/terminal.rs | 7 +-- src/engine/engine/command.rs | 30 +--------- src/engine/engine/mod.rs | 2 +- src/engine/engine/plugin.rs | 111 ++++++----------------------------- src/engine/store/config.rs | 2 - src/engine/terminal.rs | 1 - src/terminal/formatting.rs | 46 ++------------- 9 files changed, 33 insertions(+), 250 deletions(-) diff --git a/pulse-sdk/src/lib.rs b/pulse-sdk/src/lib.rs index 5147e98..60177d6 100644 --- a/pulse-sdk/src/lib.rs +++ b/pulse-sdk/src/lib.rs @@ -6,7 +6,7 @@ pub use hypersdk; use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use crate::plugin::{RiskEngineMessage, StrategyEngineMessage}; +use crate::plugin::StrategyEngineMessage; pub mod prelude { pub use crate::general::*; @@ -112,33 +112,3 @@ pub trait Strategy { unimplemented!("Strategy::candlestick") } } - -#[allow(async_fn_in_trait)] -pub trait Risk { - engine_methods!(prelude::RiskMessage); - - async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> { - match map_postcard_err(postcard::from_bytes(&data))? { - RiskEngineMessage::Initialize => self.initialize().await, - RiskEngineMessage::Command { command, args } => self.command(command, args).await, - RiskEngineMessage::WatchList(watchlist) => self.watchlist(watchlist).await, - RiskEngineMessage::Signal(signal) => self.signal_request(signal).await, - } - } - - async fn initialize(&self) -> tokio::io::Result<()> { - unimplemented!("Risk::initialize") - } - - async fn command(&self, _command: String, _args: Vec) -> tokio::io::Result<()> { - unimplemented!("Risk::command") - } - - async fn watchlist(&self, _watchlist: Vec) -> tokio::io::Result<()> { - unimplemented!("Risk::watchlist") - } - - async fn signal_request(&self, _signal: prelude::StrategySignal) -> tokio::io::Result<()> { - unimplemented!("Risk::signal_request") - } -} diff --git a/pulse-sdk/src/plugin.rs b/pulse-sdk/src/plugin.rs index e69c51c..291a480 100644 --- a/pulse-sdk/src/plugin.rs +++ b/pulse-sdk/src/plugin.rs @@ -1,7 +1,6 @@ use crate::{ general::{EventLog, Signal}, terminal::MarketItem, - units::Direction, }; use hypersdk::hypercore::{Candle, CandleInterval, Incoming, Subscription}; @@ -13,17 +12,6 @@ pub struct StrategyManifest { pub version: String, } -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub struct RiskManifest { - pub name: String, - pub description: String, - pub author: String, - pub version: String, - - pub max_loss: u8, - pub cooldown: CandleInterval, -} - #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum StrategyMessage { Log(EventLog), @@ -42,16 +30,7 @@ pub enum StrategyMessage { UnsubscribeAll, - Signal(StrategySignal), -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub enum RiskMessage { - Log(EventLog), - - GetWatchList, - - Signal(RiskSignal), + Signal(Signal), } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -73,32 +52,3 @@ pub enum StrategyEngineMessage { candles: Vec, }, } - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub enum RiskEngineMessage { - Initialize, - - WatchList(Vec), - - Command { command: String, args: Vec }, - - Signal(StrategySignal), -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub struct StrategySignal { - pub symbol: String, - pub side: Direction, - pub confidence: f32, - pub price: Option, -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub enum RiskSignal { - Approve(Signal), - Reject { - signal: StrategySignal, - rejection_confidence: f32, - reason: String, - }, -} diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 1998607..033deae 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,11 +1,11 @@ use crate::{ - general::{EventLog, MarketTrend, Position}, - plugin::{RiskManifest, RiskSignal, StrategyManifest}, + general::{EventLog, MarketTrend, Position, Signal}, + plugin::StrategyManifest, units::{Symbol, USD, Volatility}, }; use hypersdk::{Decimal, hypercore::CandleInterval}; -pub type SignalStatus = Result; +pub type SignalStatus = Result; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum TerminalServerMessage { @@ -110,7 +110,6 @@ pub enum ItemState { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Strategy { pub strategy: StrategyManifest, - pub risk: RiskManifest, pub mode: Mode, pub state: ItemState, diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index bd39825..080a394 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -69,32 +69,6 @@ impl Engine { .await?; } - "risk" | "rs" => { - set_cfg!(id, { - let id: String = id; - - if !crate::store::pulse_plugin(&id)? - .join("risk.toml") - .exists() - { - return self - .terminal_server - .error( - "config::set::risk", - &format!("Non existent risk `{id}`"), - ) - .await; - } - - self.strategy.reload_risk(id.as_str()).await?; - self.config.lock().await.risk = id; - }); - - self.terminal_server - .info("config::set", "risk set successfully, use `config save` to persist changes") - .await?; - } - "cooldown" | "cool" | "cd" => { set_cfg!(cooldown, { self.config.lock().await.cooldown = cooldown; @@ -102,7 +76,7 @@ impl Engine { self.terminal_server.broadcast( pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { strategy: self.strategy.strategy.manifest.lock().await.clone(), - risk: self.strategy.risk.manifest.lock().await.clone(), + mode: Mode::Auto, state: ItemState::Running, cooldown, @@ -114,7 +88,7 @@ impl Engine { _ => { self.terminal_server - .error("config::set", "Invalid usage, available options: watchlist, strategy, risk, cooldown") + .error("config::set", "Invalid usage, available options: watchlist, strategy, cooldown") .await?; } } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index c7d8c4c..c003d59 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -31,7 +31,7 @@ impl Engine { pub async fn new() -> tokio::io::Result> { let config = Config::new().await?; - let strategy = StrategyEngine::new(&config.strategy, &config.risk).await?; + let strategy = StrategyEngine::new(&config.strategy).await?; let accounts = Arc::new(Mutex::new(AccountList::new().await?)); let config = Arc::new(Mutex::new(config)); diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 1e5dd28..2c52ba6 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -20,7 +20,6 @@ use crate::{ pub struct StrategyEngine { pub strategy: Arc>, - pub risk: Arc>, pub engine: Weak, pub ws: WebSocket, @@ -28,9 +27,8 @@ pub struct StrategyEngine { } impl StrategyEngine { - pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { + pub async fn new(strategy_id: &str) -> tokio::io::Result { let strategy = pulse_plugin(strategy_id)?; - let risk = pulse_plugin(risk_id)?; let (strategy, strategy_manifest) = get_manifest_plugin_pair( &strategy, @@ -38,15 +36,8 @@ impl StrategyEngine { &fs::read(strategy.join("strategy.toml")).await?, )?; - let (risk, risk_manifest) = get_manifest_plugin_pair( - &risk, - &risk.join("risk.bash"), - &fs::read(risk.join("risk.toml")).await?, - )?; - Ok(Self { strategy: Arc::new(Plugin::new(strategy, strategy_manifest)), - risk: Arc::new(Plugin::new(risk, risk_manifest)), engine: Weak::new(), ws: hypercore::mainnet_ws(), subscriptions: Mutex::new(HashSet::new()), @@ -87,7 +78,24 @@ impl StrategyEngine { } Some(StrategyMessage::Signal(signal)) => { - self.risk.send(&RiskEngineMessage::Signal(signal)).await?; + match engine.execute_signal(&signal).await { + Ok(_) => engine.signals.lock().await.push(Ok(signal)), + Err(e) => { + engine.signals.lock().await.push(Err(signal)); + + engine + .terminal_server + .error("signal", &format!("Failed to execute signal: {e}")) + .await? + } + } + + engine + .terminal_server + .broadcast(TerminalServerMessage::SignalsUpdated( + engine.signals.lock().await.clone(), + )) + .await?; } Some(StrategyMessage::Subscribe(subscription)) => { @@ -153,74 +161,10 @@ impl StrategyEngine { } } - pub async fn run_risk(&self) -> tokio::io::Result<()> { - let engine = self - .engine - .upgrade() - .expect("Failed to upgrade engine (StrategyEngine)"); - - self.risk.send(&RiskEngineMessage::Initialize).await?; - - loop { - match self.risk.recv().await? { - None => {} - - Some(RiskMessage::GetWatchList) => { - self.risk - .send(&&RiskEngineMessage::WatchList( - engine.watch_list.lock().await.clone().items, - )) - .await?; - } - - Some(RiskMessage::Log(mut log)) => { - log.name.insert_str(0, "risk::"); - engine.terminal_server.log_raw(log).await?; - } - - Some(RiskMessage::Signal(ref sig @ RiskSignal::Approve(ref signal))) => { - match engine.execute_signal(signal).await { - Ok(_) => engine.signals.lock().await.push(Ok(sig.clone())), - Err(e) => { - engine.signals.lock().await.push(Err(sig.clone())); - - engine - .terminal_server - .error("signal", &format!("Failed to execute signal: {e}")) - .await? - } - } - - engine - .terminal_server - .broadcast(TerminalServerMessage::SignalsUpdated( - engine.signals.lock().await.clone(), - )) - .await?; - } - - Some(RiskMessage::Signal(reject)) => { - engine.signals.lock().await.push(Ok(reject)); - - engine - .terminal_server - .broadcast(TerminalServerMessage::SignalsUpdated( - engine.signals.lock().await.clone(), - )) - .await?; - } - } - } - } - pub async fn spawn(self: &Arc) { let engine = self.clone(); tokio::spawn(async move { engine.run_strategy().await }); - - let engine = self.clone(); - - tokio::spawn(async move { engine.run_risk().await }); } pub async fn reload_strategy(self: &Arc, id: &str) -> tokio::io::Result<()> { @@ -239,23 +183,6 @@ impl StrategyEngine { Ok(()) } - - pub async fn reload_risk(self: &Arc, id: &str) -> tokio::io::Result<()> { - let plugin = pulse_plugin(id)?; - - let (child, manifest) = get_manifest_plugin_pair( - &plugin, - &plugin.join("risk.bash"), - &fs::read(plugin.join("risk.toml")).await?, - )?; - - self.risk.reload(child, manifest).await?; - - let engine = self.clone(); - tokio::spawn(async move { engine.run_risk().await }); - - Ok(()) - } } pub fn get_manifest_plugin_pair<'de, M: serde::Deserialize<'de>>( diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index c16d916..ba42aa2 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -4,7 +4,6 @@ use hypersdk::hypercore::CandleInterval; pub struct Config { pub watchlist: Vec, pub strategy: String, - pub risk: String, pub cooldown: CandleInterval, } @@ -13,7 +12,6 @@ impl Default for Config { Self { watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()], strategy: String::new(), - risk: String::new(), cooldown: CandleInterval::ThirtyMinutes, } } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index ec2e9f1..97e8150 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -66,7 +66,6 @@ impl TerminalServer { id, pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { strategy: engine.strategy.strategy.manifest.lock().await.clone(), - risk: engine.strategy.risk.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, cooldown: engine.config.lock().await.cooldown, diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 2a2744b..05296ef 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -61,7 +61,7 @@ impl Formatted for EventLog { impl Formatted for SignalStatus { fn get_formatted(&self) -> Vec { match self { - Ok(RiskSignal::Approve(signal)) => { + Ok(signal) => { vec![ format!("\x1b[33mOK\x1b[0m"), if matches!(signal.kind, Direction::Buy) { @@ -74,25 +74,8 @@ impl Formatted for SignalStatus { format!("\x1b[33mAPR {}\x1b[0m", signal.confidence), ] } - Ok(RiskSignal::Reject { - signal, - rejection_confidence, - reason, - }) => { - vec![ - format!("\x1b[33mOK\x1b[0m"), - if matches!(signal.side, Direction::Buy) { - format!("\x1b[32mBUY\x1b[0m") - } else { - format!("\x1b[31mSELL\x1b[0m") - }, - format!("\x1b[35m{}\x1b[0m", signal.symbol), - format!("\x1b[31mREJ {}\x1b[0m", rejection_confidence), - reason.to_owned(), - ] - } - Err(RiskSignal::Approve(signal)) => { + Err(signal) => { vec![ format!("\x1b[31mERR\x1b[0m"), if matches!(signal.kind, Direction::Buy) { @@ -105,23 +88,6 @@ impl Formatted for SignalStatus { format!("\x1b[33mAPR {}\x1b[0m", signal.confidence), ] } - Err(RiskSignal::Reject { - signal, - rejection_confidence, - reason, - }) => { - vec![ - format!("\x1b[31mERR\x1b[0m"), - if matches!(signal.side, Direction::Buy) { - format!("\x1b[32mBUY\x1b[0m") - } else { - format!("\x1b[31mSELL\x1b[0m") - }, - format!("\x1b[35m{}\x1b[0m", signal.symbol), - format!("\x1b[31mREJ {}\x1b[0m", rejection_confidence), - reason.to_owned(), - ] - } } } } @@ -291,7 +257,7 @@ impl Formatted for Strategy { ), Triple( &format!("\x1b[97m{}\x1b[0m", self.strategy.name), - &format!("\x1b[93m{}\x1b[0m", self.risk.name), + &format!("\x1b[93m{}\x1b[0m", ""), &format!("\x1b[90m{}\x1b[0m", self.strategy.version), ), Triple("", "", ""), @@ -303,7 +269,7 @@ impl Formatted for Strategy { Triple( &self.mode.to_string(), &self.state.to_string(), - &format!("\x1b[90m{}\x1b[0m", self.risk.version), + &format!("\x1b[90m{}\x1b[0m", ""), ), Triple("", "", ""), Triple( @@ -314,10 +280,10 @@ impl Formatted for Strategy { Triple( &format!( "\x1b[96m{}\x1b[0m (\x1b[90m{} rec\x1b[0m)", - self.cooldown, self.risk.cooldown + self.cooldown, "" ), &format!("\x1b[96m{:?}\x1b[0m", self.strategy.author), - &format!("\x1b[93m{}%\x1b[0m", self.risk.max_loss), + &format!("\x1b[93m{}%\x1b[0m", ""), ), ] .get_formatted() From 72a014586075bdcfb8710e3adce443aad021cac2 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 30 Jul 2026 03:22:59 +0200 Subject: [PATCH 10/12] Updated formatting --- src/terminal/formatting.rs | 31 +++---------------------------- 1 file changed, 3 insertions(+), 28 deletions(-) diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 05296ef..dace9dd 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -252,38 +252,13 @@ impl Formatted for Strategy { vec![ Triple( "\x1b[2mStrategy\x1b[0m", - "\x1b[2mRisk\x1b[0m", - "\x1b[2mStrat Ver\x1b[0m", + "\x1b[2mState\x1b[0m", + "\x1b[2mMode\x1b[0m", ), Triple( &format!("\x1b[97m{}\x1b[0m", self.strategy.name), - &format!("\x1b[93m{}\x1b[0m", ""), - &format!("\x1b[90m{}\x1b[0m", self.strategy.version), - ), - Triple("", "", ""), - Triple( - "\x1b[2mMode\x1b[0m", - "\x1b[2mState\x1b[0m", - "\x1b[2mRisk Ver\x1b[0m", - ), - Triple( - &self.mode.to_string(), &self.state.to_string(), - &format!("\x1b[90m{}\x1b[0m", ""), - ), - Triple("", "", ""), - Triple( - "\x1b[2mCooldown\x1b[0m", - "\x1b[2mStrat Author\x1b[0m", - "\x1b[2mMax loss\x1b[0m", - ), - Triple( - &format!( - "\x1b[96m{}\x1b[0m (\x1b[90m{} rec\x1b[0m)", - self.cooldown, "" - ), - &format!("\x1b[96m{:?}\x1b[0m", self.strategy.author), - &format!("\x1b[93m{}%\x1b[0m", ""), + &self.mode.to_string(), ), ] .get_formatted() From 660a0cb66ac84fd4c57e6cf572c7b2723d9a6636 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 30 Jul 2026 03:42:59 +0200 Subject: [PATCH 11/12] Refactored codebase to Strategy instead of plugin --- pulse-sdk/src/lib.rs | 6 ++-- pulse-sdk/src/{plugin.rs => strategy.rs} | 0 pulse-sdk/src/terminal.rs | 2 +- src/engine/engine/mod.rs | 4 +-- src/engine/engine/{plugin.rs => strategy.rs} | 6 ++-- src/engine/store/mod.rs | 2 +- src/engine/store/{plugin.rs => strategy.rs} | 30 ++++++++++---------- 7 files changed, 25 insertions(+), 25 deletions(-) rename pulse-sdk/src/{plugin.rs => strategy.rs} (100%) rename src/engine/engine/{plugin.rs => strategy.rs} (97%) rename src/engine/store/{plugin.rs => strategy.rs} (69%) diff --git a/pulse-sdk/src/lib.rs b/pulse-sdk/src/lib.rs index 60177d6..3ddcbb7 100644 --- a/pulse-sdk/src/lib.rs +++ b/pulse-sdk/src/lib.rs @@ -1,16 +1,16 @@ pub mod general; -pub mod plugin; +pub mod strategy; pub mod terminal; pub mod units; pub use hypersdk; use tokio::io::{AsyncReadExt, AsyncWriteExt}; -use crate::plugin::StrategyEngineMessage; +use crate::strategy::StrategyEngineMessage; pub mod prelude { pub use crate::general::*; - pub use crate::plugin::*; + pub use crate::strategy::*; pub use crate::server_path; pub use crate::terminal::*; pub use crate::units::*; diff --git a/pulse-sdk/src/plugin.rs b/pulse-sdk/src/strategy.rs similarity index 100% rename from pulse-sdk/src/plugin.rs rename to pulse-sdk/src/strategy.rs diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 033deae..efac397 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,6 +1,6 @@ use crate::{ general::{EventLog, MarketTrend, Position, Signal}, - plugin::StrategyManifest, + strategy::StrategyManifest, units::{Symbol, USD, Volatility}, }; use hypersdk::{Decimal, hypercore::CandleInterval}; diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index c003d59..f6cb30b 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -1,9 +1,9 @@ pub mod command; pub mod execution; -pub mod plugin; +pub mod strategy; use crate::{ - engine::plugin::StrategyEngine, + engine::strategy::StrategyEngine, store::{accounts::AccountList, config::Config}, terminal::TerminalServer, }; diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/strategy.rs similarity index 97% rename from src/engine/engine/plugin.rs rename to src/engine/engine/strategy.rs index 2c52ba6..0ab4b0f 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/strategy.rs @@ -15,11 +15,11 @@ use tokio::{ use crate::{ engine::Engine, - store::{plugin::Plugin, pulse_plugin}, + store::{strategy::StrategyChild, pulse_plugin}, }; pub struct StrategyEngine { - pub strategy: Arc>, + pub strategy: Arc, pub engine: Weak, pub ws: WebSocket, @@ -37,7 +37,7 @@ impl StrategyEngine { )?; Ok(Self { - strategy: Arc::new(Plugin::new(strategy, strategy_manifest)), + strategy: Arc::new(StrategyChild::new(strategy, strategy_manifest)), engine: Weak::new(), ws: hypercore::mainnet_ws(), subscriptions: Mutex::new(HashSet::new()), diff --git a/src/engine/store/mod.rs b/src/engine/store/mod.rs index 73b6dee..fe051df 100644 --- a/src/engine/store/mod.rs +++ b/src/engine/store/mod.rs @@ -2,7 +2,7 @@ use std::path::PathBuf; pub mod accounts; pub mod config; -pub mod plugin; +pub mod strategy; pub fn home_dir() -> tokio::io::Result { std::env::home_dir().ok_or_else(|| { diff --git a/src/engine/store/plugin.rs b/src/engine/store/strategy.rs similarity index 69% rename from src/engine/store/plugin.rs rename to src/engine/store/strategy.rs index 62449f6..fa7a99f 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/strategy.rs @@ -1,33 +1,29 @@ -use std::marker::PhantomData; - -use pulse_sdk::map_postcard_err; -use serde::{Deserialize, Serialize}; +use pulse_sdk::{ + map_postcard_err, + strategy::{StrategyEngineMessage, StrategyManifest, StrategyMessage}, +}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, process::{Child, ChildStdout}, sync::Mutex, }; #[derive(Debug)] -pub struct Plugin Deserialize<'de>, M: for<'de> Deserialize<'de>> { - pub manifest: Mutex, +pub struct StrategyChild { + pub manifest: Mutex, pub stdout: Mutex, pub process: Mutex, - - pub _p: (PhantomData, PhantomData), } -impl Deserialize<'de>, M: for<'de> Deserialize<'de>> Plugin { - pub fn new(mut child: Child, manifest: M) -> Self { +impl StrategyChild { + pub fn new(mut child: Child, manifest: StrategyManifest) -> Self { Self { stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")), process: Mutex::new(child), manifest: Mutex::new(manifest), - - _p: (PhantomData, PhantomData), } } - pub async fn recv(&self) -> tokio::io::Result> { + pub async fn recv(&self) -> tokio::io::Result> { let mut stdout = self.stdout.lock().await; let mut len_buf = [0u8; size_of::()]; @@ -46,7 +42,7 @@ impl Deserialize<'de>, M: for<'de> Deserialize<'de>> P Ok(Some(map_postcard_err(postcard::from_bytes(&buffer))?)) } - pub async fn send(&self, msg: &S) -> tokio::io::Result<()> { + pub async fn send(&self, msg: &StrategyEngineMessage) -> tokio::io::Result<()> { self.send_raw(&map_postcard_err(postcard::to_allocvec(msg))?) .await } @@ -62,7 +58,11 @@ impl Deserialize<'de>, M: for<'de> Deserialize<'de>> P Ok(()) } - pub async fn reload(&self, mut child: Child, manifest: M) -> tokio::io::Result<()> { + pub async fn reload( + &self, + mut child: Child, + manifest: StrategyManifest, + ) -> tokio::io::Result<()> { let mut process = self.process.lock().await; process.kill().await?; From 0b169dfc3829ce71654897a1ad6eafa094265b7e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 30 Jul 2026 03:46:40 +0200 Subject: [PATCH 12/12] Renamed plugins to strategies --- src/engine/engine/command.rs | 2 +- src/engine/engine/strategy.rs | 26 +++++++++++++------------- src/engine/store/mod.rs | 8 ++++---- 3 files changed, 18 insertions(+), 18 deletions(-) diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 080a394..1facea8 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -47,7 +47,7 @@ impl Engine { set_cfg!(id, { let id: String = id; - if !crate::store::pulse_plugin(&id)? + if !crate::store::pulse_strategy(&id)? .join("strategy.toml") .exists() { diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index 0ab4b0f..8f8f074 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -15,7 +15,7 @@ use tokio::{ use crate::{ engine::Engine, - store::{strategy::StrategyChild, pulse_plugin}, + store::{pulse_strategy, strategy::StrategyChild}, }; pub struct StrategyEngine { @@ -28,9 +28,9 @@ pub struct StrategyEngine { impl StrategyEngine { pub async fn new(strategy_id: &str) -> tokio::io::Result { - let strategy = pulse_plugin(strategy_id)?; + let strategy = pulse_strategy(strategy_id)?; - let (strategy, strategy_manifest) = get_manifest_plugin_pair( + let (strategy, strategy_manifest) = get_manifest( &strategy, &strategy.join("strategy.bash"), &fs::read(strategy.join("strategy.toml")).await?, @@ -168,12 +168,12 @@ impl StrategyEngine { } pub async fn reload_strategy(self: &Arc, id: &str) -> tokio::io::Result<()> { - let plugin = pulse_plugin(id)?; + let strategy = pulse_strategy(id)?; - let (child, manifest) = get_manifest_plugin_pair( - &plugin, - &plugin.join("strategy.bash"), - &fs::read(plugin.join("strategy.toml")).await?, + let (child, manifest) = get_manifest( + &strategy, + &strategy.join("strategy.bash"), + &fs::read(strategy.join("strategy.toml")).await?, )?; self.strategy.reload(child, manifest).await?; @@ -185,15 +185,15 @@ impl StrategyEngine { } } -pub fn get_manifest_plugin_pair<'de, M: serde::Deserialize<'de>>( - plugin_dir: &PathBuf, - plugin_path: &PathBuf, +fn get_manifest<'de, M: serde::Deserialize<'de>>( + strategy_dir: &PathBuf, + strategy_path: &PathBuf, manifest: &'de [u8], ) -> tokio::io::Result<(Child, M)> { Ok(( Command::new("bash") - .arg(plugin_path) - .current_dir(plugin_dir) + .arg(strategy_path) + .current_dir(strategy_dir) .stdin(Stdio::piped()) .stdout(Stdio::piped()) .stderr(Stdio::inherit()) diff --git a/src/engine/store/mod.rs b/src/engine/store/mod.rs index fe051df..7552bd3 100644 --- a/src/engine/store/mod.rs +++ b/src/engine/store/mod.rs @@ -14,12 +14,12 @@ pub fn pulse_directory() -> tokio::io::Result { Ok(home_dir()?.join(".config").join("pulse-trader")) } -pub fn pulse_plugins_directory() -> tokio::io::Result { - Ok(pulse_directory()?.join("plugins")) +pub fn pulse_strategies_directory() -> tokio::io::Result { + Ok(pulse_directory()?.join("strategies")) } -pub fn pulse_plugin(id: &str) -> tokio::io::Result { - Ok(pulse_directory()?.join("plugins").join(id)) +pub fn pulse_strategy(id: &str) -> tokio::io::Result { + Ok(pulse_directory()?.join("strategies").join(id)) } pub fn pulse_config_file() -> tokio::io::Result {