Merge pull request #15 from selimaj-dev/strategy-engine

Strategy engine
This commit is contained in:
2026-07-27 03:49:53 +02:00
committed by GitHub
19 changed files with 1056 additions and 245 deletions
Generated
+3
View File
@@ -3882,6 +3882,7 @@ dependencies = [
name = "pulse-trader" name = "pulse-trader"
version = "0.1.0-alpha.0" version = "0.1.0-alpha.0"
dependencies = [ dependencies = [
"anyhow",
"chrono", "chrono",
"crossterm", "crossterm",
"hypersdk", "hypersdk",
@@ -3906,6 +3907,7 @@ dependencies = [
name = "pulse-wire" name = "pulse-wire"
version = "0.1.0-alpha.0" version = "0.1.0-alpha.0"
dependencies = [ dependencies = [
"hypersdk",
"pulse-macros", "pulse-macros",
"serde", "serde",
] ]
@@ -5200,6 +5202,7 @@ dependencies = [
"libc", "libc",
"mio", "mio",
"pin-project-lite", "pin-project-lite",
"signal-hook-registry",
"socket2", "socket2",
"tokio-macros", "tokio-macros",
"windows-sys 0.61.2", "windows-sys 0.61.2",
+8 -6
View File
@@ -6,8 +6,10 @@ edition = "2024"
[dependencies] [dependencies]
pulse-ui = { workspace = true } pulse-ui = { workspace = true }
pulse-wire = { workspace = true } pulse-wire = { workspace = true }
crossterm = { workspace = true }
chrono = "0.4.45" hypersdk = { workspace = true }
serde = { workspace = true }
anyhow = { workspace = true }
tokio = { workspace = true, features = [ tokio = { workspace = true, features = [
"rt-multi-thread", "rt-multi-thread",
"macros", "macros",
@@ -15,13 +17,12 @@ tokio = { workspace = true, features = [
"fs", "fs",
"io-util", "io-util",
"time", "time",
"process",
] } ] }
crossterm = { workspace = true } chrono = "0.4.45"
hypersdk = "0.2.14"
serde_json = "1" serde_json = "1"
rand = "0.8.7" rand = "0.8.7"
toml = "1.1.3" toml = "1.1.3"
serde = { workspace = true }
[workspace] [workspace]
members = ["pulse-macros", "pulse-ui", "pulse-wire"] members = ["pulse-macros", "pulse-ui", "pulse-wire"]
@@ -30,10 +31,11 @@ members = ["pulse-macros", "pulse-ui", "pulse-wire"]
pulse-macros = { path = "pulse-macros", version = "0.1.0-alpha.0" } pulse-macros = { path = "pulse-macros", version = "0.1.0-alpha.0" }
pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" } pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" }
pulse-wire = { path = "pulse-wire", version = "0.1.0-alpha.0" } pulse-wire = { path = "pulse-wire", version = "0.1.0-alpha.0" }
hypersdk = "0.2.14"
tokio = "1.52.3" tokio = "1.52.3"
crossterm = "0.29.0" crossterm = "0.29.0"
serde = { version = "1.0.229", features = ["serde_derive"] } serde = { version = "1.0.229", features = ["serde_derive"] }
anyhow = "1.0.104"
[[bin]] [[bin]]
name = "pulse-trader" name = "pulse-trader"
+1
View File
@@ -6,3 +6,4 @@ edition = "2024"
[dependencies] [dependencies]
pulse-macros = { workspace = true } pulse-macros = { workspace = true }
serde = { workspace = true } serde = { workspace = true }
hypersdk = { workspace = true }
+313
View File
@@ -0,0 +1,313 @@
use hypersdk::{
Address, Decimal,
hypercore::{Candle, CandleInterval, Subscription},
};
use crate::PulseWire;
impl PulseWire for Subscription {
fn from_com(com: &mut Vec<u8>) -> Self {
match u8::from_com(com) {
0 => Self::Bbo {
coin: PulseWire::from_com(com),
},
1 => Self::Trades {
coin: PulseWire::from_com(com),
},
2 => Self::L2Book {
coin: PulseWire::from_com(com),
n_sig_figs: PulseWire::from_com(com),
mantissa: PulseWire::from_com(com),
fast: PulseWire::from_com(com),
},
3 => Self::Candle {
coin: PulseWire::from_com(com),
interval: PulseWire::from_com(com),
},
4 => Self::AllMids {
dex: PulseWire::from_com(com),
},
5 => Self::OrderUpdates {
user: PulseWire::from_com(com),
},
6 => Self::UserFills {
user: PulseWire::from_com(com),
},
7 => Self::UserEvents {
user: PulseWire::from_com(com),
},
8 => Self::UserTwapSliceFills {
user: PulseWire::from_com(com),
},
9 => Self::UserTwapHistory {
user: PulseWire::from_com(com),
},
10 => Self::ActiveAssetCtx {
coin: PulseWire::from_com(com),
},
11 => Self::ActiveAssetData {
user: PulseWire::from_com(com),
coin: PulseWire::from_com(com),
},
12 => Self::WebData2 {
user: PulseWire::from_com(com),
dex: PulseWire::from_com(com),
},
13 => Self::ClearinghouseState {
user: PulseWire::from_com(com),
dex: PulseWire::from_com(com),
},
14 => Self::AllDexsClearinghouseState {
user: PulseWire::from_com(com),
},
15 => Self::OpenOrders {
user: PulseWire::from_com(com),
dex: PulseWire::from_com(com),
},
16 => Self::SpotState {
user: PulseWire::from_com(com),
is_portfolio_margin: PulseWire::from_com(com),
},
17 => Self::Notification {
user: PulseWire::from_com(com),
},
18 => Self::WebData3 {
user: PulseWire::from_com(com),
},
19 => Self::TwapStates {
user: PulseWire::from_com(com),
dex: Option::from_com(com),
},
20 => Self::UserFundings {
user: PulseWire::from_com(com),
},
21 => Self::UserNonFundingLedgerUpdates {
user: PulseWire::from_com(com),
},
22 => Self::AllDexsAssetCtxs,
23 => Self::FastAssetCtxs,
24 => Self::OutcomeMetaUpdates,
x => panic!("Invalid Subscription discriminant: {}", x),
}
}
fn to_com(&self) -> Vec<u8> {
let mut com = Vec::new();
match self {
Self::Bbo { coin } => {
com.push(0);
com.extend(coin.to_com());
}
Self::Trades { coin } => {
com.push(1);
com.extend(coin.to_com());
}
Self::L2Book {
coin,
n_sig_figs,
mantissa,
fast,
} => {
com.push(2);
com.extend(coin.to_com());
com.extend(n_sig_figs.to_com());
com.extend(mantissa.to_com());
com.extend(fast.to_com());
}
Self::Candle { coin, interval } => {
com.push(3);
com.extend(coin.to_com());
com.extend(interval.to_com());
}
Self::AllMids { dex } => {
com.push(4);
com.extend(dex.to_com());
}
Self::OrderUpdates { user } => {
com.push(5);
com.extend(user.to_com());
}
Self::UserFills { user } => {
com.push(6);
com.extend(user.to_com());
}
Self::UserEvents { user } => {
com.push(7);
com.extend(user.to_com());
}
Self::UserTwapSliceFills { user } => {
com.push(8);
com.extend(user.to_com());
}
Self::UserTwapHistory { user } => {
com.push(9);
com.extend(user.to_com());
}
Self::ActiveAssetCtx { coin } => {
com.push(10);
com.extend(coin.to_com());
}
Self::ActiveAssetData { user, coin } => {
com.push(11);
com.extend(user.to_com());
com.extend(coin.to_com());
}
Self::WebData2 { user, dex } => {
com.push(12);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::ClearinghouseState { user, dex } => {
com.push(13);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::AllDexsClearinghouseState { user } => {
com.push(14);
com.extend(user.to_com());
}
Self::OpenOrders { user, dex } => {
com.push(15);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::SpotState {
user,
is_portfolio_margin,
} => {
com.push(16);
com.extend(user.to_com());
com.extend(is_portfolio_margin.to_com());
}
Self::Notification { user } => {
com.push(17);
com.extend(user.to_com());
}
Self::WebData3 { user } => {
com.push(18);
com.extend(user.to_com());
}
Self::TwapStates { user, dex } => {
com.push(19);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::UserFundings { user } => {
com.push(20);
com.extend(user.to_com());
}
Self::UserNonFundingLedgerUpdates { user } => {
com.push(21);
com.extend(user.to_com());
}
Self::AllDexsAssetCtxs => {
com.push(22);
}
Self::FastAssetCtxs => {
com.push(23);
}
Self::OutcomeMetaUpdates => {
com.push(24);
}
}
com
}
}
impl PulseWire for Address {
fn from_com(com: &mut Vec<u8>) -> Self {
let bytes: [u8; 20] = com.drain(..20).collect::<Vec<_>>().try_into().unwrap();
Self::from_slice(&bytes)
}
fn to_com(&self) -> Vec<u8> {
self.as_slice().to_vec()
}
}
impl PulseWire for CandleInterval {
fn from_com(com: &mut Vec<u8>) -> Self {
match u8::from_com(com) {
0 => Self::OneMinute,
1 => Self::ThreeMinutes,
2 => Self::FiveMinutes,
3 => Self::FifteenMinutes,
4 => Self::ThirtyMinutes,
5 => Self::OneHour,
6 => Self::TwoHours,
7 => Self::FourHours,
8 => Self::EightHours,
9 => Self::TwelveHours,
10 => Self::OneDay,
11 => Self::ThreeDays,
12 => Self::OneWeek,
13 => Self::OneMonth,
x => panic!("Invalid CandleInterval discriminant: {}", x),
}
}
fn to_com(&self) -> Vec<u8> {
vec![match self {
Self::OneMinute => 0,
Self::ThreeMinutes => 1,
Self::FiveMinutes => 2,
Self::FifteenMinutes => 3,
Self::ThirtyMinutes => 4,
Self::OneHour => 5,
Self::TwoHours => 6,
Self::FourHours => 7,
Self::EightHours => 8,
Self::TwelveHours => 9,
Self::OneDay => 10,
Self::ThreeDays => 11,
Self::OneWeek => 12,
Self::OneMonth => 13,
}]
}
}
impl PulseWire for Candle {
fn from_com(com: &mut Vec<u8>) -> Self {
Self {
open_time: u64::from_com(com),
close_time: u64::from_com(com),
coin: String::from_com(com),
interval: String::from_com(com),
open: Decimal::from_com(com),
high: Decimal::from_com(com),
low: Decimal::from_com(com),
close: Decimal::from_com(com),
volume: Decimal::from_com(com),
num_trades: u64::from_com(com),
}
}
fn to_com(&self) -> Vec<u8> {
let mut com = Vec::new();
com.extend(self.open_time.to_com());
com.extend(self.close_time.to_com());
com.extend(self.coin.to_com());
com.extend(self.interval.to_com());
com.extend(self.open.to_com());
com.extend(self.high.to_com());
com.extend(self.low.to_com());
com.extend(self.close.to_com());
com.extend(self.volume.to_com());
com.extend(self.num_trades.to_com());
com
}
}
impl PulseWire for Decimal {
fn from_com(com: &mut Vec<u8>) -> Self {
Self::deserialize(com[..16].try_into().unwrap())
}
fn to_com(&self) -> Vec<u8> {
self.serialize().to_vec()
}
}
+33 -2
View File
@@ -2,17 +2,20 @@
use std::path::PathBuf; use std::path::PathBuf;
pub mod general; pub mod general;
pub mod strategy; mod hyper_types;
pub mod plugin;
pub mod terminal; pub mod terminal;
pub mod units; pub mod units;
pub use hypersdk;
pub mod prelude { pub mod prelude {
pub use crate::PulseWire; pub use crate::PulseWire;
pub use crate::general::*; pub use crate::general::*;
pub use crate::plugin::*;
pub use crate::server_path; pub use crate::server_path;
pub use crate::strategy::*;
pub use crate::terminal::*; pub use crate::terminal::*;
pub use crate::units::*; pub use crate::units::*;
pub use hypersdk;
} }
pub fn server_path() -> PathBuf { pub fn server_path() -> PathBuf {
@@ -52,6 +55,24 @@ impl<T: PulseWire> PulseWire for Vec<T> {
} }
} }
impl<T: PulseWire> PulseWire for Option<T> {
fn to_com(&self) -> Vec<u8> {
if let Some(v) = self {
vec![1].into_iter().chain(v.to_com()).collect()
} else {
vec![0]
}
}
fn from_com(com: &mut Vec<u8>) -> Self {
if com[0] > 0 {
Some(T::from_com(com))
} else {
None
}
}
}
impl PulseWire for String { impl PulseWire for String {
fn to_com(&self) -> Vec<u8> { fn to_com(&self) -> Vec<u8> {
let bytes = self.as_bytes(); let bytes = self.as_bytes();
@@ -72,6 +93,16 @@ impl PulseWire for String {
} }
} }
impl PulseWire for bool {
fn to_com(&self) -> Vec<u8> {
if *self { vec![1] } else { vec![0] }
}
fn from_com(com: &mut Vec<u8>) -> Self {
com[0] > 0
}
}
macro_rules! int_com { macro_rules! int_com {
($t:ty) => { ($t:ty) => {
impl $crate::PulseWire for $t { impl $crate::PulseWire for $t {
+96
View File
@@ -0,0 +1,96 @@
use crate::{
PulseWire,
general::{EventLog, Signal},
terminal::MarketItem,
units::Direction,
};
use hypersdk::hypercore::{Candle, CandleInterval, Subscription};
use pulse_macros::pwp;
#[pwp]
#[derive(serde::Deserialize, serde::Serialize)]
pub struct StrategyManifest {
name: String,
description: String,
author: String,
version: String,
}
#[pwp]
#[derive(serde::Deserialize, serde::Serialize)]
pub struct RiskManifest {
name: String,
description: String,
author: String,
version: String,
max_loss: u8,
cooldown: CandleInterval,
}
#[pwp]
pub enum StrategyMessage {
RequestOHLC {
symbol: String,
interval: CandleInterval,
count: u32,
},
Subscribe(Subscription),
Unsubscribe(Subscription),
UnsubscribeAll,
Signal(StrategySignal),
Log(EventLog),
}
#[pwp]
pub enum StrategyEngineMessage {
Initialize {
watchlist: Vec<MarketItem>,
},
CandleUpdate {
symbol: String,
candle: String,
},
OHLC {
symbol: String,
interval: CandleInterval,
candles: Vec<Candle>,
},
Start,
Stop,
}
#[pwp]
pub enum RiskMessage {
Approve(Signal),
Reject { reason: String },
Log(EventLog),
}
#[pwp]
pub enum RiskEngineMessage {
Initialize { strategy: String },
Signal(StrategySignal),
MarketUpdate { symbol: String, price: f64 },
}
#[pwp]
pub struct StrategySignal {
pub symbol: String,
pub side: Direction,
pub confidence: f32,
pub price: Option<f64>,
}
-25
View File
@@ -1,25 +0,0 @@
use crate::{PulseWire, units::TimeFrame};
use pulse_macros::pwp;
#[pwp]
#[derive(serde::Deserialize, serde::Serialize)]
pub struct StrategyManifest {
name: String,
description: String,
author: String,
version: String,
timeframes: Vec<TimeFrame>,
}
#[pwp]
#[derive(serde::Deserialize, serde::Serialize)]
pub struct RiskManifest {
name: String,
description: String,
author: String,
version: String,
max_loss: u8,
cooldown: TimeFrame,
}
+6 -5
View File
@@ -1,9 +1,10 @@
use crate::{ use crate::{
PulseWire, PulseWire,
general::{EventLog, MarketTrend, Position, Signal}, general::{EventLog, MarketTrend, Position, Signal},
strategy::{RiskManifest, StrategyManifest}, plugin::{RiskManifest, StrategyManifest},
units::{Symbol, TimeFrame, USD, Volatility}, units::{Symbol, USD, Volatility},
}; };
use hypersdk::hypercore::CandleInterval;
use pulse_macros::pwp; use pulse_macros::pwp;
#[pwp] #[pwp]
@@ -102,7 +103,7 @@ pub enum Mode {
#[pwp] #[pwp]
pub enum ItemState { pub enum ItemState {
Running, Running,
Off, Stopped,
Error, Error,
} }
@@ -113,7 +114,7 @@ pub struct Strategy {
mode: Mode, mode: Mode,
state: ItemState, state: ItemState,
cooldown: TimeFrame, cooldown: CandleInterval,
} }
impl std::fmt::Display for AlertLevel { impl std::fmt::Display for AlertLevel {
@@ -139,7 +140,7 @@ impl std::fmt::Display for ItemState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self { match self {
Self::Running => write!(f, "\x1b[92mRUNNING\x1b[0m"), Self::Running => write!(f, "\x1b[92mRUNNING\x1b[0m"),
Self::Off => write!(f, "\x1b[90mOFF\x1b[0m"), Self::Stopped => write!(f, "\x1b[90mSTOPPED\x1b[0m"),
Self::Error => write!(f, "\x1b[91mERROR\x1b[0m"), Self::Error => write!(f, "\x1b[91mERROR\x1b[0m"),
} }
} }
-60
View File
@@ -8,42 +8,6 @@ pub struct Symbol(pub String);
#[derive(Debug, Clone, Copy)] #[derive(Debug, Clone, Copy)]
pub struct USD(pub f64); pub struct USD(pub f64);
#[pwp]
#[derive(Copy, serde::Deserialize, serde::Serialize)]
pub enum TimeFrame {
#[serde(rename = "1m")]
M1,
#[serde(rename = "3m")]
M3,
#[serde(rename = "5m")]
M5,
#[serde(rename = "15m")]
M15,
#[serde(rename = "30m")]
M30,
#[serde(rename = "1h")]
H1,
#[serde(rename = "2h")]
H2,
#[serde(rename = "4h")]
H4,
#[serde(rename = "8h")]
H8,
#[serde(rename = "12h")]
H12,
#[serde(rename = "1d")]
D1,
#[serde(rename = "3d")]
D3,
#[serde(rename = "1w")]
W1,
#[serde(rename = "1M")]
Month1,
}
#[pwp] #[pwp]
pub enum Direction { pub enum Direction {
Buy, Buy,
@@ -148,30 +112,6 @@ pub fn format_f64(value: f64) -> String {
formatted formatted
} }
impl TimeFrame {
pub fn as_str(&self) -> &'static str {
match self {
Self::M1 => "1m",
Self::M3 => "3m",
Self::M5 => "5m",
Self::M15 => "15m",
Self::M30 => "30m",
Self::H1 => "1h",
Self::H2 => "2h",
Self::H4 => "4h",
Self::H8 => "8h",
Self::H12 => "12h",
Self::D1 => "1d",
Self::D3 => "3d",
Self::W1 => "1w",
Self::Month1 => "1M",
}
}
}
impl std::fmt::Display for Direction { impl std::fmt::Display for Direction {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self { match self {
+207
View File
@@ -0,0 +1,207 @@
use crate::{engine::Engine, store::config::Config};
use pulse_wire::terminal::{ItemState, Mode, Strategy};
use toml::Value;
impl Engine {
pub async fn execute_command(&self, command: &str, args: Vec<&str>) -> tokio::io::Result<()> {
match command {
"config" | "cfg" => {
if args.len() == 0 {
return self.invalid_command_usage("config").await;
}
match args[0] {
"reload" => {
*self.config.lock().await = Config::new().await?;
}
"save" => {
self.config.lock().await.save().await?;
}
"set" => {
macro_rules! set_cfg {
($n:ident, $v:expr) => {
if let Ok(Ok($n)) = Value::try_from(args[2]).map(|v| v.try_into()) {
$v;
} else {
self.terminal_server
.error("config::set", "Unable to parse value")
.await?;
}
};
}
if args.len() == 3 {
match args[1] {
"watchlist" | "watch" | "wl" => {
set_cfg!(v, self.config.lock().await.watchlist = v);
self.terminal_server
.info("config::set", "watchlist set successfully, use `config save` to persist changes")
.await?;
}
"strategy" | "strat" | "str" | "sg" => {
set_cfg!(id, {
let id: String = id;
if !crate::store::pulse_plugin(&id)?
.join("strategy.toml")
.exists()
{
return self
.terminal_server
.error(
"config::set::strategy",
&format!("Non existent strategy `{id}`"),
)
.await;
}
self.strategy.reload_strategy(id.as_str()).await?;
self.config.lock().await.strategy = id;
});
self.terminal_server
.info("config::set", "strategy set successfully, use `config save` to persist changes")
.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" => {
set_cfg!(cooldown, {
self.config.lock().await.cooldown = cooldown;
self.terminal_server.broadcast(
pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: self.strategy.strategy.manifest.lock().await.clone(),
risk: self.strategy.risk.manifest.lock().await.clone(),
mode: Mode::Auto,
state: ItemState::Running,
cooldown,
}),
)
.await?;
});
}
_ => {
self.terminal_server
.error("config::set", "Invalid usage, available options: watchlist, strategy, risk, cooldown")
.await?;
}
}
} else {
self.invalid_command_usage("config").await?;
}
}
_ => {
self.invalid_command_usage("config").await?;
}
}
}
"account" | "acc" => {
if args.len() == 0 {
return self.invalid_command_usage("account man").await;
}
match args[0] {
"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(())
}
}
+13 -103
View File
@@ -1,4 +1,8 @@
pub mod plugin;
pub mod command;
use crate::{ use crate::{
engine::plugin::StrategyEngine,
store::{accounts::AccountList, config::Config}, store::{accounts::AccountList, config::Config},
terminal::TerminalServer, terminal::TerminalServer,
}; };
@@ -6,20 +10,26 @@ use pulse_wire::prelude::*;
use std::sync::Arc; use std::sync::Arc;
use tokio::{sync::Mutex, task::JoinHandle}; use tokio::{sync::Mutex, task::JoinHandle};
#[derive(Debug, Clone)] #[derive(Clone)]
pub struct Engine { pub struct Engine {
pub terminal_server: Arc<TerminalServer>, pub terminal_server: Arc<TerminalServer>,
pub strategy: Arc<StrategyEngine>,
pub config: Arc<Mutex<Config>>, pub config: Arc<Mutex<Config>>,
pub accounts: Arc<Mutex<AccountList>>, pub accounts: Arc<Mutex<AccountList>>,
} }
impl Engine { impl Engine {
pub async fn new() -> tokio::io::Result<Arc<Self>> { pub async fn new() -> tokio::io::Result<Arc<Self>> {
let config = Arc::new(Mutex::new(Config::new().await?)); let config = Config::new().await?;
let strategy = StrategyEngine::new(&config.strategy, &config.risk).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));
Ok(Arc::new_cyclic(|engine| Self { Ok(Arc::new_cyclic(|engine| Self {
terminal_server: TerminalServer::new(engine.clone()), terminal_server: TerminalServer::new(engine.clone()),
strategy: strategy.initialize(engine.clone()),
config, config,
accounts, accounts,
})) }))
@@ -45,7 +55,7 @@ impl Engine {
loop { loop {
refresh.tick().await; refresh.tick().await;
let watch_list = &self.config.lock().await.watchlist.symbols; let watch_list = &self.config.lock().await.watchlist;
match crate::fetch::fetch_watch_list(&client, watch_list).await { match crate::fetch::fetch_watch_list(&client, watch_list).await {
Ok(watch_list) => { Ok(watch_list) => {
@@ -127,104 +137,4 @@ impl Engine {
pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> { pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> {
self.terminal_server.error(name, "Invalid usage").await 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(())
}
} }
+223
View File
@@ -0,0 +1,223 @@
use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket};
use pulse_wire::prelude::*;
use std::{
collections::HashSet,
path::PathBuf,
process::Stdio,
sync::{Arc, Weak},
time::{SystemTime, UNIX_EPOCH},
};
use tokio::{
fs,
process::{Child, Command},
sync::Mutex,
};
use crate::{
engine::Engine,
store::{plugin::Plugin, pulse_plugin},
};
pub struct StrategyEngine {
pub strategy: Arc<Plugin<StrategyEngineMessage, StrategyMessage, StrategyManifest>>,
pub risk: Arc<Plugin<RiskEngineMessage, RiskMessage, RiskManifest>>,
pub engine: Weak<Engine>,
pub ws: WebSocket,
pub subscriptions: Mutex<HashSet<Subscription>>,
}
impl StrategyEngine {
pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result<Self> {
let strategy = pulse_plugin(strategy_id)?;
let risk = pulse_plugin(risk_id)?;
let (strategy, strategy_manifest) = get_manifest_plugin_pair(
&strategy,
&strategy.join("strategy.bash"),
&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 {
strategy: Arc::new(Plugin::new(strategy, strategy_manifest)),
risk: Arc::new(Plugin::new(risk, risk_manifest)),
engine: Weak::new(),
ws: hypercore::mainnet_ws(),
subscriptions: Mutex::new(HashSet::new()),
})
}
pub fn initialize(mut self, engine: Weak<Engine>) -> Arc<Self> {
self.engine = engine;
Arc::new(self)
}
pub async fn run_strategy(&self) -> anyhow::Result<()> {
let engine = self
.engine
.upgrade()
.expect("Failed to upgrade engine (StrategyEngine)");
let strategy = self.strategy.clone();
let risk = self.risk.clone();
loop {
match strategy.recv().await? {
None => {}
Some(StrategyMessage::Log(mut log)) => {
log.name.insert_str(0, "strategy::");
engine.terminal_server.log_raw(log).await?;
}
Some(StrategyMessage::Signal(signal)) => {
risk.send(&RiskEngineMessage::Signal(signal)).await?;
}
Some(StrategyMessage::Subscribe(subscription)) => {
self.ws.subscribe(subscription.clone());
self.subscriptions.lock().await.insert(subscription);
}
Some(StrategyMessage::Unsubscribe(subscription)) => {
self.subscriptions.lock().await.remove(&subscription);
self.ws.unsubscribe(subscription);
}
Some(StrategyMessage::UnsubscribeAll) => {
for sub in self.subscriptions.lock().await.drain() {
self.ws.unsubscribe(sub);
}
}
Some(StrategyMessage::RequestOHLC {
symbol,
interval,
count,
}) => {
let client = hypercore::mainnet();
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
let interval_ms = match interval {
CandleInterval::OneMinute => 60_000,
CandleInterval::ThreeMinutes => 3 * 60_000,
CandleInterval::FiveMinutes => 5 * 60_000,
CandleInterval::FifteenMinutes => 15 * 60_000,
CandleInterval::ThirtyMinutes => 30 * 60_000,
CandleInterval::OneHour => 60 * 60_000,
CandleInterval::TwoHours => 2 * 60 * 60_000,
CandleInterval::FourHours => 4 * 60 * 60_000,
CandleInterval::EightHours => 8 * 60 * 60_000,
CandleInterval::TwelveHours => 12 * 60 * 60_000,
CandleInterval::OneDay => 24 * 60 * 60_000,
CandleInterval::ThreeDays => 3 * 24 * 60 * 60_000,
CandleInterval::OneWeek => 7 * 24 * 60 * 60_000,
CandleInterval::OneMonth => 30 * 24 * 60 * 60_000,
};
let start_time = now.saturating_sub(interval_ms * count as u64);
client
.candle_snapshot(symbol, interval, start_time, now)
.await?;
}
}
}
}
pub async fn run_risk(&self) -> tokio::io::Result<()> {
let engine = self
.engine
.upgrade()
.expect("Failed to upgrade engine (StrategyEngine)");
let risk = self.risk.clone();
loop {
match risk.recv().await? {
None => {}
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>) {
let engine = self.clone();
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<()> {
let plugin = pulse_plugin(id)?;
let (child, manifest) = get_manifest_plugin_pair(
&plugin,
&plugin.join("strategy.bash"),
&fs::read(plugin.join("strategy.toml")).await?,
)?;
self.strategy.reload(child, manifest).await?;
let engine = self.clone();
tokio::spawn(async move { engine.run_strategy().await });
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>>(
plugin_dir: &PathBuf,
plugin_path: &PathBuf,
manifest: &'de [u8],
) -> tokio::io::Result<(Child, M)> {
Ok((
Command::new("bash")
.arg(plugin_path)
.current_dir(plugin_dir)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::inherit())
.spawn()?,
toml::from_slice(manifest)
.map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))?,
))
}
+3 -1
View File
@@ -4,13 +4,15 @@ pub mod store;
pub mod terminal; pub mod terminal;
#[tokio::main] #[tokio::main]
async fn main() -> tokio::io::Result<()> { async fn main() -> anyhow::Result<()> {
let engine = engine::Engine::new().await?; let engine = engine::Engine::new().await?;
let terminal_server = engine.spawn_terminal_server().await; let terminal_server = engine.spawn_terminal_server().await;
let broadcaster = engine.spawn_broadcaster().await; let broadcaster = engine.spawn_broadcaster().await;
engine.strategy.spawn().await;
engine.run_engine().await?; engine.run_engine().await?;
let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster); let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster);
+25 -6
View File
@@ -1,11 +1,22 @@
#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] use hypersdk::hypercore::CandleInterval;
pub struct WatchList {
pub symbols: Vec<String>, #[derive(Debug, serde::Serialize, serde::Deserialize)]
pub struct Config {
pub watchlist: Vec<String>,
pub strategy: String,
pub risk: String,
pub cooldown: CandleInterval,
} }
#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] impl Default for Config {
pub struct Config { fn default() -> Self {
pub watchlist: WatchList, Self {
watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()],
strategy: String::new(),
risk: String::new(),
cooldown: CandleInterval::ThirtyMinutes,
}
}
} }
impl Config { impl Config {
@@ -26,6 +37,14 @@ impl Config {
Self::from_str(&output) Self::from_str(&output)
} }
pub async fn save(&self) -> tokio::io::Result<()> {
let path = crate::store::pulse_config_file()?;
tokio::fs::write(path, self.to_string()?).await?;
Ok(())
}
pub fn from_str(s: &str) -> tokio::io::Result<Self> { pub fn from_str(s: &str) -> tokio::io::Result<Self> {
toml::from_str(s) toml::from_str(s)
.map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string())) .map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))
+9
View File
@@ -2,6 +2,7 @@ use std::path::PathBuf;
pub mod accounts; pub mod accounts;
pub mod config; pub mod config;
pub mod plugin;
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(|| {
@@ -13,6 +14,14 @@ 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> {
Ok(pulse_directory()?.join("plugins"))
}
pub fn pulse_plugin(id: &str) -> tokio::io::Result<PathBuf> {
Ok(pulse_directory()?.join("plugins").join(id))
}
pub fn pulse_config_file() -> tokio::io::Result<PathBuf> { pub fn pulse_config_file() -> tokio::io::Result<PathBuf> {
Ok(pulse_directory()?.join("config.toml")) Ok(pulse_directory()?.join("config.toml"))
} }
+77
View File
@@ -0,0 +1,77 @@
use std::marker::PhantomData;
use pulse_wire::PulseWire;
use serde::Deserialize;
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
process::{Child, ChildStdout},
sync::Mutex,
};
#[derive(Debug)]
pub struct Plugin<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> {
pub manifest: Mutex<M>,
pub stdout: Mutex<ChildStdout>,
pub process: Mutex<Child>,
pub _p: (PhantomData<S>, PhantomData<R>),
}
impl<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
pub fn new(mut child: Child, manifest: M) -> Self {
Self {
stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")),
process: Mutex::new(child),
manifest: Mutex::new(manifest),
_p: (PhantomData, PhantomData),
}
}
pub async fn recv(&self) -> tokio::io::Result<Option<R>> {
let mut stdout = self.stdout.lock().await;
let mut len_buf = [0u8; size_of::<usize>()];
let size = stdout.read_exact(&mut len_buf).await?;
let len = usize::from_le_bytes(len_buf);
if size == 0 || len == 0 {
return Ok(None);
}
let mut buffer = vec![0u8; len];
stdout.read_exact(&mut buffer).await?;
Ok(Some(R::from_com(&mut buffer)))
}
pub async fn send(&self, msg: &S) -> tokio::io::Result<()> {
self.send_raw(&msg.to_com()).await
}
pub async fn send_raw(&self, msg: &[u8]) -> tokio::io::Result<()> {
let mut process = self.process.lock().await;
let stdin = process.stdin.as_mut().unwrap();
stdin.write(&msg.len().to_le_bytes()).await?;
stdin.write(msg).await?;
stdin.flush().await?;
Ok(())
}
pub async fn reload(&self, mut child: Child, manifest: M) -> tokio::io::Result<()> {
let mut process = self.process.lock().await;
process.kill().await?;
*self.stdout.lock().await = child.stdout.take().expect("Failed to obtain child stdout");
*self.manifest.lock().await = manifest;
*process = child;
Ok(())
}
}
+32 -5
View File
@@ -59,17 +59,37 @@ impl TerminalServer {
} }
} }
async fn handle_client( async fn initialize_client(self: &Arc<Self>, id: &usize) -> tokio::io::Result<()> {
self: &Arc<Self>, let engine = self.get_engine();
id: &usize,
mut reader: OwnedReadHalf, self.send_to(
) -> tokio::io::Result<()> { id,
pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: engine.strategy.strategy.manifest.lock().await.clone(),
risk: engine.strategy.risk.manifest.lock().await.clone(),
mode: Mode::Auto,
state: ItemState::Running,
cooldown: engine.config.lock().await.cooldown,
}),
)
.await?;
self.send_to( self.send_to(
id, id,
pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()),
) )
.await?; .await?;
Ok(())
}
async fn handle_client(
self: &Arc<Self>,
id: &usize,
mut reader: OwnedReadHalf,
) -> tokio::io::Result<()> {
self.initialize_client(id).await?;
loop { loop {
let mut len_buf = [0u8; size_of::<usize>()]; let mut len_buf = [0u8; size_of::<usize>()];
let size = reader.read_exact(&mut len_buf).await?; let size = reader.read_exact(&mut len_buf).await?;
@@ -169,6 +189,13 @@ impl TerminalServer {
.await .await
} }
pub async fn log_raw(self: &Arc<Self>, log: EventLog) -> tokio::io::Result<()> {
self.logs.lock().await.push(log.clone());
self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog(log))
.await
}
pub async fn info(self: &Arc<Self>, name: &str, message: &str) -> tokio::io::Result<()> { pub async fn info(self: &Arc<Self>, name: &str, message: &str) -> tokio::io::Result<()> {
self.log(LogKind::Info, name, message).await self.log(LogKind::Info, name, message).await
} }
+5 -5
View File
@@ -79,7 +79,7 @@ impl Formatted for MarketItem {
self.price.to_string(), self.price.to_string(),
self.volume_24h.to_string(), self.volume_24h.to_string(),
format!( format!(
"{} {}", "{} {}%",
if self.trend.is_sign_positive() { if self.trend.is_sign_positive() {
"\x1b[32m▲\x1b[0m" "\x1b[32m▲\x1b[0m"
} else { } else {
@@ -231,7 +231,7 @@ impl Formatted for Strategy {
fn get_formatted(&self) -> Vec<String> { fn get_formatted(&self) -> Vec<String> {
vec![ vec![
Triple( Triple(
"\x1b[2mName\x1b[0m", "\x1b[2mStrategy\x1b[0m",
"\x1b[2mRisk\x1b[0m", "\x1b[2mRisk\x1b[0m",
"\x1b[2mStrat Ver\x1b[0m", "\x1b[2mStrat Ver\x1b[0m",
), ),
@@ -254,15 +254,15 @@ impl Formatted for Strategy {
Triple("", "", ""), Triple("", "", ""),
Triple( Triple(
"\x1b[2mCooldown\x1b[0m", "\x1b[2mCooldown\x1b[0m",
"\x1b[2mTimeframes\x1b[0m", "\x1b[2mStrat Author\x1b[0m",
"\x1b[2mMax loss\x1b[0m", "\x1b[2mMax loss\x1b[0m",
), ),
Triple( Triple(
&format!( &format!(
"\x1b[96m{:?}\x1b[0m (\x1b[90m{:?} rec\x1b[0m)", "\x1b[96m{}\x1b[0m (\x1b[90m{} rec\x1b[0m)",
self.cooldown, self.risk.cooldown self.cooldown, self.risk.cooldown
), ),
&format!("\x1b[96m{:?}\x1b[0m", self.strategy.timeframes), &format!("\x1b[96m{:?}\x1b[0m", self.strategy.author),
&format!("\x1b[93m{}%\x1b[0m", self.risk.max_loss), &format!("\x1b[93m{}%\x1b[0m", self.risk.max_loss),
), ),
] ]
+2 -27
View File
@@ -144,7 +144,7 @@ impl App for PulseTradeApp {
( (
LayoutItem::Widget(Size::Flex(1)), LayoutItem::Widget(Size::Flex(1)),
Box::new( Box::new(
advanced_option_draw(&self.scroll, 3, "STRATEGY", &self.strategy).await, advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.strategy).await,
), ),
), ),
( (
@@ -181,32 +181,7 @@ async fn main() -> tokio::io::Result<()> {
signals: ctx.use_state(Vec::new()), signals: ctx.use_state(Vec::new()),
logs: ctx.use_state(Vec::new()), logs: ctx.use_state(Vec::new()),
inspect: ctx.use_state(InspectTarget::None), inspect: ctx.use_state(InspectTarget::None),
strategy: ctx.use_state(Some(Strategy { strategy: ctx.use_state(None),
strategy: StrategyManifest {
name: "Liquidity Sweep".to_string(),
description:
"Detects liquidity grabs around key support and resistance levels."
.to_string(),
author: "Klesty Selimaj".to_string(),
version: "1.0.0".to_string(),
timeframes: vec![TimeFrame::M5, TimeFrame::M15, TimeFrame::H1],
},
risk: RiskManifest {
name: "Aggressive".to_string(),
description: "High-risk profile with larger position sizing.".to_string(),
author: "Klesty Selimaj".to_string(),
version: "1.0.0".to_string(),
max_loss: 5,
cooldown: TimeFrame::M15,
},
mode: Mode::Auto,
state: ItemState::Running,
cooldown: TimeFrame::M15,
})),
status: ctx.use_state(None), status: ctx.use_state(None),
}) })
}) })