Merge pull request #17 from selimaj-dev/orders-and-strategy-redesign

Orders and strategy redesign
This commit is contained in:
2026-07-30 04:14:29 +02:00
committed by GitHub
21 changed files with 490 additions and 393 deletions
Generated
+2
View File
@@ -3947,6 +3947,7 @@ version = "0.1.0-alpha.0"
dependencies = [ dependencies = [
"hypersdk", "hypersdk",
"postcard", "postcard",
"rust_decimal",
"serde", "serde",
"tokio", "tokio",
] ]
@@ -3963,6 +3964,7 @@ dependencies = [
"pulse-sdk", "pulse-sdk",
"pulse-ui", "pulse-ui",
"rand 0.8.7", "rand 0.8.7",
"rust_decimal",
"serde", "serde",
"serde_json", "serde_json",
"tokio", "tokio",
+2
View File
@@ -11,6 +11,7 @@ hypersdk = { workspace = true }
serde = { workspace = true } serde = { workspace = true }
anyhow = { workspace = true } anyhow = { workspace = true }
postcard = { workspace = true } postcard = { workspace = true }
rust_decimal = { workspace = true }
tokio = { workspace = true, features = [ tokio = { workspace = true, features = [
"rt-multi-thread", "rt-multi-thread",
"macros", "macros",
@@ -29,6 +30,7 @@ toml = "1.1.3"
members = ["pulse-ui", "pulse-sdk"] members = ["pulse-ui", "pulse-sdk"]
[workspace.dependencies] [workspace.dependencies]
rust_decimal = { version = "1.39", features = ["serde-str"] }
postcard = { version = "1.1.3", features = ["alloc"] } postcard = { version = "1.1.3", features = ["alloc"] }
pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" } pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" }
pulse-sdk = { path = "pulse-sdk", version = "0.1.0-alpha.0" } pulse-sdk = { path = "pulse-sdk", version = "0.1.0-alpha.0" }
+2 -1
View File
@@ -7,4 +7,5 @@ edition = "2024"
serde = { workspace = true } serde = { workspace = true }
hypersdk = { workspace = true } hypersdk = { workspace = true }
postcard = { workspace = true } postcard = { workspace = true }
tokio = { workspace = true } tokio = { workspace = true, features = ["io-std"] }
rust_decimal = { workspace = true }
+3 -1
View File
@@ -1,3 +1,5 @@
use rust_decimal::Decimal;
use crate::units::{Direction, Symbol, USD}; use crate::units::{Direction, Symbol, USD};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
@@ -20,7 +22,7 @@ pub struct Signal {
pub symbol: String, pub symbol: String,
pub kind: Direction, pub kind: Direction,
pub confidence: f32, pub confidence: f32,
pub size: f64, pub size: Decimal,
pub price: USD, pub price: USD,
pub take_profit: USD, pub take_profit: USD,
pub stop_loss: USD, pub stop_loss: USD,
+96 -7
View File
@@ -1,25 +1,114 @@
#[cfg(target_os = "macos")]
use std::path::PathBuf;
pub mod general; pub mod general;
pub mod plugin; pub mod strategy;
pub mod terminal; pub mod terminal;
pub mod units; pub mod units;
pub use hypersdk; pub use hypersdk;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use crate::strategy::StrategyEngineMessage;
pub mod prelude { pub mod prelude {
pub use crate::general::*; pub use crate::general::*;
pub use crate::plugin::*; pub use crate::strategy::*;
pub use crate::server_path; pub use crate::server_path;
pub use crate::terminal::*; pub use crate::terminal::*;
pub use crate::units::*; pub use crate::units::*;
pub use hypersdk; pub use hypersdk;
pub use postcard;
} }
pub fn server_path() -> PathBuf { pub fn server_path() -> std::path::PathBuf {
PathBuf::from("/tmp/pulse-engine.sock") std::path::PathBuf::from("/tmp/pulse-engine.sock")
} }
pub fn map_postcard_err<T>(res: postcard::Result<T>) -> tokio::io::Result<T> { pub fn map_postcard_err<T>(res: postcard::Result<T>) -> tokio::io::Result<T> {
res.map_err(|e| tokio::io::Error::new(std::io::ErrorKind::Other, e)) res.map_err(|e| tokio::io::Error::new(std::io::ErrorKind::Other, e))
} }
pub async fn send_raw(data: &[u8]) -> tokio::io::Result<()> {
let mut stdout = tokio::io::stdout();
stdout.write_all(&data.len().to_le_bytes()).await?;
stdout.write_all(data).await?;
stdout.flush().await?;
Ok(())
}
macro_rules! engine_methods {
($t:ty) => {
async fn start(&self) -> tokio::io::Result<()> {
let mut stdin = tokio::io::stdin();
loop {
let mut len_buf = [0u8; size_of::<usize>()];
let size = stdin.read_exact(&mut len_buf).await?;
let len = usize::from_le_bytes(len_buf);
if size == 0 || len == 0 {
break Ok(());
}
let mut buffer = vec![0u8; len];
stdin.read_exact(&mut buffer).await?;
self.on_raw(&buffer).await?;
}
}
async fn send(&self, msg: &$t) -> tokio::io::Result<()> {
$crate::send_raw(&$crate::map_postcard_err(
$crate::prelude::postcard::to_allocvec(msg),
)?)
.await
}
};
}
#[allow(async_fn_in_trait)]
pub trait Strategy {
engine_methods!(prelude::StrategyMessage);
async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> {
match map_postcard_err(postcard::from_bytes(&data))? {
StrategyEngineMessage::Initialize => self.initialize().await,
StrategyEngineMessage::Command { command, args } => self.command(command, args).await,
StrategyEngineMessage::WatchList(watchlist) => self.watchlist(watchlist).await,
StrategyEngineMessage::Incoming(incoming) => self.incoming(incoming).await,
StrategyEngineMessage::Candlestick {
symbol,
interval,
candles,
} => self.candlestick(symbol, interval, candles).await,
}
}
async fn initialize(&self) -> tokio::io::Result<()> {
unimplemented!("Strategy::initialize")
}
async fn command(&self, _command: String, _args: Vec<String>) -> tokio::io::Result<()> {
unimplemented!("Strategy::command")
}
async fn watchlist(&self, _watchlist: Vec<prelude::MarketItem>) -> tokio::io::Result<()> {
unimplemented!("Strategy::watchlist")
}
async fn incoming(&self, _incoming: hypersdk::hypercore::Incoming) -> tokio::io::Result<()> {
unimplemented!("Strategy::event")
}
async fn candlestick(
&self,
_symbol: String,
_interval: hypersdk::hypercore::CandleInterval,
_candles: Vec<hypersdk::hypercore::Candle>,
) -> tokio::io::Result<()> {
unimplemented!("Strategy::candlestick")
}
}
-100
View File
@@ -1,100 +0,0 @@
use crate::{
general::{EventLog, Signal},
terminal::MarketItem,
units::Direction,
};
use hypersdk::hypercore::{Candle, CandleInterval, Subscription};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StrategyManifest {
pub name: String,
pub description: String,
pub author: String,
pub version: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct RiskManifest {
pub name: String,
pub description: String,
pub author: String,
pub version: String,
pub max_loss: u8,
pub cooldown: CandleInterval,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum StrategyMessage {
Log(EventLog),
GetWatchList,
GetCandlestick {
symbol: String,
interval: CandleInterval,
count: u32,
},
Subscribe(Subscription),
Unsubscribe(Subscription),
UnsubscribeAll,
Signal(StrategySignal),
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum RiskMessage {
Log(EventLog),
GetWatchList,
Approve(Signal),
Reject { reason: String },
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum StrategyEngineMessage {
Initialize,
WatchList(Vec<MarketItem>),
Command {
command: String,
args: Vec<String>,
},
CandleUpdate {
symbol: String,
interval: CandleInterval,
candle: Candle,
},
Candlestick {
symbol: String,
interval: CandleInterval,
candles: Vec<Candle>,
},
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum RiskEngineMessage {
Initialize,
WatchList(Vec<MarketItem>),
Command { command: String, args: Vec<String> },
Signal(StrategySignal),
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StrategySignal {
pub symbol: String,
pub side: Direction,
pub confidence: f32,
pub price: Option<f64>,
}
+54
View File
@@ -0,0 +1,54 @@
use crate::{
general::{EventLog, Signal},
terminal::MarketItem,
};
use hypersdk::hypercore::{Candle, CandleInterval, Incoming, Subscription};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StrategyManifest {
pub name: String,
pub description: String,
pub author: String,
pub version: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum StrategyMessage {
Log(EventLog),
GetWatchList,
GetCandlestick {
symbol: String,
interval: CandleInterval,
count: u32,
},
Subscribe(Subscription),
Unsubscribe(Subscription),
UnsubscribeAll,
Signal(Signal),
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum StrategyEngineMessage {
Initialize,
WatchList(Vec<MarketItem>),
Command {
command: String,
args: Vec<String>,
},
Incoming(Incoming),
Candlestick {
symbol: String,
interval: CandleInterval,
candles: Vec<Candle>,
},
}
+6 -5
View File
@@ -1,9 +1,11 @@
use crate::{ use crate::{
general::{EventLog, MarketTrend, Position, Signal}, general::{EventLog, MarketTrend, Position, Signal},
plugin::{RiskManifest, StrategyManifest}, strategy::StrategyManifest,
units::{Symbol, USD, Volatility}, units::{Symbol, USD, Volatility},
}; };
use hypersdk::hypercore::CandleInterval; use hypersdk::{Decimal, hypercore::CandleInterval};
pub type SignalStatus = Result<Signal, Signal>;
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum TerminalServerMessage { pub enum TerminalServerMessage {
@@ -17,7 +19,7 @@ pub enum TerminalServerMessage {
StrategyUpdated(Strategy), StrategyUpdated(Strategy),
// Signals // Signals
SignalsUpdated(Vec<Signal>), SignalsUpdated(Vec<SignalStatus>),
// Inspector // Inspector
Inspect(InspectTarget), Inspect(InspectTarget),
@@ -39,7 +41,7 @@ pub enum TerminalClientMessage {
pub struct MarketItem { pub struct MarketItem {
pub symbol: Symbol, pub symbol: Symbol,
pub price: USD, pub price: USD,
pub trend: f64, pub trend: Decimal,
pub volume_24h: USD, pub volume_24h: USD,
} }
@@ -108,7 +110,6 @@ pub enum ItemState {
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Strategy { pub struct Strategy {
pub strategy: StrategyManifest, pub strategy: StrategyManifest,
pub risk: RiskManifest,
pub mode: Mode, pub mode: Mode,
pub state: ItemState, pub state: ItemState,
+6 -4
View File
@@ -1,8 +1,10 @@
use rust_decimal::Decimal;
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Symbol(pub String); pub struct Symbol(pub String);
#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
pub struct USD(pub f64); pub struct USD(pub Decimal);
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum Direction { pub enum Direction {
@@ -25,10 +27,10 @@ impl std::fmt::Display for Symbol {
impl std::fmt::Display for USD { impl std::fmt::Display for USD {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
if self.0 > 0.0 { if self.0.is_sign_positive() {
write!(f, "\x1b[32m${}\x1b[0m", format_f64(self.0)) write!(f, "\x1b[32m${}\x1b[0m", format_f64(self.0.as_f64()))
} else { } else {
write!(f, "\x1b[31m${}\x1b[0m", format_f64(self.0)) write!(f, "\x1b[31m${}\x1b[0m", format_f64(self.0.as_f64()))
} }
} }
} }
+3 -29
View File
@@ -47,7 +47,7 @@ impl Engine {
set_cfg!(id, { set_cfg!(id, {
let id: String = id; let id: String = id;
if !crate::store::pulse_plugin(&id)? if !crate::store::pulse_strategy(&id)?
.join("strategy.toml") .join("strategy.toml")
.exists() .exists()
{ {
@@ -69,32 +69,6 @@ impl Engine {
.await?; .await?;
} }
"risk" | "rs" => {
set_cfg!(id, {
let id: String = id;
if !crate::store::pulse_plugin(&id)?
.join("risk.toml")
.exists()
{
return self
.terminal_server
.error(
"config::set::risk",
&format!("Non existent risk `{id}`"),
)
.await;
}
self.strategy.reload_risk(id.as_str()).await?;
self.config.lock().await.risk = id;
});
self.terminal_server
.info("config::set", "risk set successfully, use `config save` to persist changes")
.await?;
}
"cooldown" | "cool" | "cd" => { "cooldown" | "cool" | "cd" => {
set_cfg!(cooldown, { set_cfg!(cooldown, {
self.config.lock().await.cooldown = cooldown; self.config.lock().await.cooldown = cooldown;
@@ -102,7 +76,7 @@ impl Engine {
self.terminal_server.broadcast( self.terminal_server.broadcast(
pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: self.strategy.strategy.manifest.lock().await.clone(), strategy: self.strategy.strategy.manifest.lock().await.clone(),
risk: self.strategy.risk.manifest.lock().await.clone(),
mode: Mode::Auto, mode: Mode::Auto,
state: ItemState::Running, state: ItemState::Running,
cooldown, cooldown,
@@ -114,7 +88,7 @@ impl Engine {
_ => { _ => {
self.terminal_server self.terminal_server
.error("config::set", "Invalid usage, available options: watchlist, strategy, risk, cooldown") .error("config::set", "Invalid usage, available options: watchlist, strategy, cooldown")
.await?; .await?;
} }
} }
+98
View File
@@ -0,0 +1,98 @@
use hypersdk::hypercore::{self, BatchOrder, OrderRequest, OrderTypePlacement, TimeInForce};
use pulse_sdk::prelude::*;
use crate::engine::Engine;
impl Engine {
pub async fn execute_signal(&self, signal: &Signal) -> tokio::io::Result<()> {
let client = hypercore::mainnet();
let accounts = self.accounts.lock().await;
if let Some(acc) = accounts.get_active() {
let Some(asset_id) = self
.watch_list
.lock()
.await
.name_to_index
.get(&signal.symbol)
.cloned()
else {
self.terminal_server
.error(
"self::order",
&format!(
"Invalid Symbol: {:?}, Unable to get asset id",
signal.symbol
),
)
.await?;
return Ok(());
};
let order = BatchOrder {
orders: vec![
OrderRequest {
asset: asset_id,
is_buy: matches!(signal.kind, Direction::Buy),
limit_px: signal.price.0,
sz: signal.size,
reduce_only: false,
order_type: OrderTypePlacement::Limit {
tif: TimeInForce::Gtc,
},
cloid: Default::default(),
},
OrderRequest {
asset: asset_id,
is_buy: matches!(signal.kind, Direction::Buy),
limit_px: signal.price.0,
sz: signal.size,
reduce_only: true,
order_type: OrderTypePlacement::Trigger {
is_market: true,
trigger_px: signal.take_profit.0,
tpsl: hypercore::TpSl::Tp,
},
cloid: Default::default(),
},
OrderRequest {
asset: asset_id,
is_buy: matches!(signal.kind, Direction::Buy),
limit_px: signal.price.0,
sz: signal.size,
reduce_only: true,
order_type: OrderTypePlacement::Trigger {
is_market: true,
trigger_px: signal.stop_loss.0,
tpsl: hypercore::TpSl::Sl,
},
cloid: Default::default(),
},
],
grouping: hypercore::OrderGrouping::Na,
builder: None,
};
let nonce = chrono::Utc::now().timestamp_millis() as u64;
match client
.place(&acc.private_key.0, order, nonce, None, None)
.await
{
Ok(_) => {}
Err(e) => {
self.terminal_server
.error("self::order", &e.to_string())
.await?;
}
}
} else {
self.terminal_server
.error("Engine::order", "Unable to get active account")
.await?;
}
Ok(())
}
}
+34 -29
View File
@@ -1,29 +1,37 @@
pub mod command; pub mod command;
pub mod plugin; pub mod execution;
pub mod strategy;
use crate::{ use crate::{
engine::plugin::StrategyEngine, engine::strategy::StrategyEngine,
store::{accounts::AccountList, config::Config}, store::{accounts::AccountList, config::Config},
terminal::TerminalServer, terminal::TerminalServer,
}; };
use pulse_sdk::prelude::*; use pulse_sdk::prelude::*;
use std::sync::Arc; use std::{collections::HashMap, sync::Arc};
use tokio::{sync::Mutex, task::JoinHandle}; use tokio::{sync::Mutex, task::JoinHandle};
#[derive(Debug, Clone)]
pub struct WatchList {
pub name_to_index: HashMap<String, usize>,
pub items: Vec<MarketItem>,
}
#[derive(Clone)] #[derive(Clone)]
pub struct Engine { pub struct Engine {
pub terminal_server: Arc<TerminalServer>, pub terminal_server: Arc<TerminalServer>,
pub strategy: Arc<StrategyEngine>, pub strategy: Arc<StrategyEngine>,
pub config: Arc<Mutex<Config>>, pub config: Arc<Mutex<Config>>,
pub accounts: Arc<Mutex<AccountList>>, pub accounts: Arc<Mutex<AccountList>>,
pub watch_list: Arc<Mutex<Vec<MarketItem>>>, pub watch_list: Arc<Mutex<WatchList>>,
pub signals: Arc<Mutex<Vec<SignalStatus>>>,
} }
impl Engine { impl Engine {
pub async fn new() -> tokio::io::Result<Arc<Self>> { pub async fn new() -> tokio::io::Result<Arc<Self>> {
let config = Config::new().await?; let config = Config::new().await?;
let strategy = StrategyEngine::new(&config.strategy, &config.risk).await?; let strategy = StrategyEngine::new(&config.strategy).await?;
let accounts = Arc::new(Mutex::new(AccountList::new().await?)); let accounts = Arc::new(Mutex::new(AccountList::new().await?));
let config = Arc::new(Mutex::new(config)); let config = Arc::new(Mutex::new(config));
@@ -33,7 +41,11 @@ impl Engine {
strategy: strategy.initialize(engine.clone()), strategy: strategy.initialize(engine.clone()),
config, config,
accounts, accounts,
watch_list: Arc::new(Mutex::new(Vec::new())), watch_list: Arc::new(Mutex::new(WatchList {
name_to_index: HashMap::new(),
items: Vec::new(),
})),
signals: Arc::new(Mutex::new(Vec::new())),
})) }))
} }
@@ -59,11 +71,7 @@ impl Engine {
if let Err(error) = self if let Err(error) = self
.terminal_server .terminal_server
.broadcast( .broadcast(TerminalServerMessage::WatchListUpdated(watch_list.items))
pulse_sdk::terminal::TerminalServerMessage::WatchListUpdated(
watch_list,
),
)
.await .await
{ {
self.terminal_server self.terminal_server
@@ -91,24 +99,21 @@ impl Engine {
match client.clearinghouse_state(acc.address, None).await { match client.clearinghouse_state(acc.address, None).await {
Ok(state) => { Ok(state) => {
self.terminal_server self.terminal_server
.broadcast( .broadcast(TerminalServerMessage::PositionsUpdated(
pulse_sdk::terminal::TerminalServerMessage::PositionsUpdated( state
state .asset_positions
.asset_positions .into_iter()
.into_iter() .map(|position| Position {
.map(|position| Position { symbol: Symbol(position.position.coin),
symbol: Symbol(position.position.coin), size: position.position.szi.as_f64(),
size: position.position.szi.as_f64(), entry_price: USD(position
entry_price: USD(position .position
.position .entry_px
.entry_px .unwrap_or_default()),
.map(|px| px.as_f64()) profit: USD(position.position.unrealized_pnl),
.unwrap_or(0.0)), })
profit: USD(position.position.unrealized_pnl.as_f64()), .collect(),
}) ))
.collect(),
),
)
.await?; .await?;
} }
Err(e) => { Err(e) => {
@@ -15,12 +15,11 @@ use tokio::{
use crate::{ use crate::{
engine::Engine, engine::Engine,
store::{plugin::Plugin, pulse_plugin}, store::{pulse_strategy, strategy::StrategyChild},
}; };
pub struct StrategyEngine { pub struct StrategyEngine {
pub strategy: Arc<Plugin<StrategyEngineMessage, StrategyMessage, StrategyManifest>>, pub strategy: Arc<StrategyChild>,
pub risk: Arc<Plugin<RiskEngineMessage, RiskMessage, RiskManifest>>,
pub engine: Weak<Engine>, pub engine: Weak<Engine>,
pub ws: WebSocket, pub ws: WebSocket,
@@ -28,25 +27,17 @@ pub struct StrategyEngine {
} }
impl StrategyEngine { impl StrategyEngine {
pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result<Self> { pub async fn new(strategy_id: &str) -> tokio::io::Result<Self> {
let strategy = pulse_plugin(strategy_id)?; let strategy = pulse_strategy(strategy_id)?;
let risk = pulse_plugin(risk_id)?;
let (strategy, strategy_manifest) = get_manifest_plugin_pair( let (strategy, strategy_manifest) = get_manifest(
&strategy, &strategy,
&strategy.join("strategy.bash"), &strategy.join("strategy.bash"),
&fs::read(strategy.join("strategy.toml")).await?, &fs::read(strategy.join("strategy.toml")).await?,
)?; )?;
let (risk, risk_manifest) = get_manifest_plugin_pair(
&risk,
&risk.join("risk.bash"),
&fs::read(risk.join("risk.toml")).await?,
)?;
Ok(Self { Ok(Self {
strategy: Arc::new(Plugin::new(strategy, strategy_manifest)), strategy: Arc::new(StrategyChild::new(strategy, strategy_manifest)),
risk: Arc::new(Plugin::new(risk, risk_manifest)),
engine: Weak::new(), engine: Weak::new(),
ws: hypercore::mainnet_ws(), ws: hypercore::mainnet_ws(),
subscriptions: Mutex::new(HashSet::new()), subscriptions: Mutex::new(HashSet::new()),
@@ -76,7 +67,7 @@ impl StrategyEngine {
Some(StrategyMessage::GetWatchList) => { Some(StrategyMessage::GetWatchList) => {
self.strategy self.strategy
.send(&StrategyEngineMessage::WatchList( .send(&StrategyEngineMessage::WatchList(
engine.watch_list.lock().await.clone(), engine.watch_list.lock().await.clone().items,
)) ))
.await?; .await?;
} }
@@ -87,7 +78,24 @@ impl StrategyEngine {
} }
Some(StrategyMessage::Signal(signal)) => { Some(StrategyMessage::Signal(signal)) => {
self.risk.send(&RiskEngineMessage::Signal(signal)).await?; match engine.execute_signal(&signal).await {
Ok(_) => engine.signals.lock().await.push(Ok(signal)),
Err(e) => {
engine.signals.lock().await.push(Err(signal));
engine
.terminal_server
.error("signal", &format!("Failed to execute signal: {e}"))
.await?
}
}
engine
.terminal_server
.broadcast(TerminalServerMessage::SignalsUpdated(
engine.signals.lock().await.clone(),
))
.await?;
} }
Some(StrategyMessage::Subscribe(subscription)) => { Some(StrategyMessage::Subscribe(subscription)) => {
@@ -153,54 +161,19 @@ impl StrategyEngine {
} }
} }
pub async fn run_risk(&self) -> tokio::io::Result<()> {
let engine = self
.engine
.upgrade()
.expect("Failed to upgrade engine (StrategyEngine)");
self.risk.send(&RiskEngineMessage::Initialize).await?;
loop {
match self.risk.recv().await? {
None => {}
Some(RiskMessage::GetWatchList) => {
self.risk
.send(&&RiskEngineMessage::WatchList(
engine.watch_list.lock().await.clone(),
))
.await?;
}
Some(RiskMessage::Log(mut log)) => {
log.name.insert_str(0, "risk::");
engine.terminal_server.log_raw(log).await?;
}
Some(RiskMessage::Approve(signal)) => {}
Some(RiskMessage::Reject { reason }) => {}
}
}
}
pub async fn spawn(self: &Arc<Self>) { pub async fn spawn(self: &Arc<Self>) {
let engine = self.clone(); let engine = self.clone();
tokio::spawn(async move { engine.run_strategy().await }); tokio::spawn(async move { engine.run_strategy().await });
let engine = self.clone();
tokio::spawn(async move { engine.run_risk().await });
} }
pub async fn reload_strategy(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> { pub async fn reload_strategy(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> {
let plugin = pulse_plugin(id)?; let strategy = pulse_strategy(id)?;
let (child, manifest) = get_manifest_plugin_pair( let (child, manifest) = get_manifest(
&plugin, &strategy,
&plugin.join("strategy.bash"), &strategy.join("strategy.bash"),
&fs::read(plugin.join("strategy.toml")).await?, &fs::read(strategy.join("strategy.toml")).await?,
)?; )?;
self.strategy.reload(child, manifest).await?; self.strategy.reload(child, manifest).await?;
@@ -210,34 +183,17 @@ impl StrategyEngine {
Ok(()) Ok(())
} }
pub async fn reload_risk(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> {
let plugin = pulse_plugin(id)?;
let (child, manifest) = get_manifest_plugin_pair(
&plugin,
&plugin.join("risk.bash"),
&fs::read(plugin.join("risk.toml")).await?,
)?;
self.risk.reload(child, manifest).await?;
let engine = self.clone();
tokio::spawn(async move { engine.run_risk().await });
Ok(())
}
} }
pub fn get_manifest_plugin_pair<'de, M: serde::Deserialize<'de>>( fn get_manifest<'de, M: serde::Deserialize<'de>>(
plugin_dir: &PathBuf, strategy_dir: &PathBuf,
plugin_path: &PathBuf, strategy_path: &PathBuf,
manifest: &'de [u8], manifest: &'de [u8],
) -> tokio::io::Result<(Child, M)> { ) -> tokio::io::Result<(Child, M)> {
Ok(( Ok((
Command::new("bash") Command::new("bash")
.arg(plugin_path) .arg(strategy_path)
.current_dir(plugin_dir) .current_dir(strategy_dir)
.stdin(Stdio::piped()) .stdin(Stdio::piped())
.stdout(Stdio::piped()) .stdout(Stdio::piped())
.stderr(Stdio::inherit()) .stderr(Stdio::inherit())
+36 -23
View File
@@ -1,20 +1,23 @@
use rust_decimal::Decimal;
use pulse_sdk::prelude::*; use pulse_sdk::prelude::*;
use serde_json::Value; use serde_json::Value;
use std::collections::HashMap; use std::collections::HashMap;
fn number(value: &Value, field: &str) -> Result<f64, String> { use crate::engine::WatchList;
fn number(value: &Value, field: &str) -> Result<Decimal, String> {
let raw = value[field] let raw = value[field]
.as_str() .as_str()
.ok_or_else(|| format!("asset context field {field} must be a string"))?; .ok_or_else(|| format!("asset context field {field} must be a string"))?;
raw.parse::<f64>() raw.parse::<Decimal>()
.map_err(|error| format!("could not parse asset context field {field} ({raw}): {error}")) .map_err(|error| format!("could not parse asset context field {field} ({raw}): {error}"))
} }
pub async fn fetch_watch_list( pub async fn fetch_watch_list(
client: &hypersdk::hypercore::HttpClient, client: &hypersdk::hypercore::HttpClient,
symbols: &[String], symbols: &[String],
) -> Result<Vec<MarketItem>, String> { ) -> Result<WatchList, String> {
let response = client let response = client
.meta_and_asset_ctxs(None) .meta_and_asset_ctxs(None)
.await .await
@@ -34,6 +37,7 @@ pub async fn fetch_watch_list(
let universe = response[0]["universe"] let universe = response[0]["universe"]
.as_array() .as_array()
.ok_or("metaAndAssetCtxs response is missing meta.universe")?; .ok_or("metaAndAssetCtxs response is missing meta.universe")?;
let contexts = response[1] let contexts = response[1]
.as_array() .as_array()
.ok_or("metaAndAssetCtxs response contexts must be an array")?; .ok_or("metaAndAssetCtxs response contexts must be an array")?;
@@ -46,37 +50,46 @@ pub async fn fetch_watch_list(
)); ));
} }
let mut by_symbol = HashMap::with_capacity(universe.len()); // Build symbol -> perp asset index
let mut name_to_index = HashMap::with_capacity(universe.len());
for (meta, context) in universe.iter().zip(contexts) { for (index, meta) in universe.iter().enumerate() {
let symbol = meta["name"] let symbol = meta["name"]
.as_str() .as_str()
.ok_or("instrument metadata is missing a name")?; .ok_or("instrument metadata is missing a name")?;
name_to_index.insert(symbol.to_owned(), index);
}
// Now only process requested symbols
let mut items = Vec::with_capacity(symbols.len());
for symbol in symbols {
let index = *name_to_index
.get(symbol)
.ok_or_else(|| format!("{symbol} is not in the Hyperliquid perpetual universe"))?;
let context = &contexts[index];
let price = number(context, "markPx")?; let price = number(context, "markPx")?;
let previous_day_price = number(context, "prevDayPx")?; let previous_day_price = number(context, "prevDayPx")?;
let volume_24h = number(context, "dayNtlVlm")?; let volume_24h = number(context, "dayNtlVlm")?;
if previous_day_price <= 0.0 { if previous_day_price.is_zero() || previous_day_price.is_sign_negative() {
return Err(format!("{symbol} has an invalid previous-day price")); return Err(format!("{symbol} has an invalid previous-day price"));
} }
by_symbol.insert( items.push(MarketItem {
symbol, symbol: Symbol(symbol.clone()),
MarketItem { price: USD(price),
symbol: Symbol(symbol.to_owned()), volume_24h: USD(volume_24h),
price: USD(price), trend: ((price / previous_day_price) - <Decimal as From<i32>>::from(1))
volume_24h: USD(volume_24h), * <Decimal as From<i32>>::from(100),
trend: ((price / previous_day_price) - 1.0) * 100.0, });
},
);
} }
symbols Ok(WatchList {
.iter() items,
.map(|symbol| { name_to_index,
by_symbol })
.remove(symbol.as_str())
.ok_or_else(|| format!("{symbol} is not in the Hyperliquid perpetual universe"))
})
.collect()
} }
-2
View File
@@ -4,7 +4,6 @@ use hypersdk::hypercore::CandleInterval;
pub struct Config { pub struct Config {
pub watchlist: Vec<String>, pub watchlist: Vec<String>,
pub strategy: String, pub strategy: String,
pub risk: String,
pub cooldown: CandleInterval, pub cooldown: CandleInterval,
} }
@@ -13,7 +12,6 @@ impl Default for Config {
Self { Self {
watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()], watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()],
strategy: String::new(), strategy: String::new(),
risk: String::new(),
cooldown: CandleInterval::ThirtyMinutes, cooldown: CandleInterval::ThirtyMinutes,
} }
} }
+5 -5
View File
@@ -2,7 +2,7 @@ use std::path::PathBuf;
pub mod accounts; pub mod accounts;
pub mod config; pub mod config;
pub mod plugin; pub mod strategy;
pub fn home_dir() -> tokio::io::Result<PathBuf> { pub fn home_dir() -> tokio::io::Result<PathBuf> {
std::env::home_dir().ok_or_else(|| { std::env::home_dir().ok_or_else(|| {
@@ -14,12 +14,12 @@ pub fn pulse_directory() -> tokio::io::Result<PathBuf> {
Ok(home_dir()?.join(".config").join("pulse-trader")) Ok(home_dir()?.join(".config").join("pulse-trader"))
} }
pub fn pulse_plugins_directory() -> tokio::io::Result<PathBuf> { pub fn pulse_strategies_directory() -> tokio::io::Result<PathBuf> {
Ok(pulse_directory()?.join("plugins")) Ok(pulse_directory()?.join("strategies"))
} }
pub fn pulse_plugin(id: &str) -> tokio::io::Result<PathBuf> { pub fn pulse_strategy(id: &str) -> tokio::io::Result<PathBuf> {
Ok(pulse_directory()?.join("plugins").join(id)) Ok(pulse_directory()?.join("strategies").join(id))
} }
pub fn pulse_config_file() -> tokio::io::Result<PathBuf> { pub fn pulse_config_file() -> tokio::io::Result<PathBuf> {
@@ -1,33 +1,29 @@
use std::marker::PhantomData; use pulse_sdk::{
map_postcard_err,
use pulse_sdk::map_postcard_err; strategy::{StrategyEngineMessage, StrategyManifest, StrategyMessage},
use serde::{Deserialize, Serialize}; };
use tokio::{ use tokio::{
io::{AsyncReadExt, AsyncWriteExt}, io::{AsyncReadExt, AsyncWriteExt},
process::{Child, ChildStdout}, process::{Child, ChildStdout},
sync::Mutex, sync::Mutex,
}; };
#[derive(Debug)] #[derive(Debug)]
pub struct Plugin<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> { pub struct StrategyChild {
pub manifest: Mutex<M>, pub manifest: Mutex<StrategyManifest>,
pub stdout: Mutex<ChildStdout>, pub stdout: Mutex<ChildStdout>,
pub process: Mutex<Child>, pub process: Mutex<Child>,
pub _p: (PhantomData<S>, PhantomData<R>),
} }
impl<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> Plugin<S, R, M> { impl StrategyChild {
pub fn new(mut child: Child, manifest: M) -> Self { pub fn new(mut child: Child, manifest: StrategyManifest) -> Self {
Self { Self {
stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")), stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")),
process: Mutex::new(child), process: Mutex::new(child),
manifest: Mutex::new(manifest), manifest: Mutex::new(manifest),
_p: (PhantomData, PhantomData),
} }
} }
pub async fn recv(&self) -> tokio::io::Result<Option<R>> { pub async fn recv(&self) -> tokio::io::Result<Option<StrategyMessage>> {
let mut stdout = self.stdout.lock().await; let mut stdout = self.stdout.lock().await;
let mut len_buf = [0u8; size_of::<usize>()]; let mut len_buf = [0u8; size_of::<usize>()];
@@ -46,7 +42,7 @@ impl<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> P
Ok(Some(map_postcard_err(postcard::from_bytes(&buffer))?)) Ok(Some(map_postcard_err(postcard::from_bytes(&buffer))?))
} }
pub async fn send(&self, msg: &S) -> tokio::io::Result<()> { pub async fn send(&self, msg: &StrategyEngineMessage) -> tokio::io::Result<()> {
self.send_raw(&map_postcard_err(postcard::to_allocvec(msg))?) self.send_raw(&map_postcard_err(postcard::to_allocvec(msg))?)
.await .await
} }
@@ -62,7 +58,11 @@ impl<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> P
Ok(()) Ok(())
} }
pub async fn reload(&self, mut child: Child, manifest: M) -> tokio::io::Result<()> { pub async fn reload(
&self,
mut child: Child,
manifest: StrategyManifest,
) -> tokio::io::Result<()> {
let mut process = self.process.lock().await; let mut process = self.process.lock().await;
process.kill().await?; process.kill().await?;
+2 -3
View File
@@ -66,7 +66,6 @@ impl TerminalServer {
id, id,
pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: engine.strategy.strategy.manifest.lock().await.clone(), strategy: engine.strategy.strategy.manifest.lock().await.clone(),
risk: engine.strategy.risk.manifest.lock().await.clone(),
mode: Mode::Auto, mode: Mode::Auto,
state: ItemState::Running, state: ItemState::Running,
cooldown: engine.config.lock().await.cooldown, cooldown: engine.config.lock().await.cooldown,
@@ -164,8 +163,8 @@ impl TerminalServer {
} }
pub async fn send_to_client(client: &mut OwnedWriteHalf, msg: &[u8]) -> tokio::io::Result<()> { pub async fn send_to_client(client: &mut OwnedWriteHalf, msg: &[u8]) -> tokio::io::Result<()> {
client.write(&msg.len().to_le_bytes()).await?; client.write_all(&msg.len().to_le_bytes()).await?;
client.write(msg).await?; client.write_all(msg).await?;
client.flush().await?; client.flush().await?;
Ok(()) Ok(())
+34 -39
View File
@@ -58,17 +58,37 @@ impl Formatted for EventLog {
} }
} }
impl Formatted for Signal { impl Formatted for SignalStatus {
fn get_formatted(&self) -> Vec<String> { fn get_formatted(&self) -> Vec<String> {
vec![ match self {
if matches!(self.kind, Direction::Buy) { Ok(signal) => {
format!("\x1b[32m{}\x1b[0m", self.kind) vec![
} else { format!("\x1b[33mOK\x1b[0m"),
format!("\x1b[31m{}\x1b[0m", self.kind) if matches!(signal.kind, Direction::Buy) {
}, format!("\x1b[32mBUY\x1b[0m")
format!("\x1b[35m{}\x1b[0m", self.symbol), } else {
self.price.to_string(), format!("\x1b[31mSELL\x1b[0m")
] },
format!("\x1b[35m{}\x1b[0m", signal.symbol),
signal.price.to_string(),
format!("\x1b[33mAPR {}\x1b[0m", signal.confidence),
]
}
Err(signal) => {
vec![
format!("\x1b[31mERR\x1b[0m"),
if matches!(signal.kind, Direction::Buy) {
format!("\x1b[32mBUY\x1b[0m")
} else {
format!("\x1b[31mSELL\x1b[0m")
},
format!("\x1b[35m{}\x1b[0m", signal.symbol),
signal.price.to_string(),
format!("\x1b[33mAPR {}\x1b[0m", signal.confidence),
]
}
}
} }
} }
@@ -85,7 +105,7 @@ impl Formatted for MarketItem {
} else { } else {
"\x1b[31m▼\x1b[0m" "\x1b[31m▼\x1b[0m"
}, },
format_f64(self.trend.abs()) format_f64(self.trend.as_f64().abs())
), ),
] ]
} }
@@ -232,38 +252,13 @@ impl Formatted for Strategy {
vec![ vec![
Triple( Triple(
"\x1b[2mStrategy\x1b[0m", "\x1b[2mStrategy\x1b[0m",
"\x1b[2mRisk\x1b[0m", "\x1b[2mState\x1b[0m",
"\x1b[2mStrat Ver\x1b[0m", "\x1b[2mMode\x1b[0m",
), ),
Triple( Triple(
&format!("\x1b[97m{}\x1b[0m", self.strategy.name), &format!("\x1b[97m{}\x1b[0m", self.strategy.name),
&format!("\x1b[93m{}\x1b[0m", self.risk.name),
&format!("\x1b[90m{}\x1b[0m", self.strategy.version),
),
Triple("", "", ""),
Triple(
"\x1b[2mMode\x1b[0m",
"\x1b[2mState\x1b[0m",
"\x1b[2mRisk Ver\x1b[0m",
),
Triple(
&self.mode.to_string(),
&self.state.to_string(), &self.state.to_string(),
&format!("\x1b[90m{}\x1b[0m", self.risk.version), &self.mode.to_string(),
),
Triple("", "", ""),
Triple(
"\x1b[2mCooldown\x1b[0m",
"\x1b[2mStrat Author\x1b[0m",
"\x1b[2mMax loss\x1b[0m",
),
Triple(
&format!(
"\x1b[96m{}\x1b[0m (\x1b[90m{} rec\x1b[0m)",
self.cooldown, self.risk.cooldown
),
&format!("\x1b[96m{:?}\x1b[0m", self.strategy.author),
&format!("\x1b[93m{}%\x1b[0m", self.risk.max_loss),
), ),
] ]
.get_formatted() .get_formatted()
+1 -1
View File
@@ -35,7 +35,7 @@ pub struct PulseTradeApp {
watch_list: State<Vec<MarketItem>>, watch_list: State<Vec<MarketItem>>,
active_positions: State<Vec<Position>>, active_positions: State<Vec<Position>>,
logs: State<Vec<EventLog>>, logs: State<Vec<EventLog>>,
signals: State<Vec<Signal>>, signals: State<Vec<SignalStatus>>,
inspect: State<InspectTarget>, inspect: State<InspectTarget>,
} }
+56 -50
View File
@@ -31,8 +31,8 @@ impl TerminalClient {
message: pulse_sdk::terminal::TerminalClientMessage, message: pulse_sdk::terminal::TerminalClientMessage,
) -> tokio::io::Result<()> { ) -> tokio::io::Result<()> {
let msg = map_postcard_err(postcard::to_allocvec(&message))?; let msg = map_postcard_err(postcard::to_allocvec(&message))?;
self.writer.write(&msg.len().to_le_bytes()).await?; self.writer.write_all(&msg.len().to_le_bytes()).await?;
self.writer.write(&msg).await?; self.writer.write_all(&msg).await?;
self.writer.flush().await?; self.writer.flush().await?;
Ok(()) Ok(())
@@ -53,16 +53,22 @@ impl TerminalClient {
let status = app.status.clone(); let status = app.status.clone();
let inspect = app.inspect.clone(); let inspect = app.inspect.clone();
tokio::spawn(Self::run_client( tokio::spawn(async {
reader, if let Err(v) = Self::run_client(
watch_list, reader,
active_positions, watch_list,
logs, active_positions,
signals, logs,
market_overview, signals,
status, market_overview,
inspect, status,
)); inspect,
)
.await
{
panic!("{v}")
}
});
app.sock = Some(self); app.sock = Some(self);
@@ -75,60 +81,60 @@ impl TerminalClient {
watch_list: State<Vec<MarketItem>>, watch_list: State<Vec<MarketItem>>,
active_positions: State<Vec<Position>>, active_positions: State<Vec<Position>>,
logs: State<Vec<EventLog>>, logs: State<Vec<EventLog>>,
signals: State<Vec<Signal>>, signals: State<Vec<SignalStatus>>,
market_overview: State<Option<Strategy>>, market_overview: State<Option<Strategy>>,
status: State<Option<Status>>, status: State<Option<Status>>,
inspect: State<InspectTarget>, inspect: State<InspectTarget>,
) -> tokio::io::Result<()> { ) -> tokio::io::Result<()> {
let mut len_buf = [0u8; size_of::<usize>()]; loop {
reader let mut len_buf = [0u8; size_of::<usize>()];
.read_exact(&mut len_buf) reader
.await .read_exact(&mut len_buf)
.expect("Failed to get header length"); .await
.expect("Failed to get header length");
let len = usize::from_le_bytes(len_buf); let len = usize::from_le_bytes(len_buf);
let mut buffer = vec![0u8; len]; let mut buffer = vec![0u8; len];
reader reader
.read_exact(&mut buffer) .read_exact(&mut buffer)
.await .await
.expect("Failed to read socket"); .expect("Failed to read socket");
match map_postcard_err(postcard::from_bytes(&buffer))? { match map_postcard_err(postcard::from_bytes(&buffer))? {
TerminalServerMessage::WatchListUpdated(v) => { TerminalServerMessage::WatchListUpdated(v) => {
*watch_list.lock().await = v; *watch_list.lock().await = v;
} }
TerminalServerMessage::PositionsUpdated(v) => { TerminalServerMessage::PositionsUpdated(v) => {
*active_positions.lock().await = v; *active_positions.lock().await = v;
} }
TerminalServerMessage::StrategyUpdated(v) => { TerminalServerMessage::StrategyUpdated(v) => {
*market_overview.lock().await = Some(v); *market_overview.lock().await = Some(v);
} }
TerminalServerMessage::SignalsUpdated(v) => { TerminalServerMessage::SignalsUpdated(v) => {
*signals.lock().await = v; *signals.lock().await = v;
} }
TerminalServerMessage::Inspect(v) => { TerminalServerMessage::Inspect(v) => {
*inspect.lock().await = v; *inspect.lock().await = v;
} }
TerminalServerMessage::StatusUpdated(v) => { TerminalServerMessage::StatusUpdated(v) => {
*status.lock().await = Some(v); *status.lock().await = Some(v);
} }
TerminalServerMessage::SetLogs(v) => { TerminalServerMessage::SetLogs(v) => {
*logs.lock().await = v; *logs.lock().await = v;
} }
TerminalServerMessage::AddLog(v) => { TerminalServerMessage::AddLog(v) => {
logs.lock().await.push(v); logs.lock().await.push(v);
}
} }
} }
Ok(())
} }
} }