Refactored engine

This commit is contained in:
2026-07-31 01:26:03 +02:00
parent fccd03c3df
commit b5a1a07cf0
4 changed files with 91 additions and 88 deletions
+4
View File
@@ -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,
+1 -84
View File
@@ -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<tokio::io::Result<()>> {
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
}
}
+3 -3
View File
@@ -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(())
+83 -1
View File
@@ -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<Self>) -> tokio::io::Result<()> {
pub async fn spawn_server(self: &Arc<Self>) -> JoinHandle<tokio::io::Result<()>> {
let s = self.clone();
tokio::spawn(s.run_server())
}
pub async fn run_server(self: Arc<Self>) -> 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<Self>) -> JoinHandle<tokio::io::Result<()>> {
let s = self.clone();
tokio::spawn(s.run_broadcaster())
}
pub async fn run_broadcaster(self: Arc<Self>) -> 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<Self>,
id: &usize,