From 535c6a2964a28c9e7db2c6c894f2ca2d520e0d86 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Thu, 30 Jul 2026 05:35:03 +0200 Subject: [PATCH] Rename --- pulse-sdk/src/lib.rs | 60 ++++++++++++++++------------------- src/engine/engine/command.rs | 4 +-- src/engine/engine/strategy.rs | 25 +++++++-------- src/engine/terminal.rs | 2 +- 4 files changed, 42 insertions(+), 49 deletions(-) diff --git a/pulse-sdk/src/lib.rs b/pulse-sdk/src/lib.rs index 24b83cb..6d644bd 100644 --- a/pulse-sdk/src/lib.rs +++ b/pulse-sdk/src/lib.rs @@ -35,41 +35,37 @@ pub async fn send_raw(data: &[u8]) -> tokio::io::Result<()> { Ok(()) } -macro_rules! engine_methods { - ($t:ty) => { - async fn start(&self) -> tokio::io::Result<()> { - let mut stdin = tokio::io::stdin(); - - loop { - let mut len_buf = [0u8; size_of::()]; - let size = stdin.read_exact(&mut len_buf).await?; - - let len = usize::from_le_bytes(len_buf); - - if size == 0 || len == 0 { - break Ok(()); - } - - let mut buffer = vec![0u8; len]; - - stdin.read_exact(&mut buffer).await?; - - self.on_raw(&buffer).await?; - } - } - - async fn send(&self, msg: &$t) -> tokio::io::Result<()> { - $crate::send_raw(&$crate::map_postcard_err( - $crate::prelude::postcard::to_allocvec(msg), - )?) - .await - } - }; -} +use std::sync::{ + Arc, + atomic::{AtomicBool, Ordering}, +}; #[allow(async_fn_in_trait)] pub trait Strategy { - engine_methods!(prelude::StrategyMessage); + async fn start(&self) -> tokio::io::Result<()> { + let mut stdin = tokio::io::stdin(); + + loop { + let mut len_buf = [0u8; size_of::()]; + let size = stdin.read_exact(&mut len_buf).await?; + + let len = usize::from_le_bytes(len_buf); + + if size == 0 || len == 0 { + break Ok(()); + } + + let mut buffer = vec![0u8; len]; + + stdin.read_exact(&mut buffer).await?; + + self.on_raw(&buffer).await?; + } + } + + async fn send(&self, msg: &prelude::StrategyMessage) -> tokio::io::Result<()> { + send_raw(&map_postcard_err(prelude::postcard::to_allocvec(msg))?).await + } async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> { match map_postcard_err(postcard::from_bytes(&data))? { diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index 1facea8..3644ad3 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -60,7 +60,7 @@ impl Engine { .await; } - self.strategy.reload_strategy(id.as_str()).await?; + self.strategy.reload(id.as_str()).await?; self.config.lock().await.strategy = id; }); @@ -75,7 +75,7 @@ impl Engine { self.terminal_server.broadcast( pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { - strategy: self.strategy.strategy.manifest.lock().await.clone(), + strategy: self.strategy.child.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index 576dc68..d45db97 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -19,9 +19,8 @@ use crate::{ }; pub struct StrategyEngine { - pub strategy: Arc, pub engine: Weak, - + pub child: Arc, pub ws: WebSocket, pub subscriptions: Mutex>, } @@ -43,7 +42,7 @@ impl StrategyEngine { )?; Ok(Self { - strategy: Arc::new(StrategyChild::new(strategy, strategy_manifest)), + child: Arc::new(StrategyChild::new(strategy, strategy_manifest)), engine: Weak::new(), ws: hypercore::mainnet_ws(), subscriptions: Mutex::new(HashSet::new()), @@ -56,22 +55,20 @@ impl StrategyEngine { Arc::new(self) } - pub async fn run_strategy(&self) -> anyhow::Result<()> { + pub async fn run(&self) -> anyhow::Result<()> { let engine = self .engine .upgrade() .expect("Failed to upgrade engine (StrategyEngine)"); - self.strategy - .send(&StrategyEngineMessage::Initialize) - .await?; + self.child.send(&StrategyEngineMessage::Initialize).await?; loop { - match self.strategy.recv().await? { + match self.child.recv().await? { None => {} Some(StrategyMessage::GetWatchList) => { - self.strategy + self.child .send(&StrategyEngineMessage::WatchList( engine.watch_list.lock().await.clone().items, )) @@ -153,7 +150,7 @@ impl StrategyEngine { let start_time = now.saturating_sub(interval_ms * count as u64); - self.strategy + self.child .send(&StrategyEngineMessage::Candlestick { candles: client .candle_snapshot(&symbol, interval, start_time, now) @@ -170,10 +167,10 @@ impl StrategyEngine { pub async fn spawn(self: &Arc) { let engine = self.clone(); - tokio::spawn(async move { engine.run_strategy().await }); + tokio::spawn(async move { engine.run().await }); } - pub async fn reload_strategy(self: &Arc, id: &str) -> tokio::io::Result<()> { + pub async fn reload(self: &Arc, id: &str) -> tokio::io::Result<()> { let strategy = pulse_strategy(id)?; let (child, manifest) = get_manifest( @@ -182,10 +179,10 @@ impl StrategyEngine { &fs::read(strategy.join("strategy.toml")).await?, )?; - self.strategy.reload(child, manifest).await?; + self.child.reload(child, manifest).await?; let engine = self.clone(); - tokio::spawn(async move { engine.run_strategy().await }); + tokio::spawn(async move { engine.run().await }); Ok(()) } diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs index 97e8150..a08f521 100644 --- a/src/engine/terminal.rs +++ b/src/engine/terminal.rs @@ -65,7 +65,7 @@ impl TerminalServer { self.send_to( id, pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { - strategy: engine.strategy.strategy.manifest.lock().await.clone(), + strategy: engine.strategy.child.manifest.lock().await.clone(), mode: Mode::Auto, state: ItemState::Running, cooldown: engine.config.lock().await.cooldown,