Strategy and Risk sdk
This commit is contained in:
@@ -7,4 +7,4 @@ edition = "2024"
|
|||||||
serde = { workspace = true }
|
serde = { workspace = true }
|
||||||
hypersdk = { workspace = true }
|
hypersdk = { workspace = true }
|
||||||
postcard = { workspace = true }
|
postcard = { workspace = true }
|
||||||
tokio = { workspace = true }
|
tokio = { workspace = true, features = ["io-std"]}
|
||||||
|
|||||||
+132
-5
@@ -1,25 +1,152 @@
|
|||||||
#[cfg(target_os = "macos")]
|
|
||||||
use std::path::PathBuf;
|
|
||||||
|
|
||||||
pub mod general;
|
pub mod general;
|
||||||
pub mod plugin;
|
pub mod plugin;
|
||||||
pub mod terminal;
|
pub mod terminal;
|
||||||
pub mod units;
|
pub mod units;
|
||||||
pub use hypersdk;
|
pub use hypersdk;
|
||||||
|
|
||||||
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
|
|
||||||
|
use crate::plugin::{RiskEngineMessage, StrategyEngineMessage};
|
||||||
|
|
||||||
pub mod prelude {
|
pub mod prelude {
|
||||||
pub use crate::general::*;
|
pub use crate::general::*;
|
||||||
pub use crate::plugin::*;
|
pub use crate::plugin::*;
|
||||||
pub use crate::server_path;
|
pub use crate::server_path;
|
||||||
pub use crate::terminal::*;
|
pub use crate::terminal::*;
|
||||||
pub use crate::units::*;
|
pub use crate::units::*;
|
||||||
|
|
||||||
pub use hypersdk;
|
pub use hypersdk;
|
||||||
|
pub use postcard;
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn server_path() -> PathBuf {
|
pub fn server_path() -> std::path::PathBuf {
|
||||||
PathBuf::from("/tmp/pulse-engine.sock")
|
std::path::PathBuf::from("/tmp/pulse-engine.sock")
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn map_postcard_err<T>(res: postcard::Result<T>) -> tokio::io::Result<T> {
|
pub fn map_postcard_err<T>(res: postcard::Result<T>) -> tokio::io::Result<T> {
|
||||||
res.map_err(|e| tokio::io::Error::new(std::io::ErrorKind::Other, e))
|
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::<usize>()];
|
||||||
|
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<String>) -> tokio::io::Result<()> {
|
||||||
|
unimplemented!("Strategy::command")
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn watchlist(&self, _watchlist: Vec<prelude::MarketItem>) -> 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<hypersdk::hypercore::Candle>,
|
||||||
|
) -> 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<String>) -> tokio::io::Result<()> {
|
||||||
|
unimplemented!("Risk::command")
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn watchlist(&self, _watchlist: Vec<prelude::MarketItem>) -> tokio::io::Result<()> {
|
||||||
|
unimplemented!("Risk::watchlist")
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn signal_request(
|
||||||
|
&self,
|
||||||
|
_signal: prelude::StrategySignal,
|
||||||
|
) -> tokio::io::Result<prelude::RiskSignal> {
|
||||||
|
unimplemented!("Risk::signal_request")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+12
-9
@@ -3,7 +3,7 @@ use crate::{
|
|||||||
terminal::MarketItem,
|
terminal::MarketItem,
|
||||||
units::Direction,
|
units::Direction,
|
||||||
};
|
};
|
||||||
use hypersdk::hypercore::{Candle, CandleInterval, Subscription};
|
use hypersdk::hypercore::{Candle, CandleInterval, Incoming, Subscription};
|
||||||
|
|
||||||
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
|
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
|
||||||
pub struct StrategyManifest {
|
pub struct StrategyManifest {
|
||||||
@@ -51,9 +51,7 @@ pub enum RiskMessage {
|
|||||||
|
|
||||||
GetWatchList,
|
GetWatchList,
|
||||||
|
|
||||||
Approve(Signal),
|
Signal(RiskSignal),
|
||||||
|
|
||||||
Reject { reason: String },
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
|
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
|
||||||
@@ -67,11 +65,7 @@ pub enum StrategyEngineMessage {
|
|||||||
args: Vec<String>,
|
args: Vec<String>,
|
||||||
},
|
},
|
||||||
|
|
||||||
CandleUpdate {
|
Incoming(Incoming),
|
||||||
symbol: String,
|
|
||||||
interval: CandleInterval,
|
|
||||||
candle: Candle,
|
|
||||||
},
|
|
||||||
|
|
||||||
Candlestick {
|
Candlestick {
|
||||||
symbol: String,
|
symbol: String,
|
||||||
@@ -98,3 +92,12 @@ pub struct StrategySignal {
|
|||||||
pub confidence: f32,
|
pub confidence: f32,
|
||||||
pub price: Option<f64>,
|
pub price: Option<f64>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
|
||||||
|
pub enum RiskSignal {
|
||||||
|
Approve(Signal),
|
||||||
|
Reject {
|
||||||
|
rejection_confidence: f32,
|
||||||
|
reason: String,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|||||||
@@ -178,8 +178,7 @@ impl StrategyEngine {
|
|||||||
engine.terminal_server.log_raw(log).await?;
|
engine.terminal_server.log_raw(log).await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
Some(RiskMessage::Approve(signal)) => {}
|
Some(RiskMessage::Signal(signal)) => {}
|
||||||
Some(RiskMessage::Reject { reason }) => {}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user