Strategy and Risk websocket clients
This commit is contained in:
@@ -5,7 +5,7 @@ use hypersdk::{
|
|||||||
Decimal,
|
Decimal,
|
||||||
hypercore::{
|
hypercore::{
|
||||||
self, Cloid,
|
self, Cloid,
|
||||||
ws::{ConnectionHandle, ConnectionStream},
|
ws::{ConnectionHandle, ConnectionStream, Event},
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
use pulse_sdk::general::Signal;
|
use pulse_sdk::general::Signal;
|
||||||
@@ -75,10 +75,15 @@ impl RiskEngine {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn run(&self, mut stream: ConnectionStream) -> anyhow::Result<()> {
|
pub async fn run_event_stream(
|
||||||
|
self: Arc<Self>,
|
||||||
|
mut stream: ConnectionStream,
|
||||||
|
) -> anyhow::Result<()> {
|
||||||
|
let engine = self.get_engine();
|
||||||
|
|
||||||
while let Some(e) = stream.next().await {
|
while let Some(e) = stream.next().await {
|
||||||
match e {
|
match e {
|
||||||
hypercore::ws::Event::Message(hypercore::Incoming::OrderUpdates(order_updates)) => {
|
Event::Message(hypercore::Incoming::OrderUpdates(order_updates)) => {
|
||||||
for order in order_updates {
|
for order in order_updates {
|
||||||
if let Some(cloid) = order.order.cloid {
|
if let Some(cloid) = order.order.cloid {
|
||||||
let orders = self.orders.lock().await.clone();
|
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?;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,5 +1,9 @@
|
|||||||
use anyhow::Context;
|
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 pulse_sdk::prelude::*;
|
||||||
use std::{
|
use std::{
|
||||||
collections::HashSet,
|
collections::HashSet,
|
||||||
@@ -75,11 +79,40 @@ impl StrategyEngine {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub async fn run_event_stream(
|
||||||
|
self: Arc<Self>,
|
||||||
|
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<Self>) -> anyhow::Result<()> {
|
pub async fn run(self: &Arc<Self>) -> anyhow::Result<()> {
|
||||||
let engine = self
|
let engine = self.get_engine();
|
||||||
.engine
|
|
||||||
.upgrade()
|
|
||||||
.expect("Failed to upgrade engine (StrategyEngine)");
|
|
||||||
|
|
||||||
let mut stdout = {
|
let mut stdout = {
|
||||||
let mut child = self.strategy.lock().await;
|
let mut child = self.strategy.lock().await;
|
||||||
@@ -211,4 +244,10 @@ impl StrategyEngine {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn get_engine(&self) -> Arc<Engine> {
|
||||||
|
self.engine
|
||||||
|
.upgrade()
|
||||||
|
.expect("Failed to upgrade engine (StrategyEngine)")
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,10 +9,27 @@ async fn main() -> anyhow::Result<()> {
|
|||||||
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;
|
||||||
|
|
||||||
|
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?;
|
engine.run().await?;
|
||||||
|
|
||||||
server.await??;
|
server.await??;
|
||||||
broadcaster.await??;
|
broadcaster.await??;
|
||||||
|
|
||||||
|
risk.await??;
|
||||||
|
strategy.await??;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user