From 87e9f62c8db82fd06db42a52004751ab3959a565 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Sat, 25 Jul 2026 20:09:18 +0200 Subject: [PATCH] Plugin child process --- Cargo.lock | 1 + Cargo.toml | 9 +----- src/engine/store/plugin.rs | 60 +++++++++++++++++++++++++++++++++++++- 3 files changed, 61 insertions(+), 9 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8495645..3cf15f6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5200,6 +5200,7 @@ dependencies = [ "libc", "mio", "pin-project-lite", + "signal-hook-registry", "socket2", "tokio-macros", "windows-sys 0.61.2", diff --git a/Cargo.toml b/Cargo.toml index 8780172..8b87ecb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,14 +8,7 @@ pulse-ui = { workspace = true } pulse-wire = { workspace = true } chrono = "0.4.45" -tokio = { workspace = true, features = [ - "rt-multi-thread", - "macros", - "net", - "fs", - "io-util", - "time", -] } +tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time", "process"] } crossterm = { workspace = true } hypersdk = "0.2.14" serde_json = "1" diff --git a/src/engine/store/plugin.rs b/src/engine/store/plugin.rs index 5fd4ca3..57d03de 100644 --- a/src/engine/store/plugin.rs +++ b/src/engine/store/plugin.rs @@ -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 { + pub process: tokio::process::Child, + pub _p: (PhantomData, PhantomData), +} + +impl Plugin { + pub fn new(child: tokio::process::Child) -> Self { + Self { + process: child, + _p: (PhantomData, PhantomData), + } + } + + pub async fn recv(&mut self) -> tokio::io::Result> { + let stderr = self.process.stderr.as_mut().unwrap(); + + let mut len_buf = [0u8; size_of::()]; + 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, + pub risk: Plugin, +}