From 1a47d0c487c35046cd89f76f14498caf939aa7da Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 01:03:12 +0200 Subject: [PATCH] Subscriptions --- pulse-wire/src/hyper_types.rs | 90 ++++++++++++++++++++++++++++++++++- pulse-wire/src/plugin.rs | 18 +++---- pulse-wire/src/terminal.rs | 5 +- pulse-wire/src/units.rs | 60 ----------------------- src/engine/engine/plugin.rs | 58 ++++++++++++++++++++-- src/engine/terminal.rs | 2 +- src/terminal/formatting.rs | 4 +- 7 files changed, 159 insertions(+), 78 deletions(-) diff --git a/pulse-wire/src/hyper_types.rs b/pulse-wire/src/hyper_types.rs index dcc44c0..9e19906 100644 --- a/pulse-wire/src/hyper_types.rs +++ b/pulse-wire/src/hyper_types.rs @@ -1,4 +1,7 @@ -use hypersdk::{Address, hypercore::Subscription}; +use hypersdk::{ + Address, Decimal, + hypercore::{Candle, CandleInterval, Subscription}, +}; use crate::PulseWire; @@ -223,3 +226,88 @@ impl PulseWire for Address { 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/plugin.rs b/pulse-wire/src/plugin.rs index d4b49a8..db9139b 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -2,9 +2,9 @@ use crate::{ PulseWire, general::{EventLog, Signal}, terminal::MarketItem, - units::{Direction, TimeFrame}, + units::Direction, }; -use hypersdk::hypercore::Subscription; +use hypersdk::hypercore::{Candle, CandleInterval, Subscription}; use pulse_macros::pwp; #[pwp] @@ -14,8 +14,6 @@ pub struct StrategyManifest { description: String, author: String, version: String, - - timeframes: Vec, } #[pwp] @@ -27,19 +25,23 @@ pub struct RiskManifest { version: String, max_loss: u8, - cooldown: TimeFrame, + cooldown: String, } #[pwp] pub enum StrategyMessage { RequestOHLC { symbol: String, - timeframe: TimeFrame, + interval: CandleInterval, count: u32, }, Subscribe(Subscription), + Unsubscribe(Subscription), + + UnsubscribeAll, + Signal(StrategySignal), Log(EventLog), @@ -58,8 +60,8 @@ pub enum StrategyEngineMessage { OHLC { symbol: String, - timeframe: TimeFrame, - candles: Vec, + interval: CandleInterval, + candles: Vec, }, Start, diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index d59b3fc..006c646 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -2,8 +2,9 @@ use crate::{ PulseWire, general::{EventLog, MarketTrend, Position, Signal}, plugin::{RiskManifest, StrategyManifest}, - units::{Symbol, TimeFrame, USD, Volatility}, + units::{Symbol, USD, Volatility}, }; +use hypersdk::hypercore::CandleInterval; use pulse_macros::pwp; #[pwp] @@ -113,7 +114,7 @@ pub struct Strategy { mode: Mode, state: ItemState, - cooldown: TimeFrame, + cooldown: CandleInterval, } impl std::fmt::Display for AlertLevel { 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/plugin.rs b/src/engine/engine/plugin.rs index 86f163f..de38592 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -1,13 +1,16 @@ -use hypersdk::hypercore::{self, WebSocket}; +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::{ @@ -19,7 +22,9 @@ pub struct StrategyEngine { pub strategy: Arc>, pub risk: Arc>, pub engine: Weak, + pub ws: WebSocket, + pub subscriptions: Mutex>, } impl StrategyEngine { @@ -44,6 +49,7 @@ impl StrategyEngine { risk: Arc::new(Plugin::new(risk, risk_manifest)), engine: Weak::new(), ws: hypercore::mainnet_ws(), + subscriptions: Mutex::new(HashSet::new()), }) } @@ -75,13 +81,57 @@ impl StrategyEngine { risk.send(&RiskEngineMessage::Signal(signal)).await?; } - Some(StrategyMessage::Subscribe(val)) => {} + 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, - timeframe, + 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?; + } } } } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index f76c9d2..39eb997 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -69,7 +69,7 @@ impl TerminalServer { risk: engine.strategy.risk.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, - cooldown: TimeFrame::M15, + cooldown: hypersdk::hypercore::CandleInterval::ThirtyMinutes, }), ) .await?; diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index ed76ef9..b46ef22 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -254,7 +254,7 @@ 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( @@ -262,7 +262,7 @@ impl Formatted for Strategy { "\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), ), ]