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..c82562b 100644 --- a/pulse-sdk/src/terminal.rs +++ b/pulse-sdk/src/terminal.rs @@ -1,8 +1,8 @@ use crate::{ - general::{EventLog, Position, Signal}, + general::{EngineStatus, EventLog, Position, Signal}, strategy::StrategyManifest, }; -use hypersdk::{Decimal, hypercore::CandleInterval}; +use hypersdk::Decimal; pub type SignalStatus = Result; @@ -14,8 +14,8 @@ pub enum TerminalServerMessage { // Positions PositionsUpdated(Vec), - // Strategy - StrategyUpdated(StrategyStatus), + // Configuration + ConfigUpdated(EngineConfig), // Signals SignalsUpdated(Vec), @@ -24,7 +24,7 @@ pub enum TerminalServerMessage { Inspect(InspectTarget), // Status - StatusUpdated(Status), + StatusUpdated(EngineStatus), // Logs SetLogs(Vec), @@ -47,62 +47,21 @@ pub struct MarketItem { #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum InspectTarget { None, - Some(Vec), + Some(Vec), } +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub struct InspectRow(pub InspectItem, pub InspectItem, pub InspectItem); + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum InspectItem { - String(String), + Name(String), Symbol(String), - USD(f64), - F64(f64), + USD(Decimal), + F64(Decimal), } #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub struct Status { - pub feed: Mode, - pub exchange: String, - pub dex: String, - 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 struct EngineConfig { pub strategy: StrategyManifest, - - pub mode: Mode, - 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/pulse-ui/src/render.rs b/pulse-ui/src/render.rs index dea001d..cb58646 100644 --- a/pulse-ui/src/render.rs +++ b/pulse-ui/src/render.rs @@ -15,58 +15,79 @@ pub struct RenderScope { draw_instructions: Vec, } +struct RenderTextResult { + text: String, + line_len: u16, + lines_len: u16, +} + +impl RenderTextResult { + fn insert_newline(&mut self, height: u16) -> bool { + if self.lines_len >= height { + return true; + } + + self.text.push('\n'); + self.line_len = 0; + self.lines_len += 1; + + false + } + + fn insert(&mut self, c: char) { + self.text.push(c); + self.line_len += 1; + } +} + impl RenderScope { pub fn draw_text, T: Display>(&mut self, at: P, text: T) -> u16 { + if self.rect.width == 0 || self.rect.height == 0 { + return 0; + } + let point: Point = at.into(); - let lines = text - .to_string() - .lines() - .map(|line| { - let mut chars = line.chars().peekable(); - let mut new_line = String::new(); - let mut len = 0; + let text = text.to_string(); + + let mut chars = text.chars().peekable(); + + let mut result = RenderTextResult { + text: String::new(), + line_len: 0, + lines_len: 1, + }; + + while let Some(c) = chars.next() { + if c == '\x1b' && chars.peek() == Some(&'[') { + result.text.push(c); while let Some(c) = chars.next() { - if c == '\x1b' && chars.peek() == Some(&'[') { - new_line.push(c); + result.text.push(c); - while let Some(c) = chars.next() { - new_line.push(c); - - if c.is_ascii_alphabetic() { - break; - } - } - - continue; - } - - new_line.push(c); - len += 1; - - if len >= self.rect.width as usize { - new_line.push('\n'); - len = 0; + if c.is_ascii_alphabetic() { + break; } } + } else if c == '\n' { + if result.insert_newline(self.rect.height) { + break; + } + } else { + result.insert(c); - new_line - }) - .collect::>() - .join("\n"); - - let lines = lines - .lines() - .take((self.rect.height - point.y) as usize) - .collect::>(); - - let lines_len = lines.len(); + if result.line_len >= self.rect.width { + if result.insert_newline(self.rect.height) { + break; + } + } + } + } self.draw_instructions - .push(Instr::DrawText(point, lines.join("\n"))); + .push(Instr::DrawText(point, result.text)); - lines_len as u16 + result.lines_len } } @@ -123,6 +144,10 @@ impl RenderScope { match inst { Instr::DrawText(point, text) => { for (i, line) in text.lines().enumerate() { + if i as u16 + point.y >= self.rect.height { + continue; + } + print_at( self.rect.x + point.x, self.rect.y + point.y + i as u16, diff --git a/pulse-ui/src/widget/scroll.rs b/pulse-ui/src/widget/scroll.rs index 2a12586..62b134f 100644 --- a/pulse-ui/src/widget/scroll.rs +++ b/pulse-ui/src/widget/scroll.rs @@ -10,20 +10,17 @@ pub struct ScrollText { impl Widget for ScrollText { fn render(&self, scope: &mut crate::render::RenderScope) { - scope.draw_text(0, &self.title); + let title_lines = scope.draw_text(0, &self.title); - let title_lines = self.title.lines().count(); - - let mut y = title_lines as u16; - - for line in self + let text = self .text .lines() .skip(self.scroll) - .take(scope.rect.height as usize - title_lines) - { - y += scope.draw_text((0, y), line); - } + .map(|v| v.trim_end()) + .collect::>() + .join("\n"); + + scope.draw_text((0, title_lines), text); } } diff --git a/src/engine/engine/command.rs b/src/engine/engine/command.rs index c3c7a75..0a815e1 100644 --- a/src/engine/engine/command.rs +++ b/src/engine/engine/command.rs @@ -22,10 +22,16 @@ impl Engine { match args[0] { "reload" => { *self.config.lock().await = Config::new().await?; + self.terminal_server + .info("engine::config", "successfully reloaded") + .await?; } "save" => { self.config.lock().await.save().await?; + self.terminal_server + .info("engine::config", "successfully saved") + .await?; } "set" => { @@ -49,6 +55,12 @@ impl Engine { self.terminal_server .info("config::set", "watchlist set successfully, use `config save` to persist changes") .await?; + + self.terminal_server + .broadcast(TerminalServerMessage::ConfigUpdated( + self.get_config_status().await, + )) + .await?; } "strategy" | "strat" | "sg" => { @@ -62,7 +74,7 @@ impl Engine { return self .terminal_server .error( - "config::set::strategy", + "config::set[strategy]", &format!("Non existent strategy `{id}`"), ) .await; @@ -75,23 +87,24 @@ impl Engine { self.terminal_server .info("config::set", "strategy set successfully, use `config save` to persist changes") .await?; + + self.terminal_server + .broadcast(TerminalServerMessage::ConfigUpdated( + self.get_config_status().await, + )) + .await?; } "cooldown" | "cool" | "cd" => { set_cfg!(cooldown, { self.config.lock().await.cooldown = cooldown; - - self.terminal_server.broadcast( - pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(StrategyStatus { - strategy: self.strategy_engine.strategy.lock().await.manifest.clone(), - - mode: Mode::Auto, - state: ItemState::Running, - cooldown, - }), - ) - .await?; }); + + self.terminal_server + .broadcast(TerminalServerMessage::ConfigUpdated( + self.get_config_status().await, + )) + .await?; } _ => { @@ -101,25 +114,25 @@ impl Engine { } } } else { - self.invalid_command_usage("config").await?; + self.invalid_command_usage("engine::config").await?; } } _ => { - self.invalid_command_usage("config").await?; + self.invalid_command_usage("engine::config").await?; } } } "account" | "acc" => { if args.len() == 0 { - return self.invalid_command_usage("account man").await; + return self.invalid_command_usage("engine::account").await; } match args[0] { "list" | "ls" => { self.terminal_server - .info("account man", "ACCOUNT LIST") + .info("engine::account", "ACCOUNT LIST") .await?; let accounts = self.accounts.lock().await; @@ -127,7 +140,7 @@ impl Engine { for (name, acc) in &accounts.accounts { self.terminal_server .info( - "account man", + "engine::account", &if name == &accounts.active { format!( "{} (active) -> {}", @@ -144,7 +157,7 @@ impl Engine { "use" | "set" => { if args.len() < 2 { - return self.invalid_command_usage("account man").await; + return self.invalid_command_usage("engine::account").await; } let new_active = args[1]; @@ -154,7 +167,10 @@ impl Engine { if !accounts.accounts.contains_key(new_active) { return self .terminal_server - .error("account man", &format!("Account not found ({new_active})")) + .error( + "engine::account", + &format!("Account not found ({new_active})"), + ) .await; } @@ -162,14 +178,14 @@ impl Engine { self.terminal_server .info( - "account man", + "engine::account", &format!("Account set to {new_active} successfully!"), ) .await?; } _ => { - return self.invalid_command_usage("account man").await; + return self.invalid_command_usage("engine::account").await; } } } @@ -179,20 +195,23 @@ 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?; + } + } } _ => { self.terminal_server - .error( - "Command executor", - &format!("Command '{}' not found", command), - ) + .error("engine::cmd", &format!("Command '{}' not found", command)) .await?; } } diff --git a/src/engine/engine/mod.rs b/src/engine/engine/mod.rs index 045b85d..492dc7f 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,11 +73,33 @@ 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.update_status().await?; + self.strategy_engine.run().await?; self.terminal_server - .info("engine::main", "Strategy stopped, restarting") + .info("engine::main", "Strategy stopped") .await?; } } + + pub async fn get_config_status(&self) -> EngineConfig { + EngineConfig { + strategy: self.strategy_engine.strategy.lock().await.manifest.clone(), + } + } + + pub async fn update_status(&self) -> tokio::io::Result<()> { + self.terminal_server + .broadcast(TerminalServerMessage::StatusUpdated( + self.status.lock().await.clone(), + )) + .await + } } diff --git a/src/engine/engine/strategy.rs b/src/engine/engine/strategy.rs index b5fb0cf..b79e26e 100644 --- a/src/engine/engine/strategy.rs +++ b/src/engine/engine/strategy.rs @@ -93,6 +93,14 @@ impl StrategyEngine { self.send(&StrategyEngineMessage::Initialize).await?; + engine + .terminal_server + .info("engine::main", "Strategy Started") + .await?; + + engine.status.lock().await.strategy_state = ItemState::Running; + engine.update_status().await?; + loop { match StrategyChild::read(&mut stdout).await? { None => {} @@ -158,8 +166,6 @@ impl StrategyEngine { .unwrap() .as_millis() as u64; - println!("getting candlestick"); - let interval_ms = match interval { CandleInterval::OneMinute => 60_000, CandleInterval::ThreeMinutes => 3 * 60_000, diff --git a/src/engine/engine/terminal.rs b/src/engine/engine/terminal.rs index e308812..f2e6211 100644 --- a/src/engine/engine/terminal.rs +++ b/src/engine/engine/terminal.rs @@ -45,7 +45,11 @@ impl TerminalServer { let listener = UnixListener::bind(&path)?; - println!("Terminal server listening on {:?}", path); + self.info( + "engine::terminal", + &format!("Terminal server listening on {:?}", path), + ) + .await?; loop { let (stream, _) = listener.accept().await?; @@ -60,7 +64,23 @@ impl TerminalServer { tokio::spawn(async move { if let Err(err) = s.handle_client(&id, reader).await { - eprintln!("Terminal connection error: {err}"); + if matches!(err.kind(), tokio::io::ErrorKind::UnexpectedEof) { + Self::print_log(&EventLog { + kind: LogKind::Err, + name: "engine::terminal".to_string(), + message: format!("Terminal connection error: {err}"), + }); + + return Ok(()); + } + + s.error( + "engine::terminal", + &format!("Terminal connection error: {err}"), + ) + .await + } else { + Ok(()) } }); } @@ -71,18 +91,13 @@ impl TerminalServer { self.send_to( id, - pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(StrategyStatus { - strategy: engine - .strategy_engine - .strategy - .lock() - .await - .manifest - .clone(), - mode: Mode::Auto, - state: ItemState::Running, - cooldown: engine.config.lock().await.cooldown, - }), + TerminalServerMessage::ConfigUpdated(engine.get_config_status().await), + ) + .await?; + + self.send_to( + id, + TerminalServerMessage::StatusUpdated(engine.status.lock().await.clone()), ) .await?; @@ -147,7 +162,12 @@ impl TerminalServer { for (id, client) in clients.iter_mut() { if let Err(e) = Self::send_to_client(client, &msg).await { remove_clients.push(*id); - println!("{e:?}"); + + Self::print_log(&EventLog { + kind: LogKind::Err, + name: "engine::terminal".to_string(), + message: format!("Failed to broadcast to client: {e}"), + }); } } @@ -264,19 +284,32 @@ impl TerminalServer { name: &str, message: &str, ) -> tokio::io::Result<()> { - let log = EventLog { + self.log_raw(EventLog { kind, name: name.to_string(), message: message.to_string(), - }; + }) + .await + } - self.logs.lock().await.push(log.clone()); - - self.broadcast(pulse_sdk::terminal::TerminalServerMessage::AddLog(log)) - .await + pub fn print_log(log: &EventLog) { + println!( + "[{}{}\x1b[0m] (\x1b[37m{}\x1b[0m) {}", + match &log.kind { + LogKind::Info => "\x1b[34m", + LogKind::Debug => "\x1b[36m", + LogKind::Err => "\x1b[31m", + LogKind::Warn => "\x1b[33m", + }, + log.kind, + log.name, + log.message + ); } pub async fn log_raw(self: &Arc, log: EventLog) -> tokio::io::Result<()> { + Self::print_log(&log); + self.logs.lock().await.push(log.clone()); self.broadcast(pulse_sdk::terminal::TerminalServerMessage::AddLog(log)) diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 1fd6fdb..64744a6 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -10,15 +10,24 @@ impl Formatted for InspectTarget { match self { Self::None => vec!["\x1b[2mnothing to inspect\x1b[0m".to_string()], - Self::Some(items) => items - .iter() - .map(|item| match item { - InspectItem::F64(v) => format_f64(*v), - InspectItem::USD(v) => format!("${}", format_f64(*v)), - InspectItem::Symbol(v) => format!("\x1b[35m{v}\x1b[0m"), - InspectItem::String(v) => v.clone(), - }) - .collect(), + Self::Some(items) => { + let items = items + .iter() + .map(|row| { + ( + format_inspect_item(&row.0), + format_inspect_item(&row.1), + format_inspect_item(&row.2), + ) + }) + .collect::>(); + + items + .iter() + .map(|(a, b, c)| Triple(a, b, c)) + .collect::>() + .get_formatted() + } } } } @@ -114,28 +123,8 @@ impl Formatted for Position { } } -struct Pair<'a>(&'a str, &'a str); - struct Triple<'a>(&'a str, &'a str, &'a str); -impl Formatted for Status { - fn get_formatted(&self) -> Vec { - vec![ - Pair("Feed", &format!("{:?}", self.feed)), - Pair("Exchange", &self.exchange.clone()), - Pair("DEX", &self.dex.clone()), - Pair("Latency", &format!("{} ms", self.latency)), - ] - .get_formatted() - } -} - -impl<'a> Formatted for Pair<'a> { - fn get_formatted(&self) -> Vec { - vec![self.0.to_string(), self.1.to_string()] - } -} - impl<'a> Formatted for Triple<'a> { fn get_formatted(&self) -> Vec { vec![self.0.to_string(), self.1.to_string(), self.2.to_string()] @@ -206,19 +195,33 @@ impl Formatted for Vec { } } -impl Formatted for StrategyStatus { +impl Formatted for EngineStatus { + fn get_formatted(&self) -> Vec { + vec![ + Triple( + "\x1b[2mState\x1b[0m", + "\x1b[2mMode\x1b[0m", + "\x1b[2m..\x1b[0m", + ), + Triple( + &self.strategy_state.to_string(), + &self.strategy_mode.to_string(), + "..", + ), + ] + .get_formatted() + } +} + +impl Formatted for EngineConfig { fn get_formatted(&self) -> Vec { vec![ Triple( "\x1b[2mStrategy\x1b[0m", - "\x1b[2mState\x1b[0m", - "\x1b[2mMode\x1b[0m", - ), - Triple( - &format!("\x1b[97m{}\x1b[0m", self.strategy.name), - &self.state.to_string(), - &self.mode.to_string(), + "\x1b[2m..\x1b[0m", + "\x1b[2m..\x1b[0m", ), + Triple(&self.strategy.name, "..", ".."), ] .get_formatted() } @@ -232,6 +235,15 @@ pub fn apply_padding(mut items: Vec) -> Vec { items } +pub fn format_inspect_item(item: &InspectItem) -> String { + match item { + InspectItem::F64(v) => format_f64(v.as_f64()), + InspectItem::USD(v) => format!("${}", format_f64(v.as_f64())), + InspectItem::Symbol(v) => format!("\x1b[35m{v}\x1b[0m"), + InspectItem::Name(v) => v.clone(), + } +} + pub fn format_symbol(value: &str) -> String { format!("\x1b[35m{value}\x1b[0m") } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 070d064..dc8cb13 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -29,8 +29,8 @@ pub struct PulseTradeApp { command: State, scroll: State>, - strategy: State>, - status: State>, + config: State>, + status: State>, watch_list: State>, active_positions: State>, @@ -144,8 +144,7 @@ impl App for PulseTradeApp { ( LayoutItem::Widget(Size::Flex(1)), Box::new( - advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.strategy) - .await, + advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.config).await, ), ), ( @@ -182,7 +181,7 @@ async fn main() -> tokio::io::Result<()> { signals: ctx.use_state(Vec::new()), logs: ctx.use_state(Vec::new()), inspect: ctx.use_state(InspectTarget::None), - strategy: ctx.use_state(None), + config: ctx.use_state(None), status: ctx.use_state(None), }) }) diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs index ccc646d..9acf585 100644 --- a/src/terminal/terminal.rs +++ b/src/terminal/terminal.rs @@ -49,7 +49,7 @@ impl TerminalClient { let active_positions = app.active_positions.clone(); let logs = app.logs.clone(); let signals = app.signals.clone(); - let market_overview = app.strategy.clone(); + let market_overview = app.config.clone(); let status = app.status.clone(); let inspect = app.inspect.clone(); @@ -82,8 +82,8 @@ impl TerminalClient { active_positions: State>, logs: State>, signals: State>, - market_overview: State>, - status: State>, + config: State>, + status: State>, inspect: State, ) -> tokio::io::Result<()> { loop { @@ -111,8 +111,8 @@ impl TerminalClient { *active_positions.lock().await = v; } - TerminalServerMessage::StrategyUpdated(v) => { - *market_overview.lock().await = Some(v); + TerminalServerMessage::ConfigUpdated(v) => { + *config.lock().await = Some(v); } TerminalServerMessage::SignalsUpdated(v) => {