diff --git a/Cargo.lock b/Cargo.lock index 8495645..48d222b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3882,6 +3882,7 @@ dependencies = [ name = "pulse-trader" version = "0.1.0-alpha.0" dependencies = [ + "anyhow", "chrono", "crossterm", "hypersdk", @@ -3906,6 +3907,7 @@ dependencies = [ name = "pulse-wire" version = "0.1.0-alpha.0" dependencies = [ + "hypersdk", "pulse-macros", "serde", ] @@ -5200,6 +5202,7 @@ dependencies = [ "libc", "mio", "pin-project-lite", + "signal-hook-registry", "socket2", "tokio-macros", "windows-sys 0.61.2", diff --git a/Cargo.toml b/Cargo.toml index 8780172..5371dc0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,8 +6,10 @@ edition = "2024" [dependencies] pulse-ui = { workspace = true } pulse-wire = { workspace = true } - -chrono = "0.4.45" +crossterm = { workspace = true } +hypersdk = { workspace = true } +serde = { workspace = true } +anyhow = { workspace = true } tokio = { workspace = true, features = [ "rt-multi-thread", "macros", @@ -15,13 +17,12 @@ tokio = { workspace = true, features = [ "fs", "io-util", "time", + "process", ] } -crossterm = { workspace = true } -hypersdk = "0.2.14" +chrono = "0.4.45" serde_json = "1" rand = "0.8.7" toml = "1.1.3" -serde = { workspace = true } [workspace] 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-ui = { path = "pulse-ui", 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" crossterm = "0.29.0" serde = { version = "1.0.229", features = ["serde_derive"] } +anyhow = "1.0.104" [[bin]] name = "pulse-trader" diff --git a/pulse-wire/Cargo.toml b/pulse-wire/Cargo.toml index 81b1d0c..e96fe28 100644 --- a/pulse-wire/Cargo.toml +++ b/pulse-wire/Cargo.toml @@ -6,3 +6,4 @@ edition = "2024" [dependencies] pulse-macros = { workspace = true } serde = { workspace = true } +hypersdk = { workspace = true } diff --git a/pulse-wire/src/hyper_types.rs b/pulse-wire/src/hyper_types.rs new file mode 100644 index 0000000..9e19906 --- /dev/null +++ b/pulse-wire/src/hyper_types.rs @@ -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) -> 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 { + 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) -> Self { + let bytes: [u8; 20] = com.drain(..20).collect::>().try_into().unwrap(); + Self::from_slice(&bytes) + } + + fn to_com(&self) -> Vec { + self.as_slice().to_vec() + } +} + +impl PulseWire for CandleInterval { + fn from_com(com: &mut Vec) -> 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 { + 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) -> 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 { + 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) -> Self { + Self::deserialize(com[..16].try_into().unwrap()) + } + + fn to_com(&self) -> Vec { + self.serialize().to_vec() + } +} diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index 689db77..0ff2c32 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -2,17 +2,20 @@ use std::path::PathBuf; pub mod general; -pub mod strategy; +mod hyper_types; +pub mod plugin; pub mod terminal; pub mod units; +pub use hypersdk; pub mod prelude { pub use crate::PulseWire; pub use crate::general::*; + pub use crate::plugin::*; pub use crate::server_path; - pub use crate::strategy::*; pub use crate::terminal::*; pub use crate::units::*; + pub use hypersdk; } pub fn server_path() -> PathBuf { @@ -52,6 +55,24 @@ impl PulseWire for Vec { } } +impl PulseWire for Option { + fn to_com(&self) -> Vec { + if let Some(v) = self { + vec![1].into_iter().chain(v.to_com()).collect() + } else { + vec![0] + } + } + + fn from_com(com: &mut Vec) -> Self { + if com[0] > 0 { + Some(T::from_com(com)) + } else { + None + } + } +} + impl PulseWire for String { fn to_com(&self) -> Vec { let bytes = self.as_bytes(); @@ -72,6 +93,16 @@ impl PulseWire for String { } } +impl PulseWire for bool { + fn to_com(&self) -> Vec { + if *self { vec![1] } else { vec![0] } + } + + fn from_com(com: &mut Vec) -> Self { + com[0] > 0 + } +} + macro_rules! int_com { ($t:ty) => { impl $crate::PulseWire for $t { diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs new file mode 100644 index 0000000..bd91135 --- /dev/null +++ b/pulse-wire/src/plugin.rs @@ -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, + }, + + CandleUpdate { + symbol: String, + candle: String, + }, + + OHLC { + symbol: String, + interval: CandleInterval, + candles: Vec, + }, + + 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, +} diff --git a/pulse-wire/src/strategy.rs b/pulse-wire/src/strategy.rs deleted file mode 100644 index 757ac70..0000000 --- a/pulse-wire/src/strategy.rs +++ /dev/null @@ -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, -} - -#[pwp] -#[derive(serde::Deserialize, serde::Serialize)] -pub struct RiskManifest { - name: String, - description: String, - author: String, - version: String, - - max_loss: u8, - cooldown: TimeFrame, -} diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index 5b26c2c..072001d 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -1,9 +1,10 @@ use crate::{ PulseWire, general::{EventLog, MarketTrend, Position, Signal}, - strategy::{RiskManifest, StrategyManifest}, - units::{Symbol, TimeFrame, USD, Volatility}, + plugin::{RiskManifest, StrategyManifest}, + units::{Symbol, USD, Volatility}, }; +use hypersdk::hypercore::CandleInterval; use pulse_macros::pwp; #[pwp] @@ -102,7 +103,7 @@ pub enum Mode { #[pwp] pub enum ItemState { Running, - Off, + Stopped, Error, } @@ -113,7 +114,7 @@ pub struct Strategy { mode: Mode, state: ItemState, - cooldown: TimeFrame, + cooldown: CandleInterval, } 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 { match self { 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"), } } diff --git a/pulse-wire/src/units.rs b/pulse-wire/src/units.rs index 06e84ae..c7db0e6 100644 --- a/pulse-wire/src/units.rs +++ b/pulse-wire/src/units.rs @@ -8,42 +8,6 @@ pub struct Symbol(pub String); #[derive(Debug, Clone, Copy)] 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] pub enum Direction { Buy, @@ -148,30 +112,6 @@ pub fn format_f64(value: f64) -> String { 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 { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs new file mode 100644 index 0000000..836f104 --- /dev/null +++ b/src/engine/engine/command.rs @@ -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(()) + } +} diff --git a/src/engine/engine.rs b/src/engine/engine/mod.rs similarity index 56% rename from src/engine/engine.rs rename to src/engine/engine/mod.rs index 438b06c..6afa49f 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine/mod.rs @@ -1,4 +1,8 @@ +pub mod plugin; +pub mod command; + use crate::{ + engine::plugin::StrategyEngine, store::{accounts::AccountList, config::Config}, terminal::TerminalServer, }; @@ -6,20 +10,26 @@ use pulse_wire::prelude::*; use std::sync::Arc; use tokio::{sync::Mutex, task::JoinHandle}; -#[derive(Debug, Clone)] +#[derive(Clone)] pub struct Engine { pub terminal_server: Arc, + pub strategy: Arc, pub config: Arc>, pub accounts: Arc>, } impl Engine { pub async fn new() -> tokio::io::Result> { - 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 config = Arc::new(Mutex::new(config)); Ok(Arc::new_cyclic(|engine| Self { terminal_server: TerminalServer::new(engine.clone()), + strategy: strategy.initialize(engine.clone()), config, accounts, })) @@ -45,7 +55,7 @@ impl Engine { loop { 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 { Ok(watch_list) => { @@ -127,104 +137,4 @@ 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, 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(()) - } } diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs new file mode 100644 index 0000000..51ebbcd --- /dev/null +++ b/src/engine/engine/plugin.rs @@ -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>, + pub risk: Arc>, + pub engine: Weak, + + pub ws: WebSocket, + pub subscriptions: Mutex>, +} + +impl StrategyEngine { + pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { + 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) -> Arc { + 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) { + 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, 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, 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()))?, + )) +} diff --git a/src/engine/main.rs b/src/engine/main.rs index f4f7f64..7555ce1 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -4,13 +4,15 @@ pub mod store; pub mod terminal; #[tokio::main] -async fn main() -> tokio::io::Result<()> { +async fn main() -> anyhow::Result<()> { let engine = engine::Engine::new().await?; let terminal_server = engine.spawn_terminal_server().await; let broadcaster = engine.spawn_broadcaster().await; + engine.strategy.spawn().await; + engine.run_engine().await?; let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster); diff --git a/src/engine/store/config.rs b/src/engine/store/config.rs index 2ead6ae..c16d916 100644 --- a/src/engine/store/config.rs +++ b/src/engine/store/config.rs @@ -1,11 +1,22 @@ -#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct WatchList { - pub symbols: Vec, +use hypersdk::hypercore::CandleInterval; + +#[derive(Debug, serde::Serialize, serde::Deserialize)] +pub struct Config { + pub watchlist: Vec, + pub strategy: String, + pub risk: String, + pub cooldown: CandleInterval, } -#[derive(Debug, Default, serde::Serialize, serde::Deserialize)] -pub struct Config { - pub watchlist: WatchList, +impl Default for Config { + fn default() -> Self { + Self { + watchlist: vec!["BTC".to_string(), "SOL".to_string(), "ETH".to_string()], + strategy: String::new(), + risk: String::new(), + cooldown: CandleInterval::ThirtyMinutes, + } + } } impl Config { @@ -26,6 +37,14 @@ impl Config { 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 { toml::from_str(s) .map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string())) diff --git a/src/engine/store/mod.rs b/src/engine/store/mod.rs index 5299a77..73b6dee 100644 --- a/src/engine/store/mod.rs +++ b/src/engine/store/mod.rs @@ -2,6 +2,7 @@ use std::path::PathBuf; pub mod accounts; pub mod config; +pub mod plugin; pub fn home_dir() -> tokio::io::Result { std::env::home_dir().ok_or_else(|| { @@ -13,6 +14,14 @@ pub fn pulse_directory() -> tokio::io::Result { Ok(home_dir()?.join(".config").join("pulse-trader")) } +pub fn pulse_plugins_directory() -> tokio::io::Result { + Ok(pulse_directory()?.join("plugins")) +} + +pub fn pulse_plugin(id: &str) -> tokio::io::Result { + Ok(pulse_directory()?.join("plugins").join(id)) +} + pub fn pulse_config_file() -> tokio::io::Result { Ok(pulse_directory()?.join("config.toml")) } diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs new file mode 100644 index 0000000..6f37cb4 --- /dev/null +++ b/src/engine/store/plugin.rs @@ -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 Deserialize<'de>> { + pub manifest: Mutex, + pub stdout: Mutex, + pub process: Mutex, + + pub _p: (PhantomData, PhantomData), +} + +impl Deserialize<'de>> Plugin { + 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> { + let mut stdout = self.stdout.lock().await; + + let mut len_buf = [0u8; size_of::()]; + 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(()) + } +} diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 7de2c44..6623b6a 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -59,17 +59,37 @@ impl TerminalServer { } } - async fn handle_client( - self: &Arc, - id: &usize, - mut reader: OwnedReadHalf, - ) -> tokio::io::Result<()> { + async fn initialize_client(self: &Arc, id: &usize) -> tokio::io::Result<()> { + let engine = self.get_engine(); + + self.send_to( + 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( id, pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), ) .await?; + Ok(()) + } + + async fn handle_client( + self: &Arc, + id: &usize, + mut reader: OwnedReadHalf, + ) -> tokio::io::Result<()> { + self.initialize_client(id).await?; + loop { let mut len_buf = [0u8; size_of::()]; let size = reader.read_exact(&mut len_buf).await?; @@ -169,6 +189,13 @@ impl TerminalServer { .await } + pub async fn log_raw(self: &Arc, 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, name: &str, message: &str) -> tokio::io::Result<()> { self.log(LogKind::Info, name, message).await } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 43a6ffc..b807ead 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -79,7 +79,7 @@ impl Formatted for MarketItem { self.price.to_string(), self.volume_24h.to_string(), format!( - "{} {}", + "{} {}%", if self.trend.is_sign_positive() { "\x1b[32m▲\x1b[0m" } else { @@ -231,7 +231,7 @@ impl Formatted for Strategy { fn get_formatted(&self) -> Vec { vec![ Triple( - "\x1b[2mName\x1b[0m", + "\x1b[2mStrategy\x1b[0m", "\x1b[2mRisk\x1b[0m", "\x1b[2mStrat Ver\x1b[0m", ), @@ -254,15 +254,15 @@ impl Formatted for Strategy { Triple("", "", ""), Triple( "\x1b[2mCooldown\x1b[0m", - "\x1b[2mTimeframes\x1b[0m", + "\x1b[2mStrat Author\x1b[0m", "\x1b[2mMax loss\x1b[0m", ), Triple( &format!( - "\x1b[96m{:?}\x1b[0m (\x1b[90m{:?} rec\x1b[0m)", + "\x1b[96m{}\x1b[0m (\x1b[90m{} rec\x1b[0m)", 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), ), ] diff --git a/src/terminal/main.rs b/src/terminal/main.rs index da08d43..979c8a6 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -144,7 +144,7 @@ impl App for PulseTradeApp { ( LayoutItem::Widget(Size::Flex(1)), 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()), logs: ctx.use_state(Vec::new()), inspect: ctx.use_state(InspectTarget::None), - strategy: ctx.use_state(Some(Strategy { - 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, - })), + strategy: ctx.use_state(None), status: ctx.use_state(None), }) })