From 7b1281aa5704485a5a079aa03169f04975edefba Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 27 Jul 2026 19:30:45 +0200 Subject: [PATCH 01/12] Fix UI overflows --- pulse-ui/src/render.rs | 52 +++++++++++++++++++++++++++++++++-- pulse-ui/src/widget/scroll.rs | 13 ++++++--- 2 files changed, 59 insertions(+), 6 deletions(-) diff --git a/pulse-ui/src/render.rs b/pulse-ui/src/render.rs index bf8940d..dea001d 100644 --- a/pulse-ui/src/render.rs +++ b/pulse-ui/src/render.rs @@ -16,9 +16,57 @@ pub struct RenderScope { } impl RenderScope { - pub fn draw_text, T: Display>(&mut self, at: P, text: T) { + pub fn draw_text, T: Display>(&mut self, at: P, text: T) -> u16 { + let point: Point = at.into(); + + let lines = text + .to_string() + .lines() + .map(|line| { + let mut chars = line.chars().peekable(); + let mut new_line = String::new(); + let mut len = 0; + + while let Some(c) = chars.next() { + if c == '\x1b' && chars.peek() == Some(&'[') { + new_line.push(c); + + while let Some(c) = chars.next() { + new_line.push(c); + + if c.is_ascii_alphabetic() { + break; + } + } + + continue; + } + + new_line.push(c); + len += 1; + + if len >= self.rect.width as usize { + new_line.push('\n'); + len = 0; + } + } + + new_line + }) + .collect::>() + .join("\n"); + + let lines = lines + .lines() + .take((self.rect.height - point.y) as usize) + .collect::>(); + + let lines_len = lines.len(); + self.draw_instructions - .push(Instr::DrawText(at.into(), text.to_string())); + .push(Instr::DrawText(point, lines.join("\n"))); + + lines_len as u16 } } diff --git a/pulse-ui/src/widget/scroll.rs b/pulse-ui/src/widget/scroll.rs index be6511e..2a12586 100644 --- a/pulse-ui/src/widget/scroll.rs +++ b/pulse-ui/src/widget/scroll.rs @@ -14,21 +14,26 @@ impl Widget for ScrollText { let title_lines = self.title.lines().count(); - for (y, line) in self + let mut y = title_lines as u16; + + for line in self .text .lines() .skip(self.scroll) .take(scope.rect.height as usize - title_lines) - .enumerate() { - scope.draw_text((0, (y + title_lines) as u16), line); + y += scope.draw_text((0, y), line); } } } impl ScrollState { pub fn get_selected(&self, index: usize) -> &'static str { - if index == self.0 { "\x1b[4m\x1b[1m" } else { "\x1b[1m" } + if index == self.0 { + "\x1b[4m\x1b[1m" + } else { + "\x1b[1m" + } } pub fn scroll(&self, index: usize, title: String, text: String) -> ScrollText { From f3656ec94056b055e6823d778c89ab95852ac30b Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 00:45:38 +0200 Subject: [PATCH 02/12] Removed unnecessary run from engine --- src/engine/engine/mod.rs | 14 +------------- src/engine/main.rs | 9 ++------- 2 files changed, 3 insertions(+), 20 deletions(-) diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 6afa49f..e2c4235 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -1,5 +1,5 @@ -pub mod plugin; pub mod command; +pub mod plugin; use crate::{ engine::plugin::StrategyEngine, @@ -35,12 +35,6 @@ impl Engine { })) } - pub async fn spawn_terminal_server(&self) -> JoinHandle> { - let terminal_server = self.terminal_server.clone(); - - tokio::spawn(async move { terminal_server.run().await }) - } - pub async fn spawn_broadcaster(&self) -> JoinHandle> { let s = self.clone(); @@ -128,12 +122,6 @@ impl Engine { } } - pub async fn run_engine(&self) -> tokio::io::Result<()> { - loop { - tokio::time::sleep(tokio::time::Duration::from_millis(5000)).await; - } - } - pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> { self.terminal_server.error(name, "Invalid usage").await } diff --git a/src/engine/main.rs b/src/engine/main.rs index 7555ce1..1d45b2d 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -7,18 +7,13 @@ pub mod terminal; 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?; + engine.terminal_server.run().await?; - let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster); - - terminal_server??; - broadcaster??; + broadcaster.await??; Ok(()) } From 796bbb397b3544eac183d396e72352b1af18946e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 03:33:22 +0200 Subject: [PATCH 03/12] Strategy protocol --- pulse-wire/src/plugin.rs | 48 ++++++++++++++++++++++--------------- src/engine/engine/plugin.rs | 22 ++++++++++++++++- 2 files changed, 50 insertions(+), 20 deletions(-) diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs index bd91135..6550ceb 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -30,7 +30,11 @@ pub struct RiskManifest { #[pwp] pub enum StrategyMessage { - RequestOHLC { + Log(EventLog), + + GetWatchList, + + RequestCandlestick { symbol: String, interval: CandleInterval, count: u32, @@ -43,44 +47,50 @@ pub enum StrategyMessage { UnsubscribeAll, Signal(StrategySignal), +} +#[pwp] +pub enum RiskMessage { Log(EventLog), + + GetWatchList, + + Approve(Signal), + + Reject { reason: String }, } #[pwp] pub enum StrategyEngineMessage { - Initialize { - watchlist: Vec, - }, + Initialize, + + WatchList(Vec), CandleUpdate { symbol: String, - candle: String, + interval: CandleInterval, + candle: Candle, }, - OHLC { + Candlestick { symbol: String, interval: CandleInterval, candles: Vec, }, - Start, - - Stop, -} - -#[pwp] -pub enum RiskMessage { - Approve(Signal), - - Reject { reason: String }, - - Log(EventLog), + Command { + command: String, + args: Vec, + }, } #[pwp] pub enum RiskEngineMessage { - Initialize { strategy: String }, + Initialize, + + WatchList(Vec), + + Command { command: String, args: Vec }, Signal(StrategySignal), diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 51ebbcd..f620822 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -68,10 +68,20 @@ impl StrategyEngine { let strategy = self.strategy.clone(); let risk = self.risk.clone(); + strategy.send(&StrategyEngineMessage::Initialize).await?; + loop { match strategy.recv().await? { None => {} + Some(StrategyMessage::GetWatchList) => { + let mut v = vec![1]; + + v.extend(engine.config.lock().await.watchlist.to_com()); + + self.strategy.send_raw(&v).await?; + } + Some(StrategyMessage::Log(mut log)) => { log.name.insert_str(0, "strategy::"); engine.terminal_server.log_raw(log).await?; @@ -97,7 +107,7 @@ impl StrategyEngine { } } - Some(StrategyMessage::RequestOHLC { + Some(StrategyMessage::RequestCandlestick { symbol, interval, count, @@ -144,10 +154,20 @@ impl StrategyEngine { let risk = self.risk.clone(); + risk.send(&RiskEngineMessage::Initialize).await?; + loop { match risk.recv().await? { None => {} + Some(RiskMessage::GetWatchList) => { + let mut v = vec![1]; + + v.extend(engine.config.lock().await.watchlist.to_com()); + + self.risk.send_raw(&v).await?; + } + Some(RiskMessage::Log(mut log)) => { log.name.insert_str(0, "risk::"); engine.terminal_server.log_raw(log).await?; From dc660decee41c8f43945f5f0bc2fe8f5f7ffb41a Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 03:34:42 +0200 Subject: [PATCH 04/12] Removed unnecessary clones on strategy engine --- src/engine/engine/plugin.rs | 17 +++++++---------- 1 file changed, 7 insertions(+), 10 deletions(-) diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index f620822..990e18f 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -65,13 +65,12 @@ impl StrategyEngine { .upgrade() .expect("Failed to upgrade engine (StrategyEngine)"); - let strategy = self.strategy.clone(); - let risk = self.risk.clone(); - - strategy.send(&StrategyEngineMessage::Initialize).await?; + self.strategy + .send(&StrategyEngineMessage::Initialize) + .await?; loop { - match strategy.recv().await? { + match self.strategy.recv().await? { None => {} Some(StrategyMessage::GetWatchList) => { @@ -88,7 +87,7 @@ impl StrategyEngine { } Some(StrategyMessage::Signal(signal)) => { - risk.send(&RiskEngineMessage::Signal(signal)).await?; + self.risk.send(&RiskEngineMessage::Signal(signal)).await?; } Some(StrategyMessage::Subscribe(subscription)) => { @@ -152,12 +151,10 @@ impl StrategyEngine { .upgrade() .expect("Failed to upgrade engine (StrategyEngine)"); - let risk = self.risk.clone(); - - risk.send(&RiskEngineMessage::Initialize).await?; + self.risk.send(&RiskEngineMessage::Initialize).await?; loop { - match risk.recv().await? { + match self.risk.recv().await? { None => {} Some(RiskMessage::GetWatchList) => { From b5069954b19b97196748a5945dea4349e7e2e201 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 06:51:26 +0200 Subject: [PATCH 05/12] Small changes --- pulse-wire/src/plugin.rs | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs index 6550ceb..c0df02a 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -66,6 +66,11 @@ pub enum StrategyEngineMessage { WatchList(Vec), + Command { + command: String, + args: Vec, + }, + CandleUpdate { symbol: String, interval: CandleInterval, @@ -77,11 +82,6 @@ pub enum StrategyEngineMessage { interval: CandleInterval, candles: Vec, }, - - Command { - command: String, - args: Vec, - }, } #[pwp] @@ -93,8 +93,6 @@ pub enum RiskEngineMessage { Command { command: String, args: Vec }, Signal(StrategySignal), - - MarketUpdate { symbol: String, price: f64 }, } #[pwp] From 3e96dd5c0f4e23f2e2db7fadbe52ddb3d40ce402 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 18:25:44 +0200 Subject: [PATCH 06/12] Fix --- src/engine/store/plugin.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 6f37cb4..761fa1f 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -56,8 +56,8 @@ impl Deserialize<'de>> Plugin { 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.write_all(&msg.len().to_le_bytes()).await?; + stdin.write_all(msg).await?; stdin.flush().await?; Ok(()) From fca09449b590dc7d3e818618b7b9a2d412caeaa7 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 19:56:48 +0200 Subject: [PATCH 07/12] Fix watch list --- src/engine/engine/mod.rs | 4 ++++ src/engine/engine/plugin.rs | 10 +++++----- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index e2c4235..c38e5dd 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -16,6 +16,7 @@ pub struct Engine { pub strategy: Arc, pub config: Arc>, pub accounts: Arc>, + pub watch_list: Arc>>, } impl Engine { @@ -32,6 +33,7 @@ impl Engine { strategy: strategy.initialize(engine.clone()), config, accounts, + watch_list: Arc::new(Mutex::new(Vec::new())), })) } @@ -53,6 +55,8 @@ impl Engine { match crate::fetch::fetch_watch_list(&client, watch_list).await { Ok(watch_list) => { + *self.watch_list.lock().await = watch_list.clone(); + if let Err(error) = self .terminal_server .broadcast( diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 990e18f..3054651 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -74,11 +74,11 @@ impl StrategyEngine { None => {} Some(StrategyMessage::GetWatchList) => { - let mut v = vec![1]; - - v.extend(engine.config.lock().await.watchlist.to_com()); - - self.strategy.send_raw(&v).await?; + self.strategy + .send(&StrategyEngineMessage::WatchList( + engine.watch_list.lock().await.clone(), + )) + .await?; } Some(StrategyMessage::Log(mut log)) => { From 1fb2ab9f2b0b315680f260449813c82571191203 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 20:13:34 +0200 Subject: [PATCH 08/12] Fix GetCandlestick --- pulse-wire/src/plugin.rs | 2 +- src/engine/engine/plugin.rs | 14 +++++++++++--- 2 files changed, 12 insertions(+), 4 deletions(-) diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs index c0df02a..875413c 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -34,7 +34,7 @@ pub enum StrategyMessage { GetWatchList, - RequestCandlestick { + GetCandlestick { symbol: String, interval: CandleInterval, count: u32, diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 3054651..223e59b 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -106,7 +106,7 @@ impl StrategyEngine { } } - Some(StrategyMessage::RequestCandlestick { + Some(StrategyMessage::GetCandlestick { symbol, interval, count, @@ -118,6 +118,8 @@ impl StrategyEngine { .unwrap() .as_millis() as u64; + println!("getting candlestick"); + let interval_ms = match interval { CandleInterval::OneMinute => 60_000, CandleInterval::ThreeMinutes => 3 * 60_000, @@ -137,8 +139,14 @@ impl StrategyEngine { let start_time = now.saturating_sub(interval_ms * count as u64); - client - .candle_snapshot(symbol, interval, start_time, now) + self.strategy + .send(&StrategyEngineMessage::Candlestick { + candles: client + .candle_snapshot(&symbol, interval, start_time, now) + .await?, + symbol, + interval, + }) .await?; } } From 489ddb36e97a41dc07d7d00a98f960b51ed1497e Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 21:49:17 +0200 Subject: [PATCH 09/12] Removed PWP bc it hurt my brain --- Cargo.lock | 10 -- Cargo.toml | 3 +- pulse-macros/Cargo.toml | 12 -- pulse-macros/src/lib.rs | 170 ------------------ pulse-wire/Cargo.toml | 1 - pulse-wire/src/general.rs | 7 +- pulse-wire/src/hyper_types.rs | 313 ---------------------------------- pulse-wire/src/lib.rs | 114 ------------- 8 files changed, 2 insertions(+), 628 deletions(-) delete mode 100644 pulse-macros/Cargo.toml delete mode 100644 pulse-macros/src/lib.rs delete mode 100644 pulse-wire/src/hyper_types.rs diff --git a/Cargo.lock b/Cargo.lock index 48d222b..fe934ad 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3869,15 +3869,6 @@ dependencies = [ "syn 1.0.109", ] -[[package]] -name = "pulse-macros" -version = "0.1.0-alpha.0" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.118", -] - [[package]] name = "pulse-trader" version = "0.1.0-alpha.0" @@ -3908,7 +3899,6 @@ name = "pulse-wire" version = "0.1.0-alpha.0" dependencies = [ "hypersdk", - "pulse-macros", "serde", ] diff --git a/Cargo.toml b/Cargo.toml index 5371dc0..e809ccf 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -25,10 +25,9 @@ rand = "0.8.7" toml = "1.1.3" [workspace] -members = ["pulse-macros", "pulse-ui", "pulse-wire"] +members = ["pulse-ui", "pulse-wire"] [workspace.dependencies] -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" diff --git a/pulse-macros/Cargo.toml b/pulse-macros/Cargo.toml deleted file mode 100644 index 882d11d..0000000 --- a/pulse-macros/Cargo.toml +++ /dev/null @@ -1,12 +0,0 @@ -[package] -name = "pulse-macros" -version = "0.1.0-alpha.0" -edition = "2024" - -[lib] -proc-macro = true - -[dependencies] -syn = { version = "2", features = ["full"] } -quote = "1" -proc-macro2 = "1" diff --git a/pulse-macros/src/lib.rs b/pulse-macros/src/lib.rs deleted file mode 100644 index cac53b8..0000000 --- a/pulse-macros/src/lib.rs +++ /dev/null @@ -1,170 +0,0 @@ -use proc_macro::TokenStream; -use quote::quote; -use syn::{Fields, ItemEnum, ItemStruct, parse_macro_input}; - -#[proc_macro_attribute] -pub fn pwp(_: TokenStream, item: TokenStream) -> TokenStream { - let input = parse_macro_input!(item as syn::Item); - - match input { - syn::Item::Struct(s) => expand_struct(s), - syn::Item::Enum(e) => expand_enum(e), - _ => { - return syn::Error::new_spanned(input, "p_com only supports structs and enums") - .to_compile_error() - .into(); - } - } -} - -fn expand_struct(mut input: ItemStruct) -> TokenStream { - let name = &input.ident; - - match &mut input.fields { - Fields::Named(fields) => { - for field in fields.named.iter_mut() { - field.vis = syn::Visibility::Public(syn::token::Pub::default()); - } - } - _ => { - return syn::Error::new_spanned(input, "p_com only supports structs with named fields") - .to_compile_error() - .into(); - } - } - - let fields = match &input.fields { - Fields::Named(fields) => &fields.named, - _ => { - return syn::Error::new_spanned(input, "p_com only supports structs with named fields") - .to_compile_error() - .into(); - } - }; - - let field_names = fields.iter().map(|f| f.ident.as_ref().unwrap()); - let field_names2 = fields.iter().map(|f| f.ident.as_ref().unwrap()); - - TokenStream::from(quote! { - #[derive(Debug, Clone)] - #input - - impl PulseWire for #name { - fn to_com(&self) -> Vec { - let mut vec = Vec::new(); - - #( - vec.extend(self.#field_names.to_com()); - )* - - vec - } - - fn from_com(com: &mut Vec) -> Self { - Self { - #( - #field_names2: PulseWire::from_com(com), - )* - } - } - } - }) -} - -fn expand_enum(input: ItemEnum) -> TokenStream { - let name = &input.ident; - - let to_com = input.variants.iter().enumerate().map(|(i, variant)| { - let ident = &variant.ident; - let tag = i as u8; - - match &variant.fields { - Fields::Unit => quote! { - Self::#ident => { - vec.push(#tag); - } - }, - - Fields::Unnamed(fields) if fields.unnamed.len() == 1 => quote! { - Self::#ident(v) => { - vec.push(#tag); - vec.extend(v.to_com()); - } - }, - - Fields::Named(fields) => { - let names = fields.named.iter().map(|f| f.ident.as_ref().unwrap()); - - let names2 = fields.named.iter().map(|f| f.ident.as_ref().unwrap()); - - quote! { - Self::#ident { #( #names ),* } => { - vec.push(#tag); - #( vec.extend(#names2.to_com()); )* - } - } - } - - _ => { - panic!("tuple variants with >1 field are not supported"); - } - } - }); - - let from_com = input.variants.iter().enumerate().map(|(i, variant)| { - let ident = &variant.ident; - let tag = i as u8; - - match &variant.fields { - Fields::Unit => quote! { - #tag => Self::#ident, - }, - - Fields::Unnamed(fields) if fields.unnamed.len() == 1 => { - quote! { - #tag => Self::#ident(PulseWire::from_com(com)), - } - } - - Fields::Named(fields) => { - let names = fields.named.iter().map(|f| f.ident.as_ref().unwrap()); - - quote! { - #tag => Self::#ident { - #( - #names: PulseWire::from_com(com), - )* - }, - } - } - - _ => panic!("tuple variants with >1 field are not supported"), - } - }); - - TokenStream::from(quote! { - #[derive(Debug, Clone)] - #input - - impl PulseWire for #name { - fn to_com(&self) -> Vec { - let mut vec = Vec::new(); - - match self { - #( #to_com )* - } - - vec - } - - fn from_com(com: &mut Vec) -> Self { - let kind = com.remove(0); - - match kind { - #( #from_com )* - _ => panic!("invalid {} discriminant {}", stringify!(#name), kind), - } - } - } - }) -} diff --git a/pulse-wire/Cargo.toml b/pulse-wire/Cargo.toml index e96fe28..54d4acc 100644 --- a/pulse-wire/Cargo.toml +++ b/pulse-wire/Cargo.toml @@ -4,6 +4,5 @@ version = "0.1.0-alpha.0" edition = "2024" [dependencies] -pulse-macros = { workspace = true } serde = { workspace = true } hypersdk = { workspace = true } diff --git a/pulse-wire/src/general.rs b/pulse-wire/src/general.rs index 353cb24..97f9dfc 100644 --- a/pulse-wire/src/general.rs +++ b/pulse-wire/src/general.rs @@ -1,9 +1,4 @@ -use pulse_macros::pwp; - -use crate::{ - PulseWire, - units::{Direction, Symbol, USD}, -}; +use crate::units::{Direction, Symbol, USD}; #[pwp] pub enum MarketTrend { diff --git a/pulse-wire/src/hyper_types.rs b/pulse-wire/src/hyper_types.rs deleted file mode 100644 index 9e19906..0000000 --- a/pulse-wire/src/hyper_types.rs +++ /dev/null @@ -1,313 +0,0 @@ -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 0ff2c32..cca42fc 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -2,14 +2,12 @@ use std::path::PathBuf; pub mod general; -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; @@ -21,115 +19,3 @@ pub mod prelude { pub fn server_path() -> PathBuf { PathBuf::from("/tmp/pulse-engine.sock") } - -pub trait PulseWire { - fn to_com(&self) -> Vec; - fn from_com(_com: &mut Vec) -> Self; -} - -impl PulseWire for Vec { - fn to_com(&self) -> Vec { - let mut vec = Vec::new(); - - vec.extend_from_slice(&(self.len() as u32).to_le_bytes()); - - for item in self { - vec.extend(item.to_com()); - } - - vec - } - - fn from_com(com: &mut Vec) -> Self { - let len_bytes: [u8; 4] = com.drain(..4).collect::>().try_into().unwrap(); - - let len = u32::from_le_bytes(len_bytes) as usize; - - let mut result = Vec::with_capacity(len); - - for _ in 0..len { - result.push(T::from_com(com)); - } - - result - } -} - -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(); - let mut out = Vec::with_capacity(4 + bytes.len()); - - out.extend_from_slice(&(bytes.len() as u32).to_le_bytes()); - out.extend_from_slice(bytes); - - out - } - - fn from_com(com: &mut Vec) -> Self { - let len = u32::from_le_bytes(com[..4].try_into().unwrap()) as usize; - com.drain(..4); - - let bytes: Vec = com.drain(..len).collect(); - String::from_utf8(bytes).unwrap() - } -} - -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 { - fn to_com(&self) -> Vec { - self.to_le_bytes().to_vec() - } - - fn from_com(com: &mut Vec) -> Self { - const N: usize = std::mem::size_of::<$t>(); - let bytes: [u8; N] = com.drain(..N).collect::>().try_into().unwrap(); - <$t>::from_le_bytes(bytes) - } - } - }; -} - -int_com!(i8); -int_com!(i16); -int_com!(i32); -int_com!(i64); -int_com!(isize); - -int_com!(u8); -int_com!(u16); -int_com!(u32); -int_com!(u64); -int_com!(usize); - -int_com!(f64); -int_com!(f32); From 92fe70def2b226e401ab83f5e981d8ed03b336ee Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 21:50:30 +0200 Subject: [PATCH 10/12] Using postcard with serde --- Cargo.lock | 68 ++++++++++++++++++++++++++++++++++++++++++++++++++++++ Cargo.toml | 2 ++ 2 files changed, 70 insertions(+) diff --git a/Cargo.lock b/Cargo.lock index fe934ad..1aad095 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1069,6 +1069,15 @@ dependencies = [ "syn 3.0.3", ] +[[package]] +name = "atomic-polyfill" +version = "1.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8cf2bce30dfe09ef0bfaef228b9d414faaf7e563035494d7fe092dba54b300f4" +dependencies = [ + "critical-section", +] + [[package]] name = "atomic-waker" version = "1.1.2" @@ -1765,6 +1774,15 @@ version = "0.5.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a" +[[package]] +name = "cobs" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0fa961b519f0b462e3a3b4a34b64d119eeaca1d59af726fe450bbba07a9fc0a1" +dependencies = [ + "thiserror 2.0.19", +] + [[package]] name = "coins-ledger" version = "0.13.0" @@ -1936,6 +1954,12 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "critical-section" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" + [[package]] name = "crossbeam-utils" version = "0.8.22" @@ -2602,6 +2626,15 @@ dependencies = [ "tracing", ] +[[package]] +name = "hash32" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0c35f58762feb77d74ebe43bdbc3210f09be9fe6742234d573bacc26ed92b67" +dependencies = [ + "byteorder", +] + [[package]] name = "hashbrown" version = "0.12.3" @@ -2639,6 +2672,20 @@ dependencies = [ "serde_core", ] +[[package]] +name = "heapless" +version = "0.7.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cdc6457c0eb62c71aac4bc17216026d8410337c4126773b9c5daba343f17964f" +dependencies = [ + "atomic-polyfill", + "hash32", + "rustc_version 0.4.1", + "serde", + "spin", + "stable_deref_trait", +] + [[package]] name = "heck" version = "0.5.0" @@ -3686,6 +3733,17 @@ version = "0.3.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e" +[[package]] +name = "postcard" +version = "1.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6764c3b5dd454e283a30e6dfe78e9b31096d9e32036b5d1eaac7a6119ccb9a24" +dependencies = [ + "cobs", + "heapless", + "serde", +] + [[package]] name = "potential_utf" version = "0.1.5" @@ -3877,6 +3935,7 @@ dependencies = [ "chrono", "crossterm", "hypersdk", + "postcard", "pulse-ui", "pulse-wire", "rand 0.8.7", @@ -4952,6 +5011,15 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "spin" +version = "0.9.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3763264f6b73151db08c50ff20d7d8a0b8796e021cdea7ceedad07b80155fa0e" +dependencies = [ + "lock_api", +] + [[package]] name = "spki" version = "0.7.3" diff --git a/Cargo.toml b/Cargo.toml index e809ccf..7f050ea 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,6 +10,7 @@ crossterm = { workspace = true } hypersdk = { workspace = true } serde = { workspace = true } anyhow = { workspace = true } +postcard = { workspace = true } tokio = { workspace = true, features = [ "rt-multi-thread", "macros", @@ -28,6 +29,7 @@ toml = "1.1.3" members = ["pulse-ui", "pulse-wire"] [workspace.dependencies] +postcard = "1.1.3" 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" From 7968ae14afe99cd6e4ff9d3c5ad87fc52df28d2b Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 22:20:50 +0200 Subject: [PATCH 11/12] Using postcard --- Cargo.lock | 16 +++++ Cargo.toml | 2 +- pulse-wire/Cargo.toml | 2 + pulse-wire/src/general.rs | 10 +-- pulse-wire/src/lib.rs | 4 ++ pulse-wire/src/plugin.rs | 38 +++++------ pulse-wire/src/terminal.rs | 50 +++++++------- pulse-wire/src/units.rs | 32 ++------- src/engine/engine/plugin.rs | 10 +-- src/engine/store/plugin.rs | 15 ++-- src/engine/terminal.rs | 8 +-- src/terminal/terminal.rs | 132 +++++++++++++++++++++--------------- 12 files changed, 166 insertions(+), 153 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 1aad095..5b95774 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2264,6 +2264,18 @@ dependencies = [ "zeroize", ] +[[package]] +name = "embedded-io" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ef1a6892d9eef45c8fa6b9e0086428a2cca8491aca8f787c534a3d6d0bcb3ced" + +[[package]] +name = "embedded-io" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "edd0f118536f44f5ccd48bcb8b111bdc3de888b58c74639dfb034a357d0f206d" + [[package]] name = "encoding_rs" version = "0.8.35" @@ -3740,6 +3752,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6764c3b5dd454e283a30e6dfe78e9b31096d9e32036b5d1eaac7a6119ccb9a24" dependencies = [ "cobs", + "embedded-io 0.4.0", + "embedded-io 0.6.1", "heapless", "serde", ] @@ -3958,7 +3972,9 @@ name = "pulse-wire" version = "0.1.0-alpha.0" dependencies = [ "hypersdk", + "postcard", "serde", + "tokio", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 7f050ea..33d3eee 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -29,7 +29,7 @@ toml = "1.1.3" members = ["pulse-ui", "pulse-wire"] [workspace.dependencies] -postcard = "1.1.3" +postcard = { version = "1.1.3", features = ["alloc"]} 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" diff --git a/pulse-wire/Cargo.toml b/pulse-wire/Cargo.toml index 54d4acc..bd40410 100644 --- a/pulse-wire/Cargo.toml +++ b/pulse-wire/Cargo.toml @@ -6,3 +6,5 @@ edition = "2024" [dependencies] serde = { workspace = true } hypersdk = { workspace = true } +postcard = { workspace = true } +tokio = { workspace = true } diff --git a/pulse-wire/src/general.rs b/pulse-wire/src/general.rs index 97f9dfc..a2fffad 100644 --- a/pulse-wire/src/general.rs +++ b/pulse-wire/src/general.rs @@ -1,13 +1,13 @@ use crate::units::{Direction, Symbol, USD}; -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum MarketTrend { Bullish, Bearish, Neutral, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum LogKind { Info, Warn, @@ -15,7 +15,7 @@ pub enum LogKind { Debug, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Signal { pub symbol: String, pub kind: Direction, @@ -26,14 +26,14 @@ pub struct Signal { pub stop_loss: USD, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct EventLog { pub kind: LogKind, pub name: String, pub message: String, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Position { pub symbol: Symbol, pub size: f64, diff --git a/pulse-wire/src/lib.rs b/pulse-wire/src/lib.rs index cca42fc..dd2211e 100644 --- a/pulse-wire/src/lib.rs +++ b/pulse-wire/src/lib.rs @@ -19,3 +19,7 @@ pub mod prelude { pub fn server_path() -> PathBuf { PathBuf::from("/tmp/pulse-engine.sock") } + +pub fn map_postcard_err(res: postcard::Result) -> tokio::io::Result { + res.map_err(|e| tokio::io::Error::new(std::io::ErrorKind::Other, e)) +} diff --git a/pulse-wire/src/plugin.rs b/pulse-wire/src/plugin.rs index 875413c..821a226 100644 --- a/pulse-wire/src/plugin.rs +++ b/pulse-wire/src/plugin.rs @@ -1,34 +1,30 @@ 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)] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct StrategyManifest { - name: String, - description: String, - author: String, - version: String, + pub name: String, + pub description: String, + pub author: String, + pub version: String, } -#[pwp] -#[derive(serde::Deserialize, serde::Serialize)] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct RiskManifest { - name: String, - description: String, - author: String, - version: String, + pub name: String, + pub description: String, + pub author: String, + pub version: String, - max_loss: u8, - cooldown: CandleInterval, + pub max_loss: u8, + pub cooldown: CandleInterval, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum StrategyMessage { Log(EventLog), @@ -49,7 +45,7 @@ pub enum StrategyMessage { Signal(StrategySignal), } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum RiskMessage { Log(EventLog), @@ -60,7 +56,7 @@ pub enum RiskMessage { Reject { reason: String }, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum StrategyEngineMessage { Initialize, @@ -84,7 +80,7 @@ pub enum StrategyEngineMessage { }, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum RiskEngineMessage { Initialize, @@ -95,7 +91,7 @@ pub enum RiskEngineMessage { Signal(StrategySignal), } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct StrategySignal { pub symbol: String, pub side: Direction, diff --git a/pulse-wire/src/terminal.rs b/pulse-wire/src/terminal.rs index 072001d..a07e9f1 100644 --- a/pulse-wire/src/terminal.rs +++ b/pulse-wire/src/terminal.rs @@ -1,13 +1,11 @@ use crate::{ - PulseWire, general::{EventLog, MarketTrend, Position, Signal}, plugin::{RiskManifest, StrategyManifest}, units::{Symbol, USD, Volatility}, }; use hypersdk::hypercore::CandleInterval; -use pulse_macros::pwp; -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum TerminalServerMessage { // WatchList WatchListUpdated(Vec), @@ -32,12 +30,12 @@ pub enum TerminalServerMessage { AddLog(EventLog), } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum TerminalClientMessage { ExecuteCommand(String), } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct MarketItem { pub symbol: Symbol, pub price: USD, @@ -45,7 +43,7 @@ pub struct MarketItem { pub volume_24h: USD, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct MarketOverview { pub trend: MarketTrend, pub volatility: Volatility, @@ -53,32 +51,32 @@ pub struct MarketOverview { pub alerts: Vec, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum AlertLevel { High, Medium, Low, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Alert { - level: AlertLevel, - message: String, + pub level: AlertLevel, + pub message: String, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Balance { pub asset: String, pub amount: f64, pub value: f64, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum InspectTarget { None, Some(Vec), } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum InspectItem { String(String), Symbol(String), @@ -86,35 +84,35 @@ pub enum InspectItem { F64(f64), } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Status { - feed: Mode, - exchange: String, - dex: String, - latency: u16, + pub feed: Mode, + pub exchange: String, + pub dex: String, + pub latency: u16, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum Mode { Auto, Manual, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum ItemState { Running, Stopped, Error, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Strategy { - strategy: StrategyManifest, - risk: RiskManifest, + pub strategy: StrategyManifest, + pub risk: RiskManifest, - mode: Mode, - state: ItemState, - cooldown: CandleInterval, + pub mode: Mode, + pub state: ItemState, + pub cooldown: CandleInterval, } impl std::fmt::Display for AlertLevel { diff --git a/pulse-wire/src/units.rs b/pulse-wire/src/units.rs index c7db0e6..e2758e9 100644 --- a/pulse-wire/src/units.rs +++ b/pulse-wire/src/units.rs @@ -1,20 +1,16 @@ -use pulse_macros::pwp; - -use crate::PulseWire; - -#[derive(Debug, Clone)] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Symbol(pub String); -#[derive(Debug, Clone, Copy)] +#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] pub struct USD(pub f64); -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum Direction { Buy, Sell, } -#[pwp] +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum Volatility { Low, Medium, @@ -37,26 +33,6 @@ impl std::fmt::Display for USD { } } -impl PulseWire for Symbol { - fn from_com(com: &mut Vec) -> Self { - Self(String::from_com(com)) - } - - fn to_com(&self) -> Vec { - self.0.to_com() - } -} - -impl PulseWire for USD { - fn from_com(com: &mut Vec) -> Self { - Self(f64::from_com(com)) - } - - fn to_com(&self) -> Vec { - self.0.to_com() - } -} - pub fn format_f64(value: f64) -> String { let abs = value.abs(); diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 223e59b..4109b13 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -166,11 +166,11 @@ impl StrategyEngine { None => {} Some(RiskMessage::GetWatchList) => { - let mut v = vec![1]; - - v.extend(engine.config.lock().await.watchlist.to_com()); - - self.risk.send_raw(&v).await?; + self.risk + .send(&&RiskEngineMessage::WatchList( + engine.watch_list.lock().await.clone(), + )) + .await?; } Some(RiskMessage::Log(mut log)) => { diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 761fa1f..9763dc9 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1,16 +1,14 @@ use std::marker::PhantomData; -use pulse_wire::PulseWire; - -use serde::Deserialize; +use pulse_wire::map_postcard_err; +use serde::{Deserialize, Serialize}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, process::{Child, ChildStdout}, sync::Mutex, }; - #[derive(Debug)] -pub struct Plugin Deserialize<'de>> { +pub struct Plugin Deserialize<'de>, M: for<'de> Deserialize<'de>> { pub manifest: Mutex, pub stdout: Mutex, pub process: Mutex, @@ -18,7 +16,7 @@ pub struct Plugin Deserialize<'de>> { pub _p: (PhantomData, PhantomData), } -impl Deserialize<'de>> Plugin { +impl Deserialize<'de>, M: for<'de> 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")), @@ -45,11 +43,12 @@ impl Deserialize<'de>> Plugin { stdout.read_exact(&mut buffer).await?; - Ok(Some(R::from_com(&mut buffer))) + Ok(Some(map_postcard_err(postcard::from_bytes(&buffer))?)) } pub async fn send(&self, msg: &S) -> tokio::io::Result<()> { - self.send_raw(&msg.to_com()).await + self.send_raw(&map_postcard_err(postcard::to_allocvec(msg))?) + .await } pub async fn send_raw(&self, msg: &[u8]) -> tokio::io::Result<()> { diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 6623b6a..df601a7 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,5 +1,5 @@ use crate::engine::Engine; -use pulse_wire::prelude::*; +use pulse_wire::{map_postcard_err, prelude::*}; use std::{ collections::HashMap, sync::{Arc, Weak}, @@ -104,7 +104,7 @@ impl TerminalServer { reader.read_exact(&mut buffer).await?; - match TerminalClientMessage::from_com(&mut buffer) { + match map_postcard_err(postcard::from_bytes(&buffer))? { TerminalClientMessage::ExecuteCommand(command) => { let command = command.as_str(); @@ -126,7 +126,7 @@ impl TerminalServer { self: &Arc, message: pulse_wire::terminal::TerminalServerMessage, ) -> tokio::io::Result<()> { - let msg = message.to_com(); + let msg = map_postcard_err(postcard::to_allocvec(&message))?; let mut clients = self.clients.lock().await; @@ -158,7 +158,7 @@ impl TerminalServer { format!("Client({id}) does not exist"), ) })?, - &message.to_com(), + &map_postcard_err(postcard::to_allocvec(&message))?, ) .await } diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index 468a20d..5519562 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -1,4 +1,5 @@ -use pulse_wire::prelude::*; +use pulse_ui::state::State; +use pulse_wire::{map_postcard_err, prelude::*}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, @@ -29,7 +30,7 @@ impl TerminalClient { &mut self, message: pulse_wire::terminal::TerminalClientMessage, ) -> tokio::io::Result<()> { - let msg = message.to_com(); + let msg = map_postcard_err(postcard::to_allocvec(&message))?; self.writer.write(&msg.len().to_le_bytes()).await?; self.writer.write(&msg).await?; self.writer.flush().await?; @@ -42,7 +43,7 @@ impl TerminalClient { std::mem::swap(&mut self.reader, &mut reader); - let mut reader = reader.expect("Reader failed to swap"); + let reader = reader.expect("Reader failed to swap"); let watch_list = app.watch_list.clone(); let active_positions = app.active_positions.clone(); @@ -52,61 +53,82 @@ impl TerminalClient { let status = app.status.clone(); let inspect = app.inspect.clone(); - tokio::spawn(async move { - loop { - let mut len_buf = [0u8; size_of::()]; - reader - .read_exact(&mut len_buf) - .await - .expect("Failed to get header length"); - - let len = usize::from_le_bytes(len_buf); - - let mut buffer = vec![0u8; len]; - - reader - .read_exact(&mut buffer) - .await - .expect("Failed to read socket"); - - match TerminalServerMessage::from_com(&mut buffer) { - TerminalServerMessage::WatchListUpdated(v) => { - *watch_list.lock().await = v; - } - - TerminalServerMessage::PositionsUpdated(v) => { - *active_positions.lock().await = v; - } - - TerminalServerMessage::StrategyUpdated(v) => { - *market_overview.lock().await = Some(v); - } - - TerminalServerMessage::SignalsUpdated(v) => { - *signals.lock().await = v; - } - - TerminalServerMessage::Inspect(v) => { - *inspect.lock().await = v; - } - - TerminalServerMessage::StatusUpdated(v) => { - *status.lock().await = Some(v); - } - - TerminalServerMessage::SetLogs(v) => { - *logs.lock().await = v; - } - - TerminalServerMessage::AddLog(v) => { - logs.lock().await.push(v); - } - } - } - }); + tokio::spawn(Self::run_client( + reader, + watch_list, + active_positions, + logs, + signals, + market_overview, + status, + inspect, + )); app.sock = Some(self); app } + + pub async fn run_client( + mut reader: OwnedReadHalf, + + watch_list: State>, + active_positions: State>, + logs: State>, + signals: State>, + market_overview: State>, + status: State>, + inspect: State, + ) -> tokio::io::Result<()> { + let mut len_buf = [0u8; size_of::()]; + reader + .read_exact(&mut len_buf) + .await + .expect("Failed to get header length"); + + let len = usize::from_le_bytes(len_buf); + + let mut buffer = vec![0u8; len]; + + reader + .read_exact(&mut buffer) + .await + .expect("Failed to read socket"); + + match map_postcard_err(postcard::from_bytes(&buffer))? { + TerminalServerMessage::WatchListUpdated(v) => { + *watch_list.lock().await = v; + } + + TerminalServerMessage::PositionsUpdated(v) => { + *active_positions.lock().await = v; + } + + TerminalServerMessage::StrategyUpdated(v) => { + *market_overview.lock().await = Some(v); + } + + TerminalServerMessage::SignalsUpdated(v) => { + *signals.lock().await = v; + } + + TerminalServerMessage::Inspect(v) => { + *inspect.lock().await = v; + } + + TerminalServerMessage::StatusUpdated(v) => { + *status.lock().await = Some(v); + } + + TerminalServerMessage::SetLogs(v) => { + *logs.lock().await = v; + } + + TerminalServerMessage::AddLog(v) => { + logs.lock().await.push(v); + } + } + + Ok(()) + } } From 5c65b19e978d5c716d8374bfe08fef0889048fe5 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Tue, 28 Jul 2026 22:33:19 +0200 Subject: [PATCH 12/12] Renamed pulse wire to sdk --- Cargo.lock | 22 +++++++++++----------- Cargo.toml | 8 ++++---- {pulse-wire => pulse-sdk}/Cargo.toml | 2 +- {pulse-wire => pulse-sdk}/src/general.rs | 0 {pulse-wire => pulse-sdk}/src/lib.rs | 0 {pulse-wire => pulse-sdk}/src/plugin.rs | 0 {pulse-wire => pulse-sdk}/src/terminal.rs | 0 {pulse-wire => pulse-sdk}/src/units.rs | 0 src/engine/engine/command.rs | 4 ++-- src/engine/engine/mod.rs | 6 +++--- src/engine/engine/plugin.rs | 2 +- src/engine/fetch.rs | 2 +- src/engine/store/plugin.rs | 2 +- src/engine/terminal.rs | 16 ++++++++-------- src/terminal/command.rs | 2 +- src/terminal/formatting.rs | 2 +- src/terminal/main.rs | 5 +++-- src/terminal/terminal.rs | 4 ++-- 18 files changed, 39 insertions(+), 38 deletions(-) rename {pulse-wire => pulse-sdk}/Cargo.toml (90%) rename {pulse-wire => pulse-sdk}/src/general.rs (100%) rename {pulse-wire => pulse-sdk}/src/lib.rs (100%) rename {pulse-wire => pulse-sdk}/src/plugin.rs (100%) rename {pulse-wire => pulse-sdk}/src/terminal.rs (100%) rename {pulse-wire => pulse-sdk}/src/units.rs (100%) diff --git a/Cargo.lock b/Cargo.lock index 5b95774..eac0541 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3941,6 +3941,16 @@ dependencies = [ "syn 1.0.109", ] +[[package]] +name = "pulse-sdk" +version = "0.1.0-alpha.0" +dependencies = [ + "hypersdk", + "postcard", + "serde", + "tokio", +] + [[package]] name = "pulse-trader" version = "0.1.0-alpha.0" @@ -3950,8 +3960,8 @@ dependencies = [ "crossterm", "hypersdk", "postcard", + "pulse-sdk", "pulse-ui", - "pulse-wire", "rand 0.8.7", "serde", "serde_json", @@ -3967,16 +3977,6 @@ dependencies = [ "tokio", ] -[[package]] -name = "pulse-wire" -version = "0.1.0-alpha.0" -dependencies = [ - "hypersdk", - "postcard", - "serde", - "tokio", -] - [[package]] name = "quinn" version = "0.11.11" diff --git a/Cargo.toml b/Cargo.toml index 33d3eee..9badfe3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,7 +5,7 @@ edition = "2024" [dependencies] pulse-ui = { workspace = true } -pulse-wire = { workspace = true } +pulse-sdk = { workspace = true } crossterm = { workspace = true } hypersdk = { workspace = true } serde = { workspace = true } @@ -26,12 +26,12 @@ rand = "0.8.7" toml = "1.1.3" [workspace] -members = ["pulse-ui", "pulse-wire"] +members = ["pulse-ui", "pulse-sdk"] [workspace.dependencies] -postcard = { version = "1.1.3", features = ["alloc"]} +postcard = { version = "1.1.3", features = ["alloc"] } pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" } -pulse-wire = { path = "pulse-wire", version = "0.1.0-alpha.0" } +pulse-sdk = { path = "pulse-sdk", version = "0.1.0-alpha.0" } hypersdk = "0.2.14" tokio = "1.52.3" crossterm = "0.29.0" diff --git a/pulse-wire/Cargo.toml b/pulse-sdk/Cargo.toml similarity index 90% rename from pulse-wire/Cargo.toml rename to pulse-sdk/Cargo.toml index bd40410..bb06f35 100644 --- a/pulse-wire/Cargo.toml +++ b/pulse-sdk/Cargo.toml @@ -1,5 +1,5 @@ [package] -name = "pulse-wire" +name = "pulse-sdk" version = "0.1.0-alpha.0" edition = "2024" diff --git a/pulse-wire/src/general.rs b/pulse-sdk/src/general.rs similarity index 100% rename from pulse-wire/src/general.rs rename to pulse-sdk/src/general.rs diff --git a/pulse-wire/src/lib.rs b/pulse-sdk/src/lib.rs similarity index 100% rename from pulse-wire/src/lib.rs rename to pulse-sdk/src/lib.rs diff --git a/pulse-wire/src/plugin.rs b/pulse-sdk/src/plugin.rs similarity index 100% rename from pulse-wire/src/plugin.rs rename to pulse-sdk/src/plugin.rs diff --git a/pulse-wire/src/terminal.rs b/pulse-sdk/src/terminal.rs similarity index 100% rename from pulse-wire/src/terminal.rs rename to pulse-sdk/src/terminal.rs diff --git a/pulse-wire/src/units.rs b/pulse-sdk/src/units.rs similarity index 100% rename from pulse-wire/src/units.rs rename to pulse-sdk/src/units.rs diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 836f104..bd39825 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -1,6 +1,6 @@ use crate::{engine::Engine, store::config::Config}; -use pulse_wire::terminal::{ItemState, Mode, Strategy}; +use pulse_sdk::terminal::{ItemState, Mode, Strategy}; use toml::Value; impl Engine { @@ -100,7 +100,7 @@ impl Engine { self.config.lock().await.cooldown = cooldown; self.terminal_server.broadcast( - pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { + pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { strategy: self.strategy.strategy.manifest.lock().await.clone(), risk: self.strategy.risk.manifest.lock().await.clone(), mode: Mode::Auto, diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index c38e5dd..bb1d46f 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -6,7 +6,7 @@ use crate::{ store::{accounts::AccountList, config::Config}, terminal::TerminalServer, }; -use pulse_wire::prelude::*; +use pulse_sdk::prelude::*; use std::sync::Arc; use tokio::{sync::Mutex, task::JoinHandle}; @@ -60,7 +60,7 @@ impl Engine { if let Err(error) = self .terminal_server .broadcast( - pulse_wire::terminal::TerminalServerMessage::WatchListUpdated( + pulse_sdk::terminal::TerminalServerMessage::WatchListUpdated( watch_list, ), ) @@ -92,7 +92,7 @@ impl Engine { Ok(state) => { self.terminal_server .broadcast( - pulse_wire::terminal::TerminalServerMessage::PositionsUpdated( + pulse_sdk::terminal::TerminalServerMessage::PositionsUpdated( state .asset_positions .into_iter() diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 4109b13..2c30d1b 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -1,5 +1,5 @@ use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket}; -use pulse_wire::prelude::*; +use pulse_sdk::prelude::*; use std::{ collections::HashSet, path::PathBuf, diff --git a/src/engine/fetch.rs b/src/engine/fetch.rs index 28a3733..d8be404 100644 --- a/src/engine/fetch.rs +++ b/src/engine/fetch.rs @@ -1,4 +1,4 @@ -use pulse_wire::prelude::*; +use pulse_sdk::prelude::*; use serde_json::Value; use std::collections::HashMap; diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 9763dc9..62449f6 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1,6 +1,6 @@ use std::marker::PhantomData; -use pulse_wire::map_postcard_err; +use pulse_sdk::map_postcard_err; use serde::{Deserialize, Serialize}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index df601a7..ba14436 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -1,5 +1,5 @@ use crate::engine::Engine; -use pulse_wire::{map_postcard_err, prelude::*}; +use pulse_sdk::{map_postcard_err, prelude::*}; use std::{ collections::HashMap, sync::{Arc, Weak}, @@ -30,7 +30,7 @@ impl TerminalServer { } pub async fn run(self: &Arc) -> tokio::io::Result<()> { - let path = pulse_wire::server_path(); + let path = pulse_sdk::server_path(); if path.exists() { tokio::fs::remove_file(&path).await?; @@ -64,7 +64,7 @@ impl TerminalServer { self.send_to( id, - pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { + pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { strategy: engine.strategy.strategy.manifest.lock().await.clone(), risk: engine.strategy.risk.manifest.lock().await.clone(), mode: Mode::Auto, @@ -76,7 +76,7 @@ impl TerminalServer { self.send_to( id, - pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), + pulse_sdk::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), ) .await?; @@ -124,7 +124,7 @@ impl TerminalServer { pub async fn broadcast( self: &Arc, - message: pulse_wire::terminal::TerminalServerMessage, + message: pulse_sdk::terminal::TerminalServerMessage, ) -> tokio::io::Result<()> { let msg = map_postcard_err(postcard::to_allocvec(&message))?; @@ -149,7 +149,7 @@ impl TerminalServer { pub async fn send_to( self: &Arc, id: &usize, - message: pulse_wire::terminal::TerminalServerMessage, + message: pulse_sdk::terminal::TerminalServerMessage, ) -> tokio::io::Result<()> { Self::send_to_client( self.clients.lock().await.get_mut(id).ok_or_else(|| { @@ -185,14 +185,14 @@ impl TerminalServer { self.logs.lock().await.push(log.clone()); - self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog(log)) + self.broadcast(pulse_sdk::terminal::TerminalServerMessage::AddLog(log)) .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)) + self.broadcast(pulse_sdk::terminal::TerminalServerMessage::AddLog(log)) .await } diff --git a/src/terminal/command.rs b/src/terminal/command.rs index b8849df..c0e01c3 100644 --- a/src/terminal/command.rs +++ b/src/terminal/command.rs @@ -16,7 +16,7 @@ impl PulseTradeApp { .sock .as_mut() .unwrap() - .send(pulse_wire::terminal::TerminalClientMessage::ExecuteCommand( + .send(pulse_sdk::terminal::TerminalClientMessage::ExecuteCommand( command.to_string(), )) .await diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index b807ead..50211d6 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -1,4 +1,4 @@ -use pulse_wire::prelude::*; +use pulse_sdk::prelude::*; pub trait Formatted { fn get_formatted(&self) -> Vec; diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 979c8a6..25ae2f5 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -21,7 +21,7 @@ use pulse_ui::{ use crate::formatting::{Formatted, apply_padding}; -use pulse_wire::prelude::*; +use pulse_sdk::prelude::*; pub struct PulseTradeApp { sock: Option, @@ -144,7 +144,8 @@ impl App for PulseTradeApp { ( LayoutItem::Widget(Size::Flex(1)), Box::new( - advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.strategy).await, + advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.strategy) + .await, ), ), ( diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index 5519562..3657ce2 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -1,5 +1,5 @@ +use pulse_sdk::{map_postcard_err, prelude::*}; use pulse_ui::state::State; -use pulse_wire::{map_postcard_err, prelude::*}; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, @@ -28,7 +28,7 @@ impl TerminalClient { pub async fn send( &mut self, - message: pulse_wire::terminal::TerminalClientMessage, + message: pulse_sdk::terminal::TerminalClientMessage, ) -> tokio::io::Result<()> { let msg = map_postcard_err(postcard::to_allocvec(&message))?; self.writer.write(&msg.len().to_le_bytes()).await?;