Implementing hyper ws and types

This commit is contained in:
2026-07-27 00:17:43 +02:00
parent c8be31c88c
commit 7d0075e2f9
8 changed files with 249 additions and 12 deletions
Generated
+1
View File
@@ -3906,6 +3906,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",
] ]
+2 -2
View File
@@ -10,7 +10,7 @@ pulse-wire = { workspace = true }
chrono = "0.4.45" chrono = "0.4.45"
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time", "process"] } tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time", "process"] }
crossterm = { workspace = true } crossterm = { workspace = true }
hypersdk = "0.2.14" hypersdk = { workspace = true }
serde_json = "1" serde_json = "1"
rand = "0.8.7" rand = "0.8.7"
toml = "1.1.3" 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-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"] }
+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 }
+225
View File
@@ -0,0 +1,225 @@
use hypersdk::{Address, hypercore::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()
}
}
+13 -2
View File
@@ -2,6 +2,7 @@
use std::path::PathBuf; use std::path::PathBuf;
pub mod general; pub mod general;
mod hyper_types;
pub mod plugin; pub mod plugin;
pub mod terminal; pub mod terminal;
pub mod units; pub mod units;
@@ -9,8 +10,8 @@ pub mod units;
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::server_path;
pub use crate::plugin::*; pub use crate::plugin::*;
pub use crate::server_path;
pub use crate::terminal::*; pub use crate::terminal::*;
pub use crate::units::*; pub use crate::units::*;
} }
@@ -55,7 +56,7 @@ impl<T: PulseWire> PulseWire for Vec<T> {
impl<T: PulseWire> PulseWire for Option<T> { impl<T: PulseWire> PulseWire for Option<T> {
fn to_com(&self) -> Vec<u8> { fn to_com(&self) -> Vec<u8> {
if let Some(v) = self { if let Some(v) = self {
v.to_com() vec![1].into_iter().chain(v.to_com()).collect()
} else { } else {
vec![0] vec![0]
} }
@@ -90,6 +91,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 {
+2 -4
View File
@@ -4,6 +4,7 @@ use crate::{
terminal::MarketItem, terminal::MarketItem,
units::{Direction, TimeFrame}, units::{Direction, TimeFrame},
}; };
use hypersdk::hypercore::Subscription;
use pulse_macros::pwp; use pulse_macros::pwp;
#[pwp] #[pwp]
@@ -37,10 +38,7 @@ pub enum StrategyMessage {
count: u32, count: u32,
}, },
SubscribeCandle { Subscribe(Subscription),
symbol: String,
timeframe: TimeFrame,
},
Signal(StrategySignal), Signal(StrategySignal),
+1 -1
View File
@@ -9,7 +9,7 @@ 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 strategy: Arc<StrategyEngine>,
+4 -3
View File
@@ -1,3 +1,4 @@
use hypersdk::hypercore::{self, WebSocket};
use pulse_wire::prelude::*; use pulse_wire::prelude::*;
use std::{ use std::{
path::PathBuf, path::PathBuf,
@@ -14,12 +15,11 @@ use crate::{
store::{plugin::Plugin, pulse_plugin}, store::{plugin::Plugin, pulse_plugin},
}; };
#[derive(Debug)]
pub struct StrategyEngine { pub struct StrategyEngine {
pub strategy: Arc<Plugin<StrategyEngineMessage, StrategyMessage, StrategyManifest>>, pub strategy: Arc<Plugin<StrategyEngineMessage, StrategyMessage, StrategyManifest>>,
pub risk: Arc<Plugin<RiskEngineMessage, RiskMessage, RiskManifest>>, pub risk: Arc<Plugin<RiskEngineMessage, RiskMessage, RiskManifest>>,
pub engine: Weak<Engine>, pub engine: Weak<Engine>,
pub ws: WebSocket,
} }
impl StrategyEngine { impl StrategyEngine {
@@ -43,6 +43,7 @@ impl StrategyEngine {
strategy: Arc::new(Plugin::new(strategy, strategy_manifest)), strategy: Arc::new(Plugin::new(strategy, strategy_manifest)),
risk: Arc::new(Plugin::new(risk, risk_manifest)), risk: Arc::new(Plugin::new(risk, risk_manifest)),
engine: Weak::new(), engine: Weak::new(),
ws: hypercore::mainnet_ws(),
}) })
} }
@@ -74,7 +75,7 @@ impl StrategyEngine {
risk.send(&RiskEngineMessage::Signal(signal)).await?; risk.send(&RiskEngineMessage::Signal(signal)).await?;
} }
Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} Some(StrategyMessage::Subscribe(val)) => {}
Some(StrategyMessage::RequestOHLC { Some(StrategyMessage::RequestOHLC {
symbol, symbol,