Running strategy engine
This commit is contained in:
@@ -9,23 +9,25 @@ use tokio::{sync::Mutex, task::JoinHandle};
|
|||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct Engine {
|
pub struct Engine {
|
||||||
pub terminal_server: Arc<TerminalServer>,
|
pub terminal_server: Arc<TerminalServer>,
|
||||||
|
pub strategy: Arc<StrategyEngine>,
|
||||||
pub config: Arc<Mutex<Config>>,
|
pub config: Arc<Mutex<Config>>,
|
||||||
pub accounts: Arc<Mutex<AccountList>>,
|
pub accounts: Arc<Mutex<AccountList>>,
|
||||||
pub strategy: Arc<StrategyEngine>,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Engine {
|
impl Engine {
|
||||||
pub async fn new() -> tokio::io::Result<Arc<Self>> {
|
pub async fn new() -> tokio::io::Result<Arc<Self>> {
|
||||||
let config = Config::new().await?;
|
let config = Config::new().await?;
|
||||||
|
|
||||||
|
let strategy = StrategyEngine::new(&config.strategy, &config.risk).await?;
|
||||||
|
|
||||||
let accounts = Arc::new(Mutex::new(AccountList::new().await?));
|
let accounts = Arc::new(Mutex::new(AccountList::new().await?));
|
||||||
let strategy = Arc::new(StrategyEngine::new(&config.strategy, &config.risk).await?);
|
|
||||||
let config = Arc::new(Mutex::new(config));
|
let config = Arc::new(Mutex::new(config));
|
||||||
|
|
||||||
Ok(Arc::new_cyclic(|engine| Self {
|
Ok(Arc::new_cyclic(|engine| Self {
|
||||||
terminal_server: TerminalServer::new(engine.clone()),
|
terminal_server: TerminalServer::new(engine.clone()),
|
||||||
|
strategy: strategy.initialize(engine.clone()),
|
||||||
config,
|
config,
|
||||||
accounts,
|
accounts,
|
||||||
strategy,
|
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -11,6 +11,8 @@ async fn main() -> tokio::io::Result<()> {
|
|||||||
|
|
||||||
let broadcaster = engine.spawn_broadcaster().await;
|
let broadcaster = engine.spawn_broadcaster().await;
|
||||||
|
|
||||||
|
engine.strategy.spawn().await;
|
||||||
|
|
||||||
engine.run_engine().await?;
|
engine.run_engine().await?;
|
||||||
|
|
||||||
let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster);
|
let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster);
|
||||||
@@ -18,5 +20,18 @@ async fn main() -> tokio::io::Result<()> {
|
|||||||
terminal_server??;
|
terminal_server??;
|
||||||
broadcaster??;
|
broadcaster??;
|
||||||
|
|
||||||
|
let strategy = engine.strategy.strategy_handle.lock().await.take();
|
||||||
|
let risk = engine.strategy.risk_handle.lock().await.take();
|
||||||
|
|
||||||
|
drop(engine);
|
||||||
|
|
||||||
|
if let Some(strategy) = strategy {
|
||||||
|
strategy.into_future().await??;
|
||||||
|
}
|
||||||
|
|
||||||
|
if let Some(risk) = risk {
|
||||||
|
risk.into_future().await??;
|
||||||
|
}
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,4 +1,7 @@
|
|||||||
use std::marker::PhantomData;
|
use std::{
|
||||||
|
marker::PhantomData,
|
||||||
|
sync::{Arc, Weak},
|
||||||
|
};
|
||||||
|
|
||||||
use pulse_wire::{
|
use pulse_wire::{
|
||||||
PulseWire,
|
PulseWire,
|
||||||
@@ -16,7 +19,7 @@ use tokio::{
|
|||||||
task::JoinHandle,
|
task::JoinHandle,
|
||||||
};
|
};
|
||||||
|
|
||||||
use crate::store::pulse_plugin;
|
use crate::{engine::Engine, store::pulse_plugin};
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub struct Plugin<S: PulseWire, R: PulseWire> {
|
pub struct Plugin<S: PulseWire, R: PulseWire> {
|
||||||
@@ -108,8 +111,10 @@ impl StrategyPair {
|
|||||||
pub struct StrategyEngine {
|
pub struct StrategyEngine {
|
||||||
pub pair: Mutex<StrategyPair>,
|
pub pair: Mutex<StrategyPair>,
|
||||||
|
|
||||||
pub strategy_handle: Mutex<Option<JoinHandle<()>>>,
|
pub strategy_handle: Mutex<Option<JoinHandle<tokio::io::Result<()>>>>,
|
||||||
pub risk_handle: Mutex<Option<JoinHandle<()>>>,
|
pub risk_handle: Mutex<Option<JoinHandle<tokio::io::Result<()>>>>,
|
||||||
|
|
||||||
|
pub engine: Weak<Engine>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl StrategyEngine {
|
impl StrategyEngine {
|
||||||
@@ -118,9 +123,49 @@ impl StrategyEngine {
|
|||||||
pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?),
|
pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?),
|
||||||
strategy_handle: Mutex::new(None),
|
strategy_handle: Mutex::new(None),
|
||||||
risk_handle: Mutex::new(None),
|
risk_handle: Mutex::new(None),
|
||||||
|
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)");
|
||||||
|
|
||||||
|
loop {}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn run_risk(&self) -> tokio::io::Result<()> {
|
||||||
|
let engine = self
|
||||||
|
.engine
|
||||||
|
.upgrade()
|
||||||
|
.expect("Failed to upgrade engine (StrategyEngine)");
|
||||||
|
|
||||||
|
loop {}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn spawn(self: &Arc<Self>) {
|
||||||
|
let engine = self.clone();
|
||||||
|
|
||||||
|
*self.strategy_handle.lock().await =
|
||||||
|
Some(tokio::spawn(async move { engine.run_strategy().await }));
|
||||||
|
|
||||||
|
let engine = self.clone();
|
||||||
|
|
||||||
|
*self.risk_handle.lock().await = Some(tokio::spawn(async move { engine.run_risk().await }));
|
||||||
|
}
|
||||||
|
|
||||||
pub async fn reload(&self, strategy_id: &str, risk_id: &str) -> tokio::io::Result<()> {
|
pub async fn reload(&self, strategy_id: &str, risk_id: &str) -> tokio::io::Result<()> {
|
||||||
if let Some(handle) = &*self.strategy_handle.lock().await {
|
if let Some(handle) = &*self.strategy_handle.lock().await {
|
||||||
handle.abort();
|
handle.abort();
|
||||||
|
|||||||
Reference in New Issue
Block a user