Risk stuff and trying to test different ws archs

This commit is contained in:
2026-08-01 06:56:15 +02:00
parent a9ccd84633
commit f48e239d05
3 changed files with 108 additions and 29 deletions
+18 -8
View File
@@ -25,7 +25,6 @@ pub struct 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 risk_engine: Arc<RiskEngine>,
pub ws_stream: Arc<Mutex<ConnectionStream>>,
// data // data
pub config: Arc<Mutex<ConfigManager>>, pub config: Arc<Mutex<ConfigManager>>,
@@ -37,23 +36,29 @@ pub struct Engine {
pub signals: Arc<Mutex<Vec<SignalStatus>>>, pub signals: Arc<Mutex<Vec<SignalStatus>>>,
} }
pub struct EngineSides {
pub strategy_stream: ConnectionStream,
pub risk_stream: ConnectionStream,
}
impl Engine { impl Engine {
pub async fn new() -> tokio::io::Result<Arc<Self>> { pub async fn new() -> tokio::io::Result<(Arc<Self>, EngineSides)> {
let config = ConfigManager::new().await?; 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 accounts = Arc::new(Mutex::new(AccountList::new().await?));
let config = Arc::new(Mutex::new(config)); let config = Arc::new(Mutex::new(config));
Ok(Arc::new_cyclic(|engine| Self { Ok((
Arc::new_cyclic(|engine| Self {
// 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()), risk_engine: RiskEngine::new(engine.clone(), risk_handle),
ws_stream: Arc::new(Mutex::new(ws_stream)),
// data // data
config, config,
@@ -69,7 +74,12 @@ impl Engine {
items: Vec::new(), items: Vec::new(),
})), })),
signals: Arc::new(Mutex::new(Vec::new())), signals: Arc::new(Mutex::new(Vec::new())),
})) }),
EngineSides {
strategy_stream,
risk_stream,
},
))
} }
/// Starts the main strategy server /// Starts the main strategy server
+73 -4
View File
@@ -1,6 +1,13 @@
use std::sync::{Arc, Weak}; 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 pulse_sdk::general::Signal;
use tokio::sync::Mutex; 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 { pub struct RiskEngine {
// Orders made by the engine // Orders made by the engine
pub orders: Mutex<Vec<OrderIds>>, pub orders: Mutex<Vec<OrderIds>>,
pub ws: Mutex<WebSocket>, pub handle: Mutex<ConnectionHandle>,
pub state: Mutex<RiskState>,
pub engine: Weak<Engine>, pub engine: Weak<Engine>,
} }
impl RiskEngine { impl RiskEngine {
pub fn new(engine: Weak<Engine>) -> Arc<Self> { pub fn new(engine: Weak<Engine>, handle: ConnectionHandle) -> Arc<Self> {
Arc::new(Self { Arc::new(Self {
orders: Mutex::new(Vec::new()), 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, 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 { pub async fn validate_signal(&self, _signal: &mut Signal) -> bool {
true true
} }
+1 -1
View File
@@ -4,7 +4,7 @@ pub mod store;
#[tokio::main] #[tokio::main]
async fn main() -> anyhow::Result<()> { 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 server = engine.terminal_server.spawn_server().await;
let broadcaster = engine.terminal_server.spawn_broadcaster().await; let broadcaster = engine.terminal_server.spawn_broadcaster().await;