diff --git a/src/engine/engine/risk.rs b/src/engine/engine/risk.rs index 58eb06d..428c034 100644 --- a/src/engine/engine/risk.rs +++ b/src/engine/engine/risk.rs @@ -5,7 +5,7 @@ use hypersdk::{ Decimal, hypercore::{ self, Cloid, - ws::{ConnectionHandle, ConnectionStream}, + ws::{ConnectionHandle, ConnectionStream, Event}, }, }; use pulse_sdk::general::Signal; @@ -75,10 +75,15 @@ impl RiskEngine { Ok(()) } - pub async fn run(&self, mut stream: ConnectionStream) -> anyhow::Result<()> { + pub async fn run_event_stream( + self: Arc, + mut stream: ConnectionStream, + ) -> anyhow::Result<()> { + let engine = self.get_engine(); + while let Some(e) = stream.next().await { match e { - hypercore::ws::Event::Message(hypercore::Incoming::OrderUpdates(order_updates)) => { + 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(); @@ -101,7 +106,30 @@ impl RiskEngine { } } } - _ => {} + + Event::Connected => { + engine + .terminal_server + .info("engine::risk", "HyperLiquid WebSocket connected") + .await?; + } + + Event::Disconnected => { + engine + .terminal_server + .warn("engine::risk", "HyperLiquid WebSocket disconnected") + .await?; + } + + _ => { + engine + .terminal_server + .warn( + "engine::risk", + &format!("Unexpected event from HyperLiquid WebSocket: {e:?}"), + ) + .await?; + } } } diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index e358aee..56621fc 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -1,5 +1,9 @@ use anyhow::Context; -use hypersdk::hypercore::{self, CandleInterval, Subscription, ws::ConnectionHandle}; +use futures::StreamExt; +use hypersdk::hypercore::{ + self, CandleInterval, Subscription, + ws::{ConnectionHandle, ConnectionStream, Event}, +}; use pulse_sdk::prelude::*; use std::{ collections::HashSet, @@ -75,11 +79,40 @@ impl StrategyEngine { Ok(()) } + pub async fn run_event_stream( + self: Arc, + mut stream: ConnectionStream, + ) -> anyhow::Result<()> { + let engine = self.get_engine(); + + while let Some(e) = stream.next().await { + match e { + Event::Message(incoming) => { + self.send(&StrategyEngineMessage::Incoming(incoming)) + .await?; + } + + Event::Connected => { + engine + .terminal_server + .info("engine::strategy", "HyperLiquid WebSocket connected") + .await?; + } + + Event::Disconnected => { + engine + .terminal_server + .warn("engine::strategy", "HyperLiquid WebSocket disconnected") + .await?; + } + } + } + + Ok(()) + } + pub async fn run(self: &Arc) -> anyhow::Result<()> { - let engine = self - .engine - .upgrade() - .expect("Failed to upgrade engine (StrategyEngine)"); + let engine = self.get_engine(); let mut stdout = { let mut child = self.strategy.lock().await; @@ -211,4 +244,10 @@ impl StrategyEngine { } } } + + pub fn get_engine(&self) -> Arc { + self.engine + .upgrade() + .expect("Failed to upgrade engine (StrategyEngine)") + } } diff --git a/src/engine/main.rs b/src/engine/main.rs index 8531472..e8be395 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -9,10 +9,27 @@ async fn main() -> anyhow::Result<()> { let server = engine.terminal_server.spawn_server().await; let broadcaster = engine.terminal_server.spawn_broadcaster().await; + let risk = tokio::spawn( + engine + .risk_engine + .clone() + .run_event_stream(sides.risk_stream), + ); + + let strategy = tokio::spawn( + engine + .strategy_engine + .clone() + .run_event_stream(sides.strategy_stream), + ); + engine.run().await?; server.await??; broadcaster.await??; + risk.await??; + strategy.await??; + Ok(()) }