Signal forwarding
This commit is contained in:
@@ -62,11 +62,11 @@ impl<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
|
|||||||
Ok(Some(R::from_com(&mut buffer)))
|
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
|
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 mut process = self.process.lock().await;
|
||||||
let stdin = process.stdin.as_mut().unwrap();
|
let stdin = process.stdin.as_mut().unwrap();
|
||||||
|
|
||||||
@@ -170,11 +170,14 @@ impl StrategyEngine {
|
|||||||
match pair.strategy.recv().await? {
|
match pair.strategy.recv().await? {
|
||||||
None => {}
|
None => {}
|
||||||
|
|
||||||
Some(StrategyMessage::Log(log)) => {
|
Some(StrategyMessage::Log(mut log)) => {
|
||||||
|
log.name.insert_str(0, "strategy::");
|
||||||
engine.terminal_server.log_raw(log).await?;
|
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 }) => {}
|
Some(StrategyMessage::SubscribeCandle { symbol, timeframe }) => {}
|
||||||
|
|
||||||
@@ -193,7 +196,21 @@ impl StrategyEngine {
|
|||||||
.upgrade()
|
.upgrade()
|
||||||
.expect("Failed to upgrade engine (StrategyEngine)");
|
.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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user