Files
pulse-trader/src/engine/engine.rs
T

236 lines
8.3 KiB
Rust

use std::sync::Arc;
use pulse_wire::{
general::Position,
units::{Symbol, USD},
};
use tokio::{sync::Mutex, task::JoinHandle};
use crate::{
store::{accounts::AccountList, config::Config},
terminal::TerminalServer,
};
#[derive(Debug, Clone)]
pub struct Engine {
pub terminal_server: Arc<TerminalServer>,
pub config: Arc<Mutex<Config>>,
pub accounts: Arc<Mutex<AccountList>>,
}
impl Engine {
pub async fn new() -> tokio::io::Result<Arc<Self>> {
let config = Arc::new(Mutex::new(Config::new().await?));
let accounts = Arc::new(Mutex::new(AccountList::new().await?));
Ok(Arc::new_cyclic(|engine| Self {
terminal_server: TerminalServer::new(engine.clone()),
config,
accounts,
}))
}
pub async fn spawn_terminal_server(&self) -> JoinHandle<tokio::io::Result<()>> {
let terminal_server = self.terminal_server.clone();
tokio::spawn(async move { terminal_server.run().await })
}
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.symbols;
match crate::fetch::fetch_watch_list(&client, watch_list).await {
Ok(watch_list) => {
if let Err(error) = self
.terminal_server
.broadcast(
pulse_wire::terminal::TerminalServerMessage::WatchListUpdated(
watch_list,
),
)
.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(
pulse_wire::terminal::TerminalServerMessage::PositionsUpdated(
state
.asset_positions
.into_iter()
.map(|position| Position {
symbol: Symbol(position.position.coin),
size: position.position.szi.as_f64(),
entry_price: USD(position
.position
.entry_px
.map(|px| px.as_f64())
.unwrap_or(0.0)),
profit: USD(position.position.unrealized_pnl.as_f64()),
})
.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 run_engine(&self) -> tokio::io::Result<()> {
loop {
tokio::time::sleep(tokio::time::Duration::from_millis(5000)).await;
}
}
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, args: Vec<&str>) -> tokio::io::Result<()> {
match command {
"config" | "cfg" => {
if args.len() != 1 {
return self.invalid_command_usage("config").await;
}
match args[0] {
"reload" => {
*self.config.lock().await = Config::new().await?;
}
_ => {
self.terminal_server
.error("config", "Invalid usage")
.await?;
}
}
}
"account" | "acc" => {
if args.len() == 0 {
return self.invalid_command_usage("account man").await;
}
match args[0] {
"reload" => {
*self.accounts.lock().await = AccountList::new().await?;
}
"list" | "ls" => {
self.terminal_server
.info("account man", "ACCOUNT LIST")
.await?;
let accounts = self.accounts.lock().await;
for (name, acc) in &accounts.accounts {
self.terminal_server
.info(
"account man",
&if name == &accounts.active {
format!(
"{} (active) -> {}",
name,
acc.get_truncated_address()
)
} else {
format!("{} -> {}", name, acc.get_truncated_address())
},
)
.await?;
}
}
"use" | "set" => {
if args.len() < 2 {
return self.invalid_command_usage("account man").await;
}
let new_active = args[1];
let mut accounts = self.accounts.lock().await;
if !accounts.accounts.contains_key(new_active) {
return self
.terminal_server
.error("account man", &format!("Account not found ({new_active})"))
.await;
}
accounts.active = new_active.to_string();
self.terminal_server
.info(
"account man",
&format!("Account set to {new_active} successfully!"),
)
.await?;
}
_ => {
return self.invalid_command_usage("account man").await;
}
}
}
_ => {
self.terminal_server
.error(
"Command executor",
&format!("Command '{}' not found", command),
)
.await?;
}
}
Ok(())
}
}