From 9c29695728fb1b3157e68b4936b49b0d5bc28c48 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Wed, 29 Jul 2026 01:35:36 +0200 Subject: [PATCH] Strategy and Risk sdk --- pulse-sdk/Cargo.toml | 2 +- pulse-sdk/src/lib.rs | 137 ++++++++++++++++++++++++++++++++++-- pulse-sdk/src/plugin.rs | 21 +++--- src/engine/engine/plugin.rs | 3 +- 4 files changed, 146 insertions(+), 17 deletions(-) diff --git a/pulse-sdk/Cargo.toml b/pulse-sdk/Cargo.toml index bb06f35..f8b9607 100644 --- a/pulse-sdk/Cargo.toml +++ b/pulse-sdk/Cargo.toml @@ -7,4 +7,4 @@ edition = "2024" serde = { workspace = true } hypersdk = { workspace = true } postcard = { workspace = true } -tokio = { workspace = true } +tokio = { workspace = true, features = ["io-std"]} diff --git a/pulse-sdk/src/lib.rs b/pulse-sdk/src/lib.rs index dd2211e..189435a 100644 --- a/pulse-sdk/src/lib.rs +++ b/pulse-sdk/src/lib.rs @@ -1,25 +1,152 @@ -#[cfg(target_os = "macos")] -use std::path::PathBuf; - pub mod general; pub mod plugin; pub mod terminal; pub mod units; pub use hypersdk; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; + +use crate::plugin::{RiskEngineMessage, StrategyEngineMessage}; + pub mod prelude { pub use crate::general::*; pub use crate::plugin::*; pub use crate::server_path; pub use crate::terminal::*; pub use crate::units::*; + pub use hypersdk; + pub use postcard; } -pub fn server_path() -> PathBuf { - PathBuf::from("/tmp/pulse-engine.sock") +pub fn server_path() -> std::path::PathBuf { + std::path::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)) } + +pub async fn send_raw(data: &[u8]) -> tokio::io::Result<()> { + let mut stdout = tokio::io::stdout(); + + stdout.write_all(&data.len().to_le_bytes()).await?; + stdout.write_all(data).await?; + stdout.flush().await?; + + 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 + } + }; +} + +#[allow(async_fn_in_trait)] +pub trait Strategy { + engine_methods!(prelude::StrategyMessage); + + async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> { + match map_postcard_err(postcard::from_bytes(&data))? { + StrategyEngineMessage::Initialize => self.initialize().await, + StrategyEngineMessage::Command { command, args } => self.command(command, args).await, + StrategyEngineMessage::WatchList(watchlist) => self.watchlist(watchlist).await, + StrategyEngineMessage::Incoming(incoming) => self.incoming(incoming).await, + StrategyEngineMessage::Candlestick { + symbol, + interval, + candles, + } => self.candlestick(symbol, interval, candles).await, + } + } + + async fn initialize(&self) -> tokio::io::Result<()> { + unimplemented!("Strategy::initialize") + } + + async fn command(&self, _command: String, _args: Vec) -> tokio::io::Result<()> { + unimplemented!("Strategy::command") + } + + async fn watchlist(&self, _watchlist: Vec) -> tokio::io::Result<()> { + unimplemented!("Strategy::watchlist") + } + + async fn incoming(&self, _incoming: hypersdk::hypercore::Incoming) -> tokio::io::Result<()> { + unimplemented!("Strategy::event") + } + + async fn candlestick( + &self, + _symbol: String, + _interval: hypersdk::hypercore::CandleInterval, + _candles: Vec, + ) -> tokio::io::Result<()> { + unimplemented!("Strategy::candlestick") + } +} + +#[allow(async_fn_in_trait)] +pub trait Risk { + engine_methods!(prelude::RiskMessage); + + async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> { + match map_postcard_err(postcard::from_bytes(&data))? { + RiskEngineMessage::Initialize => self.initialize().await, + RiskEngineMessage::Command { command, args } => self.command(command, args).await, + RiskEngineMessage::WatchList(watchlist) => self.watchlist(watchlist).await, + RiskEngineMessage::Signal(signal) => { + self.send(&prelude::RiskMessage::Signal( + self.signal_request(signal).await?, + )) + .await + } + } + } + + async fn initialize(&self) -> tokio::io::Result<()> { + unimplemented!("Risk::initialize") + } + + async fn command(&self, _command: String, _args: Vec) -> tokio::io::Result<()> { + unimplemented!("Risk::command") + } + + async fn watchlist(&self, _watchlist: Vec) -> tokio::io::Result<()> { + unimplemented!("Risk::watchlist") + } + + async fn signal_request( + &self, + _signal: prelude::StrategySignal, + ) -> tokio::io::Result { + unimplemented!("Risk::signal_request") + } +} diff --git a/pulse-sdk/src/plugin.rs b/pulse-sdk/src/plugin.rs index 821a226..fc54bf7 100644 --- a/pulse-sdk/src/plugin.rs +++ b/pulse-sdk/src/plugin.rs @@ -3,7 +3,7 @@ use crate::{ terminal::MarketItem, units::Direction, }; -use hypersdk::hypercore::{Candle, CandleInterval, Subscription}; +use hypersdk::hypercore::{Candle, CandleInterval, Incoming, Subscription}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct StrategyManifest { @@ -51,9 +51,7 @@ pub enum RiskMessage { GetWatchList, - Approve(Signal), - - Reject { reason: String }, + Signal(RiskSignal), } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] @@ -67,11 +65,7 @@ pub enum StrategyEngineMessage { args: Vec, }, - CandleUpdate { - symbol: String, - interval: CandleInterval, - candle: Candle, - }, + Incoming(Incoming), Candlestick { symbol: String, @@ -98,3 +92,12 @@ pub struct StrategySignal { pub confidence: f32, pub price: Option, } + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub enum RiskSignal { + Approve(Signal), + Reject { + rejection_confidence: f32, + reason: String, + }, +} diff --git a/src/engine/engine/plugin.rs b/src/engine/engine/plugin.rs index 2c30d1b..5c69307 100644 --- a/src/engine/engine/plugin.rs +++ b/src/engine/engine/plugin.rs @@ -178,8 +178,7 @@ impl StrategyEngine { engine.terminal_server.log_raw(log).await?; } - Some(RiskMessage::Approve(signal)) => {} - Some(RiskMessage::Reject { reason }) => {} + Some(RiskMessage::Signal(signal)) => {} } } }