Strategy spawning and maniest
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
use crate::{
|
||||
store::{accounts::AccountList, config::Config},
|
||||
store::{accounts::AccountList, config::Config, plugin::StrategyPair},
|
||||
terminal::TerminalServer,
|
||||
};
|
||||
use pulse_wire::prelude::*;
|
||||
@@ -11,17 +11,23 @@ pub struct Engine {
|
||||
pub terminal_server: Arc<TerminalServer>,
|
||||
pub config: Arc<Mutex<Config>>,
|
||||
pub accounts: Arc<Mutex<AccountList>>,
|
||||
pub strategy: Arc<Mutex<StrategyPair>>,
|
||||
}
|
||||
|
||||
impl Engine {
|
||||
pub async fn new() -> tokio::io::Result<Arc<Self>> {
|
||||
let config = Arc::new(Mutex::new(Config::new().await?));
|
||||
let config = Config::new().await?;
|
||||
let accounts = Arc::new(Mutex::new(AccountList::new().await?));
|
||||
let strategy = Arc::new(Mutex::new(
|
||||
StrategyPair::new(&config.strategy, &config.risk).await?,
|
||||
));
|
||||
let config = Arc::new(Mutex::new(config));
|
||||
|
||||
Ok(Arc::new_cyclic(|engine| Self {
|
||||
terminal_server: TerminalServer::new(engine.clone()),
|
||||
config,
|
||||
accounts,
|
||||
strategy,
|
||||
}))
|
||||
}
|
||||
|
||||
|
||||
@@ -6,6 +6,8 @@ pub struct WatchList {
|
||||
#[derive(Debug, Default, serde::Serialize, serde::Deserialize)]
|
||||
pub struct Config {
|
||||
pub watchlist: WatchList,
|
||||
pub strategy: String,
|
||||
pub risk: String,
|
||||
}
|
||||
|
||||
impl Config {
|
||||
|
||||
@@ -2,10 +2,20 @@ use std::marker::PhantomData;
|
||||
|
||||
use pulse_wire::{
|
||||
PulseWire,
|
||||
plugin::{RiskEngineMessage, RiskMessage, StrategyEngineMessage, StrategyMessage},
|
||||
plugin::{
|
||||
RiskEngineMessage, RiskManifest, RiskMessage, StrategyEngineMessage, StrategyManifest,
|
||||
StrategyMessage,
|
||||
},
|
||||
};
|
||||
use tokio::{
|
||||
fs,
|
||||
io::{AsyncReadExt, AsyncWriteExt},
|
||||
process::Command,
|
||||
};
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
|
||||
use crate::store::pulse_plugin;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct Plugin<S: PulseWire, R: PulseWire> {
|
||||
pub process: tokio::process::Child,
|
||||
pub _p: (PhantomData<S>, PhantomData<R>),
|
||||
@@ -53,7 +63,40 @@ impl<S: PulseWire, R: PulseWire> Plugin<S, R> {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct StrategyPair {
|
||||
pub strategy: Plugin<StrategyEngineMessage, StrategyMessage>,
|
||||
pub risk: Plugin<RiskEngineMessage, RiskMessage>,
|
||||
|
||||
pub strategy_manifest: StrategyManifest,
|
||||
pub risk_manifest: RiskManifest,
|
||||
}
|
||||
|
||||
impl StrategyPair {
|
||||
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)?;
|
||||
|
||||
Ok(Self {
|
||||
strategy: Plugin::new(
|
||||
Command::new("bash")
|
||||
.arg(strategy.join("strategy.bash"))
|
||||
.current_dir(&strategy)
|
||||
.spawn()?,
|
||||
),
|
||||
risk: Plugin::new(
|
||||
Command::new("bash")
|
||||
.arg(strategy.join("risk.bash"))
|
||||
.current_dir(&strategy)
|
||||
.spawn()?,
|
||||
),
|
||||
strategy_manifest: toml::from_slice(&fs::read(strategy.join("strategy.toml")).await?)
|
||||
.map_err(|v| {
|
||||
tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string())
|
||||
})?,
|
||||
risk_manifest: toml::from_slice(&fs::read(risk.join("risk.toml")).await?).map_err(
|
||||
|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()),
|
||||
)?,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user