Improved engine architecture, and plugin architecture
This commit is contained in:
@@ -1,5 +1,8 @@
|
|||||||
|
pub mod plugin;
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
store::{accounts::AccountList, config::Config, plugin::StrategyEngine},
|
engine::plugin::StrategyEngine,
|
||||||
|
store::{accounts::AccountList, config::Config},
|
||||||
terminal::TerminalServer,
|
terminal::TerminalServer,
|
||||||
};
|
};
|
||||||
use pulse_wire::prelude::*;
|
use pulse_wire::prelude::*;
|
||||||
@@ -0,0 +1,174 @@
|
|||||||
|
use pulse_wire::prelude::*;
|
||||||
|
use std::{
|
||||||
|
path::PathBuf,
|
||||||
|
process::Stdio,
|
||||||
|
sync::{Arc, Weak},
|
||||||
|
};
|
||||||
|
use tokio::{
|
||||||
|
fs,
|
||||||
|
process::{Child, Command},
|
||||||
|
};
|
||||||
|
|
||||||
|
use crate::{
|
||||||
|
engine::Engine,
|
||||||
|
store::{plugin::Plugin, pulse_plugin},
|
||||||
|
};
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
pub struct StrategyEngine {
|
||||||
|
pub strategy: Arc<Plugin<StrategyEngineMessage, StrategyMessage, StrategyManifest>>,
|
||||||
|
pub risk: Arc<Plugin<RiskEngineMessage, RiskMessage, RiskManifest>>,
|
||||||
|
|
||||||
|
pub engine: Weak<Engine>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl StrategyEngine {
|
||||||
|
pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result<Self> {
|
||||||
|
let strategy = pulse_plugin(strategy_id)?;
|
||||||
|
let risk = pulse_plugin(risk_id)?;
|
||||||
|
|
||||||
|
let (strategy, strategy_manifest) = get_manifest_plugin_pair(
|
||||||
|
&strategy,
|
||||||
|
&strategy.join("strategy.bash"),
|
||||||
|
&fs::read(strategy.join("strategy.toml")).await?,
|
||||||
|
)?;
|
||||||
|
|
||||||
|
let (risk, risk_manifest) = get_manifest_plugin_pair(
|
||||||
|
&risk,
|
||||||
|
&risk.join("risk.bash"),
|
||||||
|
&fs::read(risk.join("risk.toml")).await?,
|
||||||
|
)?;
|
||||||
|
|
||||||
|
Ok(Self {
|
||||||
|
strategy: Arc::new(Plugin::new(strategy, strategy_manifest)),
|
||||||
|
risk: Arc::new(Plugin::new(risk, risk_manifest)),
|
||||||
|
engine: Weak::new(),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn initialize(mut self, engine: Weak<Engine>) -> Arc<Self> {
|
||||||
|
self.engine = engine;
|
||||||
|
|
||||||
|
Arc::new(self)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn run_strategy(&self) -> tokio::io::Result<()> {
|
||||||
|
let engine = self
|
||||||
|
.engine
|
||||||
|
.upgrade()
|
||||||
|
.expect("Failed to upgrade engine (StrategyEngine)");
|
||||||
|
|
||||||
|
let strategy = self.strategy.clone();
|
||||||
|
let risk = self.risk.clone();
|
||||||
|
|
||||||
|
loop {
|
||||||
|
match strategy.recv().await? {
|
||||||
|
None => {}
|
||||||
|
|
||||||
|
Some(StrategyMessage::Log(mut log)) => {
|
||||||
|
log.name.insert_str(0, "strategy::");
|
||||||
|
engine.terminal_server.log_raw(log).await?;
|
||||||
|
}
|
||||||
|
|
||||||
|
Some(StrategyMessage::Signal(signal)) => {
|
||||||
|
risk.send(&RiskEngineMessage::Signal(signal)).await?;
|
||||||
|
}
|
||||||
|
|
||||||
|
Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {}
|
||||||
|
|
||||||
|
Some(StrategyMessage::RequestOHLC {
|
||||||
|
symbol,
|
||||||
|
timeframe,
|
||||||
|
count,
|
||||||
|
}) => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn run_risk(&self) -> tokio::io::Result<()> {
|
||||||
|
let engine = self
|
||||||
|
.engine
|
||||||
|
.upgrade()
|
||||||
|
.expect("Failed to upgrade engine (StrategyEngine)");
|
||||||
|
|
||||||
|
let risk = self.risk.clone();
|
||||||
|
|
||||||
|
loop {
|
||||||
|
match risk.recv().await? {
|
||||||
|
None => {}
|
||||||
|
|
||||||
|
Some(RiskMessage::Log(mut log)) => {
|
||||||
|
log.name.insert_str(0, "strategy::");
|
||||||
|
engine.terminal_server.log_raw(log).await?;
|
||||||
|
}
|
||||||
|
|
||||||
|
Some(RiskMessage::Approve(signal)) => {}
|
||||||
|
Some(RiskMessage::Reject { reason }) => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn spawn(self: &Arc<Self>) {
|
||||||
|
let engine = self.clone();
|
||||||
|
|
||||||
|
tokio::spawn(async move { engine.run_strategy().await });
|
||||||
|
|
||||||
|
let engine = self.clone();
|
||||||
|
|
||||||
|
tokio::spawn(async move { engine.run_risk().await });
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn reload_strategy(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> {
|
||||||
|
let plugin = pulse_plugin(id)?;
|
||||||
|
|
||||||
|
let (child, manifest) = get_manifest_plugin_pair(
|
||||||
|
&plugin,
|
||||||
|
&plugin.join("strategy.bash"),
|
||||||
|
&fs::read(plugin.join("strategy.toml")).await?,
|
||||||
|
)?;
|
||||||
|
|
||||||
|
self.strategy.reload(child, manifest).await?;
|
||||||
|
|
||||||
|
let engine = self.clone();
|
||||||
|
tokio::spawn(async move { engine.run_strategy().await });
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn reload_risk(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> {
|
||||||
|
let plugin = pulse_plugin(id)?;
|
||||||
|
|
||||||
|
let (child, manifest) = get_manifest_plugin_pair(
|
||||||
|
&plugin,
|
||||||
|
&plugin.join("risk.bash"),
|
||||||
|
&fs::read(plugin.join("risk.toml")).await?,
|
||||||
|
)?;
|
||||||
|
|
||||||
|
self.risk.reload(child, manifest).await?;
|
||||||
|
|
||||||
|
let engine = self.clone();
|
||||||
|
tokio::spawn(async move { engine.run_risk().await });
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn get_manifest_plugin_pair<'de, M: serde::Deserialize<'de>>(
|
||||||
|
plugin_dir: &PathBuf,
|
||||||
|
plugin_path: &PathBuf,
|
||||||
|
manifest: &'de [u8],
|
||||||
|
) -> tokio::io::Result<(Child, M)> {
|
||||||
|
Ok((
|
||||||
|
Command::new("bash")
|
||||||
|
.arg(plugin_path)
|
||||||
|
.current_dir(plugin_dir)
|
||||||
|
.stdin(Stdio::piped())
|
||||||
|
.stdout(Stdio::piped())
|
||||||
|
.stderr(Stdio::inherit())
|
||||||
|
.spawn()?,
|
||||||
|
toml::from_slice(manifest)
|
||||||
|
.map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))?,
|
||||||
|
))
|
||||||
|
}
|
||||||
+3
-195
@@ -1,28 +1,14 @@
|
|||||||
use std::{
|
use std::marker::PhantomData;
|
||||||
marker::PhantomData,
|
|
||||||
path::PathBuf,
|
|
||||||
process::Stdio,
|
|
||||||
sync::{Arc, Weak},
|
|
||||||
};
|
|
||||||
|
|
||||||
use pulse_wire::{
|
use pulse_wire::PulseWire;
|
||||||
PulseWire,
|
|
||||||
plugin::{
|
|
||||||
RiskEngineMessage, RiskManifest, RiskMessage, StrategyEngineMessage, StrategyManifest,
|
|
||||||
StrategyMessage,
|
|
||||||
},
|
|
||||||
};
|
|
||||||
|
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
use tokio::{
|
use tokio::{
|
||||||
fs,
|
|
||||||
io::{AsyncReadExt, AsyncWriteExt},
|
io::{AsyncReadExt, AsyncWriteExt},
|
||||||
process::{Child, ChildStdout, Command},
|
process::{Child, ChildStdout},
|
||||||
sync::Mutex,
|
sync::Mutex,
|
||||||
};
|
};
|
||||||
|
|
||||||
use crate::{engine::Engine, store::pulse_plugin};
|
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub struct Plugin<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> {
|
pub struct Plugin<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> {
|
||||||
pub manifest: Mutex<M>,
|
pub manifest: Mutex<M>,
|
||||||
@@ -89,181 +75,3 @@ impl<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn get_manifest_plugin_pair<'de, M: Deserialize<'de>>(
|
|
||||||
plugin_dir: &PathBuf,
|
|
||||||
plugin_path: &PathBuf,
|
|
||||||
manifest: &'de [u8],
|
|
||||||
) -> tokio::io::Result<(Child, M)> {
|
|
||||||
Ok((
|
|
||||||
Command::new("bash")
|
|
||||||
.arg(plugin_path)
|
|
||||||
.current_dir(plugin_dir)
|
|
||||||
.stdin(Stdio::piped())
|
|
||||||
.stdout(Stdio::piped())
|
|
||||||
.stderr(Stdio::inherit())
|
|
||||||
.spawn()?,
|
|
||||||
toml::from_slice(manifest)
|
|
||||||
.map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))?,
|
|
||||||
))
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Debug)]
|
|
||||||
pub struct StrategyPair {
|
|
||||||
pub strategy: Plugin<StrategyEngineMessage, StrategyMessage, StrategyManifest>,
|
|
||||||
pub risk: Plugin<RiskEngineMessage, RiskMessage, RiskManifest>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl StrategyPair {
|
|
||||||
pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result<Arc<Self>> {
|
|
||||||
let strategy = pulse_plugin(strategy_id)?;
|
|
||||||
let risk = pulse_plugin(risk_id)?;
|
|
||||||
|
|
||||||
let (strategy, strategy_manifest) = get_manifest_plugin_pair(
|
|
||||||
&strategy,
|
|
||||||
&strategy.join("strategy.bash"),
|
|
||||||
&fs::read(strategy.join("strategy.toml")).await?,
|
|
||||||
)?;
|
|
||||||
|
|
||||||
let (risk, risk_manifest) = get_manifest_plugin_pair(
|
|
||||||
&risk,
|
|
||||||
&risk.join("risk.bash"),
|
|
||||||
&fs::read(risk.join("risk.toml")).await?,
|
|
||||||
)?;
|
|
||||||
|
|
||||||
Ok(Arc::new(Self {
|
|
||||||
strategy: Plugin::new(strategy, strategy_manifest),
|
|
||||||
risk: Plugin::new(risk, risk_manifest),
|
|
||||||
}))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Debug)]
|
|
||||||
pub struct StrategyEngine {
|
|
||||||
pub pair: Mutex<Arc<StrategyPair>>,
|
|
||||||
pub engine: Weak<Engine>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl StrategyEngine {
|
|
||||||
pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result<Self> {
|
|
||||||
Ok(Self {
|
|
||||||
pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?),
|
|
||||||
engine: Weak::new(),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn initialize(mut self, engine: Weak<Engine>) -> Arc<Self> {
|
|
||||||
self.engine = engine;
|
|
||||||
|
|
||||||
Arc::new(self)
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn run_strategy(&self) -> tokio::io::Result<()> {
|
|
||||||
let engine = self
|
|
||||||
.engine
|
|
||||||
.upgrade()
|
|
||||||
.expect("Failed to upgrade engine (StrategyEngine)");
|
|
||||||
|
|
||||||
let pair = self.pair.lock().await.clone();
|
|
||||||
|
|
||||||
loop {
|
|
||||||
match pair.strategy.recv().await? {
|
|
||||||
None => {}
|
|
||||||
|
|
||||||
Some(StrategyMessage::Log(mut log)) => {
|
|
||||||
log.name.insert_str(0, "strategy::");
|
|
||||||
engine.terminal_server.log_raw(log).await?;
|
|
||||||
}
|
|
||||||
|
|
||||||
Some(StrategyMessage::Signal(signal)) => {
|
|
||||||
pair.risk.send(&RiskEngineMessage::Signal(signal)).await?;
|
|
||||||
}
|
|
||||||
|
|
||||||
Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {}
|
|
||||||
|
|
||||||
Some(StrategyMessage::RequestOHLC {
|
|
||||||
symbol,
|
|
||||||
timeframe,
|
|
||||||
count,
|
|
||||||
}) => {}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn run_risk(&self) -> tokio::io::Result<()> {
|
|
||||||
let engine = self
|
|
||||||
.engine
|
|
||||||
.upgrade()
|
|
||||||
.expect("Failed to upgrade engine (StrategyEngine)");
|
|
||||||
|
|
||||||
let pair = self.pair.lock().await.clone();
|
|
||||||
|
|
||||||
loop {
|
|
||||||
match pair.risk.recv().await? {
|
|
||||||
None => {}
|
|
||||||
|
|
||||||
Some(RiskMessage::Log(mut log)) => {
|
|
||||||
log.name.insert_str(0, "strategy::");
|
|
||||||
engine.terminal_server.log_raw(log).await?;
|
|
||||||
}
|
|
||||||
|
|
||||||
Some(RiskMessage::Approve(signal)) => {}
|
|
||||||
Some(RiskMessage::Reject { reason }) => {}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn spawn(self: &Arc<Self>) {
|
|
||||||
let engine = self.clone();
|
|
||||||
|
|
||||||
tokio::spawn(async move { engine.run_strategy().await });
|
|
||||||
|
|
||||||
let engine = self.clone();
|
|
||||||
|
|
||||||
tokio::spawn(async move { engine.run_risk().await });
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn reload_strategy(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> {
|
|
||||||
{
|
|
||||||
let pair = self.pair.lock().await;
|
|
||||||
|
|
||||||
let plugin = pulse_plugin(id)?;
|
|
||||||
|
|
||||||
let (child, manifest) = get_manifest_plugin_pair(
|
|
||||||
&plugin,
|
|
||||||
&plugin.join("strategy.bash"),
|
|
||||||
&fs::read(plugin.join("strategy.toml")).await?,
|
|
||||||
)?;
|
|
||||||
|
|
||||||
pair.strategy.reload(child, manifest).await?;
|
|
||||||
}
|
|
||||||
|
|
||||||
let engine = self.clone();
|
|
||||||
tokio::spawn(async move { engine.run_strategy().await });
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
|
|
||||||
pub async fn reload_risk(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> {
|
|
||||||
{
|
|
||||||
let pair = self.pair.lock().await;
|
|
||||||
|
|
||||||
let plugin = pulse_plugin(id)?;
|
|
||||||
|
|
||||||
let (child, manifest) = get_manifest_plugin_pair(
|
|
||||||
&plugin,
|
|
||||||
&plugin.join("risk.bash"),
|
|
||||||
&fs::read(plugin.join("risk.toml")).await?,
|
|
||||||
)?;
|
|
||||||
|
|
||||||
pair.risk.reload(child, manifest).await?;
|
|
||||||
}
|
|
||||||
|
|
||||||
let engine = self.clone();
|
|
||||||
tokio::spawn(async move { engine.run_risk().await });
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -61,13 +61,12 @@ impl TerminalServer {
|
|||||||
|
|
||||||
async fn initialize_client(self: &Arc<Self>, id: &usize) -> tokio::io::Result<()> {
|
async fn initialize_client(self: &Arc<Self>, id: &usize) -> tokio::io::Result<()> {
|
||||||
let engine = self.get_engine();
|
let engine = self.get_engine();
|
||||||
let pair = engine.strategy.pair.lock().await;
|
|
||||||
|
|
||||||
self.send_to(
|
self.send_to(
|
||||||
id,
|
id,
|
||||||
pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
|
pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
|
||||||
strategy: pair.strategy.manifest.lock().await.clone(),
|
strategy: engine.strategy.strategy.manifest.lock().await.clone(),
|
||||||
risk: pair.risk.manifest.lock().await.clone(),
|
risk: engine.strategy.risk.manifest.lock().await.clone(),
|
||||||
mode: Mode::Auto,
|
mode: Mode::Auto,
|
||||||
state: ItemState::Running,
|
state: ItemState::Running,
|
||||||
cooldown: TimeFrame::M15,
|
cooldown: TimeFrame::M15,
|
||||||
|
|||||||
Reference in New Issue
Block a user