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? } }