From 4bf47afaa404075aef7b62c16e8e7ab214519c57 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sun, 26 Jul 2026 02:56:08 +0200 Subject: [PATCH] Fix deadlock --- src/engine/store/plugin.rs | 85 ++++++++++++++++++++++++++------------ 1 file changed, 58 insertions(+), 27 deletions(-) diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index de0fe97..3b2bcb7 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -71,11 +71,11 @@ impl Plugin { #[derive(Debug)] pub struct StrategyPair { - pub strategy: Plugin, - pub risk: Plugin, + pub strategy: Mutex>, + pub risk: Mutex>, - pub strategy_manifest: StrategyManifest, - pub risk_manifest: RiskManifest, + pub strategy_manifest: Mutex, + pub risk_manifest: Mutex, } impl StrategyPair { @@ -84,32 +84,35 @@ impl StrategyPair { let risk = pulse_plugin(risk_id)?; Ok(Self { - strategy: Plugin::new( + strategy: Mutex::new(Plugin::new( Command::new("bash") .arg(strategy.join("strategy.bash")) .current_dir(&strategy) .spawn()?, - ), - risk: Plugin::new( + )), + risk: Mutex::new(Plugin::new( Command::new("bash") .arg(strategy.join("risk.bash")) .current_dir(&strategy) .spawn()?, + )), + strategy_manifest: Mutex::new( + 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: Mutex::new( + toml::from_slice(&fs::read(risk.join("risk.toml")).await?).map_err(|v| { + tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()) + })?, ), - 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()), - )?, }) } } #[derive(Debug)] pub struct StrategyEngine { - pub pair: Mutex, + pub pair: StrategyPair, pub strategy_handle: Mutex>>>, pub risk_handle: Mutex>>>, @@ -120,7 +123,7 @@ pub struct StrategyEngine { impl StrategyEngine { pub async fn new(strategy_id: &str, risk_id: &str) -> tokio::io::Result { Ok(Self { - pair: Mutex::new(StrategyPair::new(strategy_id, risk_id).await?), + pair: StrategyPair::new(strategy_id, risk_id).await?, strategy_handle: Mutex::new(None), risk_handle: Mutex::new(None), engine: Weak::new(), @@ -166,21 +169,49 @@ impl StrategyEngine { *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(); + pub async fn reload( + self: &Arc, + strategy_id: &str, + risk_id: &str, + ) -> tokio::io::Result<()> { + { + if let Some(handle) = &*self.strategy_handle.lock().await { + handle.abort(); + } + + if let Some(handle) = &*self.risk_handle.lock().await { + handle.abort(); + } + + self.pair.strategy.lock().await.process.kill().await?; + self.pair.risk.lock().await.process.kill().await?; } - if let Some(handle) = &*self.risk_handle.lock().await { - handle.abort(); + { + let strategy = StrategyPair::new(strategy_id, risk_id).await?; + + std::mem::swap( + &mut *self.pair.strategy.lock().await, + &mut *strategy.strategy.lock().await, + ); + + std::mem::swap( + &mut *self.pair.strategy_manifest.lock().await, + &mut *strategy.strategy_manifest.lock().await, + ); + + std::mem::swap( + &mut *self.pair.risk.lock().await, + &mut *strategy.risk.lock().await, + ); + + std::mem::swap( + &mut *self.pair.risk_manifest.lock().await, + &mut *strategy.risk_manifest.lock().await, + ); } - let mut pair = self.pair.lock().await; - - pair.strategy.process.kill().await?; - pair.risk.process.kill().await?; - - *pair = StrategyPair::new(strategy_id, risk_id).await?; + self.spawn().await; Ok(()) }