diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 5cee75e..c3c7a75 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -4,6 +4,10 @@ use pulse_sdk::prelude::*; use toml::Value; impl Engine { + pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> { + self.terminal_server.error(name, "Invalid usage").await + } + pub async fn execute_command( &self, command: &str, diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 89bc936..d0d6f5c 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -10,7 +10,7 @@ use crate::{ use hypersdk::hypercore::ws::ConnectionStream; use pulse_sdk::prelude::*; use std::{collections::HashMap, sync::Arc}; -use tokio::{sync::Mutex, task::JoinHandle}; +use tokio::sync::Mutex; #[derive(Debug, Clone)] pub struct WatchList { @@ -56,87 +56,4 @@ impl Engine { signals: Arc::new(Mutex::new(Vec::new())), })) } - - pub async fn spawn_broadcaster(&self) -> JoinHandle> { - let s = self.clone(); - - tokio::spawn(async move { s.run_broadcaster().await }) - } - - pub async fn run_broadcaster(&self) -> tokio::io::Result<()> { - let mut refresh = tokio::time::interval(tokio::time::Duration::from_secs(5)); - - let client = hypersdk::hypercore::mainnet(); - - loop { - refresh.tick().await; - - let watch_list = &self.config.lock().await.watchlist; - - match crate::fetch::fetch_watch_list(&client, watch_list).await { - Ok(watch_list) => { - *self.watch_list.lock().await = watch_list.clone(); - - if let Err(error) = self - .terminal_server - .broadcast(TerminalServerMessage::WatchListUpdated(watch_list.items)) - .await - { - self.terminal_server - .error( - "Broadcaster", - &format!("Failed to broadcast HyperLiquid watch list: {error}"), - ) - .await? - } - } - - Err(error) => { - self.terminal_server - .error( - "Broadcaster", - &format!("Failed to refresh HyperLiquid watch list: {error}"), - ) - .await? - } - } - - let accounts = self.accounts.lock().await; - - if let Some(acc) = accounts.get_active() { - match client.clearinghouse_state(acc.address, None).await { - Ok(state) => { - self.terminal_server - .broadcast(TerminalServerMessage::PositionsUpdated( - state - .asset_positions - .into_iter() - .map(|position| Position { - symbol: position.position.coin, - size: position.position.szi, - entry_price: position.position.entry_px.unwrap_or_default(), - pnl: position.position.unrealized_pnl, - }) - .collect(), - )) - .await?; - } - Err(e) => { - self.terminal_server - .error("orders", &format!("Unable to get open orders: {e}")) - .await?; - } - } - } else { - self.terminal_server.error( - "orders", - "Unable to get active account, make sure you have configured accounts properly", - ).await?; - } - } - } - - pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> { - self.terminal_server.error(name, "Invalid usage").await - } } diff --git a/src/engine/main.rs b/src/engine/main.rs index 1393e79..af859e4 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -7,12 +7,12 @@ pub mod terminal; async fn main() -> anyhow::Result<()> { let engine = engine::Engine::new().await?; - let broadcaster = engine.spawn_broadcaster().await; + let server = engine.terminal_server.spawn_server().await; + let broadcaster = engine.terminal_server.spawn_broadcaster().await; engine.strategy_engine.spawn().await; - engine.terminal_server.run().await?; - + server.await??; broadcaster.await??; Ok(()) diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 619d564..e308812 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -11,6 +11,7 @@ use tokio::{ unix::{OwnedReadHalf, OwnedWriteHalf}, }, sync::Mutex, + task::JoinHandle, }; #[derive(Debug)] @@ -29,7 +30,13 @@ impl TerminalServer { }) } - pub async fn run(self: &Arc) -> tokio::io::Result<()> { + pub async fn spawn_server(self: &Arc) -> JoinHandle> { + let s = self.clone(); + + tokio::spawn(s.run_server()) + } + + pub async fn run_server(self: Arc) -> tokio::io::Result<()> { let path = pulse_sdk::server_path(); if path.exists() { @@ -151,6 +158,81 @@ impl TerminalServer { Ok(()) } + pub async fn spawn_broadcaster(self: &Arc) -> JoinHandle> { + let s = self.clone(); + + tokio::spawn(s.run_broadcaster()) + } + + pub async fn run_broadcaster(self: Arc) -> tokio::io::Result<()> { + let engine = self.get_engine(); + + let mut refresh = tokio::time::interval(tokio::time::Duration::from_secs(5)); + + let client = hypersdk::hypercore::mainnet(); + + loop { + refresh.tick().await; + + match crate::fetch::fetch_watch_list(&client, &engine.config.lock().await.watchlist) + .await + { + Ok(watch_list) => { + *engine.watch_list.lock().await = watch_list.clone(); + + if let Err(error) = self + .broadcast(TerminalServerMessage::WatchListUpdated(watch_list.items)) + .await + { + self.error( + "Broadcaster", + &format!("Failed to broadcast HyperLiquid watch list: {error}"), + ) + .await? + } + } + + Err(error) => { + self.error( + "Broadcaster", + &format!("Failed to refresh HyperLiquid watch list: {error}"), + ) + .await? + } + } + + if let Some(acc) = engine.accounts.lock().await.get_active() { + match client.clearinghouse_state(acc.address, None).await { + Ok(state) => { + self.broadcast(TerminalServerMessage::PositionsUpdated( + state + .asset_positions + .into_iter() + .map(|position| Position { + symbol: position.position.coin, + size: position.position.szi, + entry_price: position.position.entry_px.unwrap_or_default(), + pnl: position.position.unrealized_pnl, + }) + .collect(), + )) + .await?; + } + Err(e) => { + self.error("orders", &format!("Unable to get open orders: {e}")) + .await?; + } + } + } else { + self.error( + "orders", + "Unable to get active account, make sure you have configured accounts properly", + ) + .await?; + } + } + } + pub async fn send_to( self: &Arc, id: &usize,