From 73a47c43f73b2cc3512c523d1778819116275526 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 02:38:51 +0200 Subject: [PATCH] Running strategy engine --- src/engine/engine.rs | 8 +++--- src/engine/main.rs | 15 +++++++++++ src/engine/store/plugin.rs | 53 +++++++++++++++++++++++++++++++++++--- 3 files changed, 69 insertions(+), 7 deletions(-) diff --git a/src/engine/engine.rs b/src/engine/engine.rs index c3ce821..798ebb3 100644 --- a/src/engine/engine.rs +++ b/src/engine/engine.rs @@ -9,23 +9,25 @@ use tokio::{sync::Mutex, task::JoinHandle}; #[derive(Debug, Clone)] pub struct Engine { pub terminal_server: Arc, + pub strategy: Arc, pub config: Arc>, pub accounts: Arc>, - pub strategy: Arc, } impl Engine { pub async fn new() -> tokio::io::Result> { let config = Config::new().await?; + + let strategy = StrategyEngine::new(&config.strategy, &config.risk).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)); Ok(Arc::new_cyclic(|engine| Self { terminal_server: TerminalServer::new(engine.clone()), + strategy: strategy.initialize(engine.clone()), config, accounts, - strategy, })) } diff --git a/src/engine/main.rs b/src/engine/main.rs index f4f7f64..184a811 100644 --- a/src/engine/main.rs +++ b/src/engine/main.rs @@ -11,6 +11,8 @@ async fn main() -> tokio::io::Result<()> { let broadcaster = engine.spawn_broadcaster().await; + engine.strategy.spawn().await; + engine.run_engine().await?; let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster); @@ -18,5 +20,18 @@ async fn main() -> tokio::io::Result<()> { terminal_server??; 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(()) } diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index a589dfe..de0fe97 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -1,4 +1,7 @@ -use std::marker::PhantomData; +use std::{ + marker::PhantomData, + sync::{Arc, Weak}, +}; use pulse_wire::{ PulseWire, @@ -16,7 +19,7 @@ use tokio::{ task::JoinHandle, }; -use crate::store::pulse_plugin; +use crate::{engine::Engine, store::pulse_plugin}; #[derive(Debug)] pub struct Plugin { @@ -108,8 +111,10 @@ impl StrategyPair { pub struct StrategyEngine { pub pair: Mutex, - pub strategy_handle: Mutex>>, - pub risk_handle: Mutex>>, + pub strategy_handle: Mutex>>>, + pub risk_handle: Mutex>>>, + + pub engine: Weak, } impl StrategyEngine { @@ -118,9 +123,49 @@ impl StrategyEngine { pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), strategy_handle: Mutex::new(None), risk_handle: Mutex::new(None), + engine: Weak::new(), }) } + pub fn initialize(mut self, engine: Weak) -> Arc { + 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) { + 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<()> { if let Some(handle) = &*self.strategy_handle.lock().await { handle.abort();