pub mod general; pub mod strategy; pub mod terminal; pub use hypersdk; use hypersdk::Decimal; use rust_decimal::RoundingStrategy; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use crate::{general::LogKind, strategy::StrategyEngineMessage}; pub mod prelude { pub use crate::Strategy; pub use crate::general::*; pub use crate::strategy::*; pub use crate::terminal::*; pub use crate::{round_price, server_path}; pub use hypersdk; pub use postcard; } 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 fn round_price(price: Decimal) -> Decimal { let digits = price.trunc().to_string().len() as u32; let decimal_places = 5_i32 - digits as i32; if decimal_places < 0 { let factor = Decimal::from(10_u64.pow((-decimal_places) as u32)); (price / factor).round() * factor } else { price.round_dp_with_strategy( decimal_places as u32, RoundingStrategy::MidpointAwayFromZero, ) } } 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(()) } #[allow(async_fn_in_trait)] pub trait Strategy { 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 log(&self, kind: LogKind, name: &str, message: &str) -> tokio::io::Result<()> { self.send(&prelude::StrategyMessage::Log(prelude::EventLog { kind, name: name.to_owned(), message: message.to_owned(), })) .await } 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") } }