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()