diff --git a/pulse-sdk/src/general.rs b/pulse-sdk/src/general.rs index 5b5a21f..03582af 100644 --- a/pulse-sdk/src/general.rs +++ b/pulse-sdk/src/general.rs @@ -16,6 +16,20 @@ pub enum LogKind { Debug, } +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub enum ItemState { + Starting, + Running, + Stopped, + Error, +} + +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub enum Mode { + Auto, + Manual, +} + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct Signal { pub symbol: String, @@ -41,6 +55,12 @@ pub struct Position { pub pnl: Decimal, } +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub struct EngineStatus { + pub strategy_mode: Mode, + pub strategy_state: ItemState, +} + impl std::fmt::Display for MarketTrend { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { @@ -61,3 +81,23 @@ impl std::fmt::Display for LogKind { } } } + +impl std::fmt::Display for Mode { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Auto => write!(f, "\x1b[96mAUTO\x1b[0m"), + Self::Manual => write!(f, "\x1b[97mMANUAL\x1b[0m"), + } + } +} + +impl std::fmt::Display for ItemState { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Starting => write!(f, "\x1b[93mSTARTING\x1b[0m"), + Self::Running => write!(f, "\x1b[92mRUNNING\x1b[0m"), + Self::Stopped => write!(f, "\x1b[90mSTOPPED\x1b[0m"), + Self::Error => write!(f, "\x1b[91mERROR\x1b[0m"), + } + } +} diff --git a/pulse-sdk/src/terminal.rs b/pulse-sdk/src/terminal.rs index 6d52182..d111e0b 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,6 +1,5 @@ use crate::{ - general::{EventLog, Position, Signal}, - strategy::StrategyManifest, + general::{EventLog, ItemState, Mode, Position, Signal}, strategy::StrategyManifest, }; use hypersdk::{Decimal, hypercore::CandleInterval}; @@ -66,19 +65,6 @@ pub struct Status { pub latency: u16, } -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub enum Mode { - Auto, - Manual, -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub enum ItemState { - Running, - Stopped, - Error, -} - #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct StrategyStatus { pub strategy: StrategyManifest, @@ -87,22 +73,3 @@ pub struct StrategyStatus { pub state: ItemState, pub cooldown: CandleInterval, } - -impl std::fmt::Display for Mode { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - match self { - Self::Auto => write!(f, "\x1b[96mAUTO\x1b[0m"), - Self::Manual => write!(f, "\x1b[97mMANUAL\x1b[0m"), - } - } -} - -impl std::fmt::Display for ItemState { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - match self { - Self::Running => write!(f, "\x1b[92mRUNNING\x1b[0m"), - Self::Stopped => write!(f, "\x1b[90mSTOPPED\x1b[0m"), - Self::Error => write!(f, "\x1b[91mERROR\x1b[0m"), - } - } -} diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index c3c7a75..dfe98cd 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -179,12 +179,18 @@ impl Engine { return self.invalid_command_usage("strategy").await; }; - self.strategy_engine - .send(&StrategyEngineMessage::Command { - command: strategy_command.to_owned(), - args: args.into_iter().map(Into::into).collect(), - }) - .await?; + match strategy_command { + "start" => {} + + _ => { + self.strategy_engine + .send(&StrategyEngineMessage::Command { + command: strategy_command.to_owned(), + args: args.into_iter().map(Into::into).collect(), + }) + .await?; + } + } } _ => { diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 045b85d..7487d17 100644 --- a/src/engine/engine/mod.rs +++ b/src/engine/engine/mod.rs @@ -28,6 +28,9 @@ pub struct Engine { // data pub config: Arc>, pub accounts: Arc>, + + // live data / status + pub status: Arc>, pub watch_list: Arc>, pub signals: Arc>>, } @@ -44,11 +47,20 @@ impl Engine { let config = Arc::new(Mutex::new(config)); Ok(Arc::new_cyclic(|engine| Self { - config, - accounts, - ws_stream: Arc::new(Mutex::new(ws_stream)), + // engine terminal_server: TerminalServer::new(engine.clone()), strategy_engine: strategy.initialize(engine.clone()), + ws_stream: Arc::new(Mutex::new(ws_stream)), + + // data + config, + accounts, + + // live data / status + status: Arc::new(Mutex::new(EngineStatus { + strategy_mode: Mode::Auto, + strategy_state: ItemState::Stopped, + })), watch_list: Arc::new(Mutex::new(WatchList { name_to_index: HashMap::new(), items: Vec::new(), @@ -61,10 +73,16 @@ impl Engine { /// When a strategy is reloaded it restarts automatically pub async fn run(&self) -> anyhow::Result<()> { loop { + self.terminal_server + .info("engine::main", "Strategy starting") + .await?; + + self.status.lock().await.strategy_state = ItemState::Starting; + self.strategy_engine.run().await?; self.terminal_server - .info("engine::main", "Strategy stopped, restarting") + .info("engine::main", "Strategy stopped") .await?; } } diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index b5fb0cf..b91690a 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -93,6 +93,13 @@ impl StrategyEngine { self.send(&StrategyEngineMessage::Initialize).await?; + engine + .terminal_server + .info("engine::main", "Strategy Started") + .await?; + + engine.status.lock().await.strategy_state = ItemState::Running; + loop { match StrategyChild::read(&mut stdout).await? { None => {}