Engine status
This commit is contained in:
@@ -16,6 +16,20 @@ pub enum LogKind {
|
|||||||
Debug,
|
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)]
|
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
|
||||||
pub struct Signal {
|
pub struct Signal {
|
||||||
pub symbol: String,
|
pub symbol: String,
|
||||||
@@ -41,6 +55,12 @@ pub struct Position {
|
|||||||
pub pnl: Decimal,
|
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 {
|
impl std::fmt::Display for MarketTrend {
|
||||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||||
match self {
|
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"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,6 +1,5 @@
|
|||||||
use crate::{
|
use crate::{
|
||||||
general::{EventLog, Position, Signal},
|
general::{EventLog, ItemState, Mode, Position, Signal}, strategy::StrategyManifest,
|
||||||
strategy::StrategyManifest,
|
|
||||||
};
|
};
|
||||||
use hypersdk::{Decimal, hypercore::CandleInterval};
|
use hypersdk::{Decimal, hypercore::CandleInterval};
|
||||||
|
|
||||||
@@ -66,19 +65,6 @@ pub struct Status {
|
|||||||
pub latency: u16,
|
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)]
|
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
|
||||||
pub struct StrategyStatus {
|
pub struct StrategyStatus {
|
||||||
pub strategy: StrategyManifest,
|
pub strategy: StrategyManifest,
|
||||||
@@ -87,22 +73,3 @@ pub struct StrategyStatus {
|
|||||||
pub state: ItemState,
|
pub state: ItemState,
|
||||||
pub cooldown: CandleInterval,
|
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"),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -179,6 +179,10 @@ impl Engine {
|
|||||||
return self.invalid_command_usage("strategy").await;
|
return self.invalid_command_usage("strategy").await;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
match strategy_command {
|
||||||
|
"start" => {}
|
||||||
|
|
||||||
|
_ => {
|
||||||
self.strategy_engine
|
self.strategy_engine
|
||||||
.send(&StrategyEngineMessage::Command {
|
.send(&StrategyEngineMessage::Command {
|
||||||
command: strategy_command.to_owned(),
|
command: strategy_command.to_owned(),
|
||||||
@@ -186,6 +190,8 @@ impl Engine {
|
|||||||
})
|
})
|
||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
_ => {
|
_ => {
|
||||||
self.terminal_server
|
self.terminal_server
|
||||||
|
|||||||
@@ -28,6 +28,9 @@ pub struct Engine {
|
|||||||
// data
|
// data
|
||||||
pub config: Arc<Mutex<Config>>,
|
pub config: Arc<Mutex<Config>>,
|
||||||
pub accounts: Arc<Mutex<AccountList>>,
|
pub accounts: Arc<Mutex<AccountList>>,
|
||||||
|
|
||||||
|
// live data / status
|
||||||
|
pub status: Arc<Mutex<EngineStatus>>,
|
||||||
pub watch_list: Arc<Mutex<WatchList>>,
|
pub watch_list: Arc<Mutex<WatchList>>,
|
||||||
pub signals: Arc<Mutex<Vec<SignalStatus>>>,
|
pub signals: Arc<Mutex<Vec<SignalStatus>>>,
|
||||||
}
|
}
|
||||||
@@ -44,11 +47,20 @@ impl Engine {
|
|||||||
let config = Arc::new(Mutex::new(config));
|
let config = Arc::new(Mutex::new(config));
|
||||||
|
|
||||||
Ok(Arc::new_cyclic(|engine| Self {
|
Ok(Arc::new_cyclic(|engine| Self {
|
||||||
config,
|
// engine
|
||||||
accounts,
|
|
||||||
ws_stream: Arc::new(Mutex::new(ws_stream)),
|
|
||||||
terminal_server: TerminalServer::new(engine.clone()),
|
terminal_server: TerminalServer::new(engine.clone()),
|
||||||
strategy_engine: strategy.initialize(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 {
|
watch_list: Arc::new(Mutex::new(WatchList {
|
||||||
name_to_index: HashMap::new(),
|
name_to_index: HashMap::new(),
|
||||||
items: Vec::new(),
|
items: Vec::new(),
|
||||||
@@ -61,10 +73,16 @@ impl Engine {
|
|||||||
/// When a strategy is reloaded it restarts automatically
|
/// When a strategy is reloaded it restarts automatically
|
||||||
pub async fn run(&self) -> anyhow::Result<()> {
|
pub async fn run(&self) -> anyhow::Result<()> {
|
||||||
loop {
|
loop {
|
||||||
|
self.terminal_server
|
||||||
|
.info("engine::main", "Strategy starting")
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
self.status.lock().await.strategy_state = ItemState::Starting;
|
||||||
|
|
||||||
self.strategy_engine.run().await?;
|
self.strategy_engine.run().await?;
|
||||||
|
|
||||||
self.terminal_server
|
self.terminal_server
|
||||||
.info("engine::main", "Strategy stopped, restarting")
|
.info("engine::main", "Strategy stopped")
|
||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -93,6 +93,13 @@ impl StrategyEngine {
|
|||||||
|
|
||||||
self.send(&StrategyEngineMessage::Initialize).await?;
|
self.send(&StrategyEngineMessage::Initialize).await?;
|
||||||
|
|
||||||
|
engine
|
||||||
|
.terminal_server
|
||||||
|
.info("engine::main", "Strategy Started")
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
engine.status.lock().await.strategy_state = ItemState::Running;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
match StrategyChild::read(&mut stdout).await? {
|
match StrategyChild::read(&mut stdout).await? {
|
||||||
None => {}
|
None => {}
|
||||||
|
|||||||
Reference in New Issue
Block a user