This commit is contained in:
2026-07-30 05:35:03 +02:00
parent 6c0bfe7c74
commit 535c6a2964
4 changed files with 42 additions and 49 deletions
+28 -32
View File
@@ -35,41 +35,37 @@ pub async fn send_raw(data: &[u8]) -> tokio::io::Result<()> {
Ok(()) Ok(())
} }
macro_rules! engine_methods { use std::sync::{
($t:ty) => { Arc,
async fn start(&self) -> tokio::io::Result<()> { atomic::{AtomicBool, Ordering},
let mut stdin = tokio::io::stdin(); };
loop {
let mut len_buf = [0u8; size_of::<usize>()];
let size = stdin.read_exact(&mut len_buf).await?;
let len = usize::from_le_bytes(len_buf);
if size == 0 || len == 0 {
break Ok(());
}
let mut buffer = vec![0u8; len];
stdin.read_exact(&mut buffer).await?;
self.on_raw(&buffer).await?;
}
}
async fn send(&self, msg: &$t) -> tokio::io::Result<()> {
$crate::send_raw(&$crate::map_postcard_err(
$crate::prelude::postcard::to_allocvec(msg),
)?)
.await
}
};
}
#[allow(async_fn_in_trait)] #[allow(async_fn_in_trait)]
pub trait Strategy { pub trait Strategy {
engine_methods!(prelude::StrategyMessage); async fn start(&self) -> tokio::io::Result<()> {
let mut stdin = tokio::io::stdin();
loop {
let mut len_buf = [0u8; size_of::<usize>()];
let size = stdin.read_exact(&mut len_buf).await?;
let len = usize::from_le_bytes(len_buf);
if size == 0 || len == 0 {
break Ok(());
}
let mut buffer = vec![0u8; len];
stdin.read_exact(&mut buffer).await?;
self.on_raw(&buffer).await?;
}
}
async fn send(&self, msg: &prelude::StrategyMessage) -> tokio::io::Result<()> {
send_raw(&map_postcard_err(prelude::postcard::to_allocvec(msg))?).await
}
async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> { async fn on_raw(&self, data: &[u8]) -> tokio::io::Result<()> {
match map_postcard_err(postcard::from_bytes(&data))? { match map_postcard_err(postcard::from_bytes(&data))? {
+2 -2
View File
@@ -60,7 +60,7 @@ impl Engine {
.await; .await;
} }
self.strategy.reload_strategy(id.as_str()).await?; self.strategy.reload(id.as_str()).await?;
self.config.lock().await.strategy = id; self.config.lock().await.strategy = id;
}); });
@@ -75,7 +75,7 @@ impl Engine {
self.terminal_server.broadcast( self.terminal_server.broadcast(
pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: self.strategy.strategy.manifest.lock().await.clone(), strategy: self.strategy.child.manifest.lock().await.clone(),
mode: Mode::Auto, mode: Mode::Auto,
state: ItemState::Running, state: ItemState::Running,
+11 -14
View File
@@ -19,9 +19,8 @@ use crate::{
}; };
pub struct StrategyEngine { pub struct StrategyEngine {
pub strategy: Arc<StrategyChild>,
pub engine: Weak<Engine>, pub engine: Weak<Engine>,
pub child: Arc<StrategyChild>,
pub ws: WebSocket, pub ws: WebSocket,
pub subscriptions: Mutex<HashSet<Subscription>>, pub subscriptions: Mutex<HashSet<Subscription>>,
} }
@@ -43,7 +42,7 @@ impl StrategyEngine {
)?; )?;
Ok(Self { Ok(Self {
strategy: Arc::new(StrategyChild::new(strategy, strategy_manifest)), child: Arc::new(StrategyChild::new(strategy, strategy_manifest)),
engine: Weak::new(), engine: Weak::new(),
ws: hypercore::mainnet_ws(), ws: hypercore::mainnet_ws(),
subscriptions: Mutex::new(HashSet::new()), subscriptions: Mutex::new(HashSet::new()),
@@ -56,22 +55,20 @@ impl StrategyEngine {
Arc::new(self) Arc::new(self)
} }
pub async fn run_strategy(&self) -> anyhow::Result<()> { pub async fn run(&self) -> anyhow::Result<()> {
let engine = self let engine = self
.engine .engine
.upgrade() .upgrade()
.expect("Failed to upgrade engine (StrategyEngine)"); .expect("Failed to upgrade engine (StrategyEngine)");
self.strategy self.child.send(&StrategyEngineMessage::Initialize).await?;
.send(&StrategyEngineMessage::Initialize)
.await?;
loop { loop {
match self.strategy.recv().await? { match self.child.recv().await? {
None => {} None => {}
Some(StrategyMessage::GetWatchList) => { Some(StrategyMessage::GetWatchList) => {
self.strategy self.child
.send(&StrategyEngineMessage::WatchList( .send(&StrategyEngineMessage::WatchList(
engine.watch_list.lock().await.clone().items, engine.watch_list.lock().await.clone().items,
)) ))
@@ -153,7 +150,7 @@ impl StrategyEngine {
let start_time = now.saturating_sub(interval_ms * count as u64); let start_time = now.saturating_sub(interval_ms * count as u64);
self.strategy self.child
.send(&StrategyEngineMessage::Candlestick { .send(&StrategyEngineMessage::Candlestick {
candles: client candles: client
.candle_snapshot(&symbol, interval, start_time, now) .candle_snapshot(&symbol, interval, start_time, now)
@@ -170,10 +167,10 @@ impl StrategyEngine {
pub async fn spawn(self: &Arc<Self>) { pub async fn spawn(self: &Arc<Self>) {
let engine = self.clone(); let engine = self.clone();
tokio::spawn(async move { engine.run_strategy().await }); tokio::spawn(async move { engine.run().await });
} }
pub async fn reload_strategy(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> { pub async fn reload(self: &Arc<Self>, id: &str) -> tokio::io::Result<()> {
let strategy = pulse_strategy(id)?; let strategy = pulse_strategy(id)?;
let (child, manifest) = get_manifest( let (child, manifest) = get_manifest(
@@ -182,10 +179,10 @@ impl StrategyEngine {
&fs::read(strategy.join("strategy.toml")).await?, &fs::read(strategy.join("strategy.toml")).await?,
)?; )?;
self.strategy.reload(child, manifest).await?; self.child.reload(child, manifest).await?;
let engine = self.clone(); let engine = self.clone();
tokio::spawn(async move { engine.run_strategy().await }); tokio::spawn(async move { engine.run().await });
Ok(()) Ok(())
} }
+1 -1
View File
@@ -65,7 +65,7 @@ impl TerminalServer {
self.send_to( self.send_to(
id, id,
pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy { pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: engine.strategy.strategy.manifest.lock().await.clone(), strategy: engine.strategy.child.manifest.lock().await.clone(),
mode: Mode::Auto, mode: Mode::Auto,
state: ItemState::Running, state: ItemState::Running,
cooldown: engine.config.lock().await.cooldown, cooldown: engine.config.lock().await.cooldown,