Plugin child process
This commit is contained in:
@@ -1 +1,59 @@
|
||||
pub struct Plugin();
|
||||
use std::marker::PhantomData;
|
||||
|
||||
use pulse_wire::{
|
||||
PulseWire,
|
||||
plugin::{RiskEngineMessage, RiskMessage, StrategyEngineMessage, StrategyMessage},
|
||||
};
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
|
||||
pub struct Plugin<S: PulseWire, R: PulseWire> {
|
||||
pub process: tokio::process::Child,
|
||||
pub _p: (PhantomData<S>, PhantomData<R>),
|
||||
}
|
||||
|
||||
impl<S: PulseWire, R: PulseWire> Plugin<S, R> {
|
||||
pub fn new(child: tokio::process::Child) -> Self {
|
||||
Self {
|
||||
process: child,
|
||||
_p: (PhantomData, PhantomData),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn recv(&mut self) -> tokio::io::Result<Option<R>> {
|
||||
let stderr = self.process.stderr.as_mut().unwrap();
|
||||
|
||||
let mut len_buf = [0u8; size_of::<usize>()];
|
||||
let size = stderr.read_exact(&mut len_buf).await?;
|
||||
|
||||
let len = usize::from_le_bytes(len_buf);
|
||||
|
||||
if size == 0 || len == 0 {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let mut buffer = vec![0u8; len];
|
||||
|
||||
stderr.read_exact(&mut buffer).await?;
|
||||
|
||||
Ok(Some(R::from_com(&mut buffer)))
|
||||
}
|
||||
|
||||
pub async fn send(&mut 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<()> {
|
||||
let stdin = self.process.stdin.as_mut().unwrap();
|
||||
|
||||
stdin.write(&msg.len().to_le_bytes()).await?;
|
||||
stdin.write(msg).await?;
|
||||
stdin.flush().await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
pub struct StrategyPair {
|
||||
pub strategy: Plugin<StrategyEngineMessage, StrategyMessage>,
|
||||
pub risk: Plugin<RiskEngineMessage, RiskMessage>,
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user