diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index a681b6c..3a83bd9 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -62,11 +62,11 @@ impl Deserialize<'de>> Plugin { Ok(Some(R::from_com(&mut buffer))) } - pub async fn send(&mut self, msg: &S) -> tokio::io::Result<()> { + pub async fn send(&self, msg: &S) -> tokio::io::Result<()> { self.send_raw(&msg.to_com()).await } - pub async fn send_raw(&mut self, msg: &[u8]) -> tokio::io::Result<()> { + pub async fn send_raw(&self, msg: &[u8]) -> tokio::io::Result<()> { let mut process = self.process.lock().await; let stdin = process.stdin.as_mut().unwrap(); @@ -170,11 +170,14 @@ impl StrategyEngine { match pair.strategy.recv().await? { None => {} - Some(StrategyMessage::Log(log)) => { + Some(StrategyMessage::Log(mut log)) => { + log.name.insert_str(0, "strategy::"); engine.terminal_server.log_raw(log).await?; } - Some(StrategyMessage::Signal(signal)) => {} + Some(StrategyMessage::Signal(signal)) => { + pair.risk.send(&RiskEngineMessage::Signal(signal)).await?; + } Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {} @@ -193,7 +196,21 @@ impl StrategyEngine { .upgrade() .expect("Failed to upgrade engine (StrategyEngine)"); - // loop {} + 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(()) }