114 lines
3.4 KiB
Rust
114 lines
3.4 KiB
Rust
pub mod general;
|
|
pub mod strategy;
|
|
pub mod terminal;
|
|
pub use hypersdk;
|
|
|
|
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::server_path;
|
|
pub use crate::strategy::*;
|
|
pub use crate::terminal::*;
|
|
|
|
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<T>(res: postcard::Result<T>) -> tokio::io::Result<T> {
|
|
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(())
|
|
}
|
|
|
|
#[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::<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: &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<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")
|
|
}
|
|
}
|