Risk engine
This commit is contained in:
@@ -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 pulse_sdk::prelude::*;
|
||||||
|
|
||||||
use crate::engine::Engine;
|
use crate::engine::{Engine, risk::OrderIds};
|
||||||
|
|
||||||
impl Engine {
|
impl Engine {
|
||||||
pub async fn execute_signal(&self, signal: &Signal) -> tokio::io::Result<()> {
|
pub async fn execute_signal(&self, signal: &Signal) -> tokio::io::Result<()> {
|
||||||
@@ -30,6 +32,8 @@ impl Engine {
|
|||||||
return Ok(());
|
return Ok(());
|
||||||
};
|
};
|
||||||
|
|
||||||
|
let order_ids = OrderIds::new();
|
||||||
|
|
||||||
let order = BatchOrder {
|
let order = BatchOrder {
|
||||||
orders: vec![
|
orders: vec![
|
||||||
OrderRequest {
|
OrderRequest {
|
||||||
@@ -41,7 +45,7 @@ impl Engine {
|
|||||||
order_type: OrderTypePlacement::Limit {
|
order_type: OrderTypePlacement::Limit {
|
||||||
tif: TimeInForce::Gtc,
|
tif: TimeInForce::Gtc,
|
||||||
},
|
},
|
||||||
cloid: Default::default(),
|
cloid: order_ids.entry,
|
||||||
},
|
},
|
||||||
OrderRequest {
|
OrderRequest {
|
||||||
asset: asset_id,
|
asset: asset_id,
|
||||||
@@ -54,7 +58,7 @@ impl Engine {
|
|||||||
trigger_px: signal.take_profit,
|
trigger_px: signal.take_profit,
|
||||||
tpsl: hypercore::TpSl::Tp,
|
tpsl: hypercore::TpSl::Tp,
|
||||||
},
|
},
|
||||||
cloid: Default::default(),
|
cloid: order_ids.take_profit,
|
||||||
},
|
},
|
||||||
OrderRequest {
|
OrderRequest {
|
||||||
asset: asset_id,
|
asset: asset_id,
|
||||||
@@ -67,7 +71,7 @@ impl Engine {
|
|||||||
trigger_px: signal.stop_loss,
|
trigger_px: signal.stop_loss,
|
||||||
tpsl: hypercore::TpSl::Sl,
|
tpsl: hypercore::TpSl::Sl,
|
||||||
},
|
},
|
||||||
cloid: Default::default(),
|
cloid: order_ids.stop_loss,
|
||||||
},
|
},
|
||||||
],
|
],
|
||||||
grouping: hypercore::OrderGrouping::Na,
|
grouping: hypercore::OrderGrouping::Na,
|
||||||
@@ -80,7 +84,31 @@ impl Engine {
|
|||||||
.place(&acc.private_key.0, order, nonce, None, None)
|
.place(&acc.private_key.0, order, nonce, None, None)
|
||||||
.await
|
.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::<Vec<_>>()
|
||||||
|
.join(", ")
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
|
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
self.terminal_server
|
self.terminal_server
|
||||||
.error("self::order", &e.to_string())
|
.error("self::order", &e.to_string())
|
||||||
|
|||||||
@@ -1,10 +1,11 @@
|
|||||||
pub mod command;
|
pub mod command;
|
||||||
pub mod execution;
|
pub mod execution;
|
||||||
|
pub mod risk;
|
||||||
pub mod strategy;
|
pub mod strategy;
|
||||||
pub mod terminal;
|
pub mod terminal;
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
engine::{strategy::StrategyEngine, terminal::TerminalServer},
|
engine::{risk::RiskEngine, strategy::StrategyEngine, terminal::TerminalServer},
|
||||||
store::{accounts::AccountList, config::ConfigManager},
|
store::{accounts::AccountList, config::ConfigManager},
|
||||||
};
|
};
|
||||||
use hypersdk::hypercore::ws::ConnectionStream;
|
use hypersdk::hypercore::ws::ConnectionStream;
|
||||||
@@ -23,6 +24,7 @@ pub struct Engine {
|
|||||||
// engine
|
// engine
|
||||||
pub terminal_server: Arc<TerminalServer>,
|
pub terminal_server: Arc<TerminalServer>,
|
||||||
pub strategy_engine: Arc<StrategyEngine>,
|
pub strategy_engine: Arc<StrategyEngine>,
|
||||||
|
pub risk_engine: Arc<RiskEngine>,
|
||||||
pub ws_stream: Arc<Mutex<ConnectionStream>>,
|
pub ws_stream: Arc<Mutex<ConnectionStream>>,
|
||||||
|
|
||||||
// data
|
// data
|
||||||
@@ -50,6 +52,7 @@ impl Engine {
|
|||||||
// engine
|
// engine
|
||||||
terminal_server: TerminalServer::new(engine.clone()),
|
terminal_server: TerminalServer::new(engine.clone()),
|
||||||
strategy_engine: strategy.initialize(engine.clone()),
|
strategy_engine: strategy.initialize(engine.clone()),
|
||||||
|
risk_engine: RiskEngine::new(engine.clone()),
|
||||||
ws_stream: Arc::new(Mutex::new(ws_stream)),
|
ws_stream: Arc::new(Mutex::new(ws_stream)),
|
||||||
|
|
||||||
// data
|
// data
|
||||||
|
|||||||
@@ -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<Vec<OrderIds>>,
|
||||||
|
pub ws: Mutex<WebSocket>,
|
||||||
|
pub engine: Weak<Engine>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl RiskEngine {
|
||||||
|
pub fn new(engine: Weak<Engine>) -> Arc<Self> {
|
||||||
|
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<Engine> {
|
||||||
|
self.engine
|
||||||
|
.upgrade()
|
||||||
|
.expect("Failed to upgrade engine(Weak) to Arc")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -117,7 +117,18 @@ impl StrategyEngine {
|
|||||||
engine.terminal_server.log_raw(log).await?;
|
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 {
|
match engine.execute_signal(&signal).await {
|
||||||
Ok(_) => engine.signals.lock().await.push(Ok(signal)),
|
Ok(_) => engine.signals.lock().await.push(Ok(signal)),
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
@@ -125,7 +136,10 @@ impl StrategyEngine {
|
|||||||
|
|
||||||
engine
|
engine
|
||||||
.terminal_server
|
.terminal_server
|
||||||
.error("signal", &format!("Failed to execute signal: {e}"))
|
.error(
|
||||||
|
"engine::strategy",
|
||||||
|
&format!("Failed to execute signal: {e}"),
|
||||||
|
)
|
||||||
.await?
|
.await?
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user