diff --git a/Cargo.lock b/Cargo.lock index 3cf15f6..288a837 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3906,6 +3906,7 @@ dependencies = [ name = "pulse-wire" version = "0.1.0-alpha.0" dependencies = [ + "hypersdk", "pulse-macros", "serde", ] diff --git a/Cargo.toml b/Cargo.toml index 8b87ecb..79ece7b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,7 +10,7 @@ pulse-wire = { workspace = true } chrono = "0.4.45" tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time", "process"] } crossterm = { workspace = true } -hypersdk = "0.2.14" +hypersdk = { workspace = true } serde_json = "1" rand = "0.8.7" toml = "1.1.3" @@ -23,7 +23,7 @@ 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"] } 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..dcc44c0 --- /dev/null +++ b/pulse-wire/src/hyper_types.rs @@ -0,0 +1,225 @@ +use hypersdk::{Address, hypercore::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() + } +} diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index 0d8b6d8..be1a05a 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -2,6 +2,7 @@ use std::path::PathBuf; pub mod general; +mod hyper_types; pub mod plugin; pub mod terminal; pub mod units; @@ -9,8 +10,8 @@ pub mod units; pub mod prelude { pub use crate::PulseWire; pub use crate::general::*; - pub use crate::server_path; pub use crate::plugin::*; + pub use crate::server_path; pub use crate::terminal::*; pub use crate::units::*; } @@ -55,7 +56,7 @@ impl PulseWire for Vec { impl PulseWire for Option { fn to_com(&self) -> Vec { if let Some(v) = self { - v.to_com() + vec![1].into_iter().chain(v.to_com()).collect() } else { vec![0] } @@ -90,6 +91,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 index 07d2419..d4b49a8 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -4,6 +4,7 @@ use crate::{ terminal::MarketItem, units::{Direction, TimeFrame}, }; +use hypersdk::hypercore::Subscription; use pulse_macros::pwp; #[pwp] @@ -37,10 +38,7 @@ pub enum StrategyMessage { count: u32, }, - SubscribeCandle { - symbol: String, - timeframe: TimeFrame, - }, + Subscribe(Subscription), Signal(StrategySignal), diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index c510c6c..aa7e18e 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -9,7 +9,7 @@ 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, diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 6fdbf1d..86f163f 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -1,3 +1,4 @@ +use hypersdk::hypercore::{self, WebSocket}; use pulse_wire::prelude::*; use std::{ path::PathBuf, @@ -14,12 +15,11 @@ use crate::{ store::{plugin::Plugin, pulse_plugin}, }; -#[derive(Debug)] pub struct StrategyEngine { pub strategy: Arc>, pub risk: Arc>, - pub engine: Weak, + pub ws: WebSocket, } impl StrategyEngine { @@ -43,6 +43,7 @@ impl StrategyEngine { strategy: Arc::new(Plugin::new(strategy, strategy_manifest)), risk: Arc::new(Plugin::new(risk, risk_manifest)), engine: Weak::new(), + ws: hypercore::mainnet_ws(), }) } @@ -74,7 +75,7 @@ impl StrategyEngine { risk.send(&RiskEngineMessage::Signal(signal)).await?; } - Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} + Some(StrategyMessage::Subscribe(val)) => {} Some(StrategyMessage::RequestOHLC { symbol,