69 lines
1.9 KiB
Rust
69 lines
1.9 KiB
Rust
use pulse_sdk::{map_postcard_err, prelude::*};
|
|
|
|
use std::{path::PathBuf, process::Stdio};
|
|
use tokio::{
|
|
io::{AsyncReadExt, AsyncWriteExt},
|
|
process::{Child, ChildStdout, Command},
|
|
};
|
|
#[derive(Debug)]
|
|
pub struct StrategyChild {
|
|
pub manifest: StrategyManifest,
|
|
pub child: Child,
|
|
}
|
|
|
|
impl StrategyChild {
|
|
pub fn new(child: Child, manifest: StrategyManifest) -> Self {
|
|
Self { child, manifest }
|
|
}
|
|
|
|
pub async fn read(stdout: &mut ChildStdout) -> tokio::io::Result<Option<StrategyMessage>> {
|
|
let mut len_buf = [0u8; size_of::<usize>()];
|
|
let size = stdout.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];
|
|
|
|
stdout.read_exact(&mut buffer).await?;
|
|
|
|
Ok(Some(map_postcard_err(postcard::from_bytes(&buffer))?))
|
|
}
|
|
|
|
pub async fn send(&mut self, msg: &StrategyEngineMessage) -> tokio::io::Result<()> {
|
|
self.send_raw(&map_postcard_err(postcard::to_allocvec(msg))?)
|
|
.await
|
|
}
|
|
|
|
pub async fn send_raw(&mut self, msg: &[u8]) -> tokio::io::Result<()> {
|
|
let stdin = self.child.stdin.as_mut().unwrap();
|
|
|
|
stdin.write_all(&msg.len().to_le_bytes()).await?;
|
|
stdin.write_all(msg).await?;
|
|
stdin.flush().await?;
|
|
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
pub fn get_manifest<'de, M: serde::Deserialize<'de>>(
|
|
strategy_dir: &PathBuf,
|
|
strategy_path: &PathBuf,
|
|
manifest: &'de [u8],
|
|
) -> tokio::io::Result<(Child, M)> {
|
|
Ok((
|
|
Command::new("bash")
|
|
.arg(strategy_path)
|
|
.current_dir(strategy_dir)
|
|
.stdin(Stdio::piped())
|
|
.stdout(Stdio::piped())
|
|
.stderr(Stdio::inherit())
|
|
.spawn()?,
|
|
toml::from_slice(manifest)
|
|
.map_err(|v| tokio::io::Error::new(std::io::ErrorKind::InvalidInput, v.to_string()))?,
|
|
))
|
|
}
|