diff --git a/Cargo.lock b/Cargo.lock index 1f87829..9caa2a4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -29,6 +29,12 @@ version = "3.20.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" +[[package]] +name = "bytes" +version = "1.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" + [[package]] name = "cc" version = "1.2.67" @@ -319,8 +325,8 @@ version = "0.1.0-alpha.0" dependencies = [ "chrono", "crossterm", - "pulse-macros", "pulse-ui", + "pulse-wire", "tokio", ] @@ -332,6 +338,13 @@ dependencies = [ "tokio", ] +[[package]] +name = "pulse-wire" +version = "0.1.0-alpha.0" +dependencies = [ + "pulse-macros", +] + [[package]] name = "quote" version = "1.0.46" @@ -439,6 +452,16 @@ version = "1.15.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8ed6a63f02c8539c91a8685a86f4099661ba3da017932f6ebbea6de3f0fa7c90" +[[package]] +name = "socket2" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51" +dependencies = [ + "libc", + "windows-sys", +] + [[package]] name = "syn" version = "2.0.118" @@ -456,8 +479,13 @@ version = "1.52.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8fc7f01b389ac15039e4dc9531aa973a135d7a4135281b12d7c1bc79fd57fffe" dependencies = [ + "bytes", + "libc", + "mio", "pin-project-lite", + "socket2", "tokio-macros", + "windows-sys", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 572e0f1..60ec14c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,18 +4,21 @@ version = "0.1.0-alpha.0" edition = "2024" [dependencies] -chrono = "0.4.45" pulse-ui = { workspace = true } -pulse-macros = { workspace = true } -tokio = { workspace = true, features = ["rt-multi-thread", "macros"] } +pulse-wire = { workspace = true } + +chrono = "0.4.45" +tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time"] } crossterm = { workspace = true } [workspace] -members = ["pulse-macros", "pulse-ui"] +members = ["pulse-macros", "pulse-ui", "pulse-wire"] [workspace.dependencies] pulse-macros = { path = "pulse-macros", version = "0.1.0-alpha.0" } pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" } +pulse-wire = { path = "pulse-wire", version = "0.1.0-alpha.0" } + tokio = "1.52.3" crossterm = "0.29.0" @@ -24,5 +27,5 @@ name = "pulse-trader" path = "src/terminal/main.rs" [[bin]] -name = "pulse-trader-daemon" -path = "src/daemon/main.rs" +name = "pulse-trader-engine" +path = "src/engine/main.rs" diff --git a/pulse-macros/src/lib.rs b/pulse-macros/src/lib.rs index d554b2d..cac53b8 100644 --- a/pulse-macros/src/lib.rs +++ b/pulse-macros/src/lib.rs @@ -3,7 +3,7 @@ use quote::quote; use syn::{Fields, ItemEnum, ItemStruct, parse_macro_input}; #[proc_macro_attribute] -pub fn p_com(_: TokenStream, item: TokenStream) -> TokenStream { +pub fn pwp(_: TokenStream, item: TokenStream) -> TokenStream { let input = parse_macro_input!(item as syn::Item); match input { @@ -49,7 +49,7 @@ fn expand_struct(mut input: ItemStruct) -> TokenStream { #[derive(Debug, Clone)] #input - impl PulseCom for #name { + impl PulseWire for #name { fn to_com(&self) -> Vec { let mut vec = Vec::new(); @@ -63,7 +63,7 @@ fn expand_struct(mut input: ItemStruct) -> TokenStream { fn from_com(com: &mut Vec) -> Self { Self { #( - #field_names2: PulseCom::from_com(com), + #field_names2: PulseWire::from_com(com), )* } } @@ -122,7 +122,7 @@ fn expand_enum(input: ItemEnum) -> TokenStream { Fields::Unnamed(fields) if fields.unnamed.len() == 1 => { quote! { - #tag => Self::#ident(PulseCom::from_com(com)), + #tag => Self::#ident(PulseWire::from_com(com)), } } @@ -132,7 +132,7 @@ fn expand_enum(input: ItemEnum) -> TokenStream { quote! { #tag => Self::#ident { #( - #names: PulseCom::from_com(com), + #names: PulseWire::from_com(com), )* }, } @@ -146,7 +146,7 @@ fn expand_enum(input: ItemEnum) -> TokenStream { #[derive(Debug, Clone)] #input - impl PulseCom for #name { + impl PulseWire for #name { fn to_com(&self) -> Vec { let mut vec = Vec::new(); diff --git a/pulse-ui/src/state.rs b/pulse-ui/src/state.rs index f84f952..4d1d536 100644 --- a/pulse-ui/src/state.rs +++ b/pulse-ui/src/state.rs @@ -83,7 +83,15 @@ impl<'a, T> Drop for StateGuard<'a, T> { } impl Context { + pub async fn event(&self, event: E) { + self.tx.send(Box::new(event)).await.unwrap() + } + pub async fn close(&self) { - self.tx.send(Box::new(Close)).await.unwrap() + self.event(Close).await + } + + pub async fn refresh(&self) { + self.event(Refresh).await } } diff --git a/pulse-ui/src/widget/mod.rs b/pulse-ui/src/widget/mod.rs index dc75b5f..b5cc038 100644 --- a/pulse-ui/src/widget/mod.rs +++ b/pulse-ui/src/widget/mod.rs @@ -1,8 +1,9 @@ pub mod align; -pub mod outline; -pub mod spaced; pub mod dynamic; pub mod input; +pub mod outline; +pub mod scroll; +pub mod spaced; use crate::render::RenderScope; diff --git a/pulse-ui/src/widget/scroll.rs b/pulse-ui/src/widget/scroll.rs new file mode 100644 index 0000000..be6511e --- /dev/null +++ b/pulse-ui/src/widget/scroll.rs @@ -0,0 +1,67 @@ +use crate::widget::Widget; + +pub struct ScrollState(pub usize, pub [usize; N]); + +pub struct ScrollText { + pub scroll: usize, + pub title: String, + pub text: String, +} + +impl Widget for ScrollText { + fn render(&self, scope: &mut crate::render::RenderScope) { + scope.draw_text(0, &self.title); + + let title_lines = self.title.lines().count(); + + for (y, line) in self + .text + .lines() + .skip(self.scroll) + .take(scope.rect.height as usize - title_lines) + .enumerate() + { + scope.draw_text((0, (y + title_lines) as u16), line); + } + } +} + +impl ScrollState { + pub fn get_selected(&self, index: usize) -> &'static str { + if index == self.0 { "\x1b[4m\x1b[1m" } else { "\x1b[1m" } + } + + pub fn scroll(&self, index: usize, title: String, text: String) -> ScrollText { + ScrollText { + scroll: self.1[index], + title, + text, + } + } + + pub fn up(&mut self) { + if self.1[self.0] > 0 { + self.1[self.0] -= 1; + } + } + + pub fn down(&mut self) { + self.1[self.0] += 1; + } + + pub fn back_tab(&mut self) { + if self.0 > 0 { + self.0 -= 1; + } else { + self.0 = N - 1; + } + } + + pub fn tab(&mut self) { + if self.0 + 1 < N { + self.0 += 1; + } else { + self.0 = 0; + } + } +} diff --git a/pulse-wire/Cargo.toml b/pulse-wire/Cargo.toml new file mode 100644 index 0000000..449b94f --- /dev/null +++ b/pulse-wire/Cargo.toml @@ -0,0 +1,7 @@ +[package] +name = "pulse-wire" +version = "0.1.0-alpha.0" +edition = "2024" + +[dependencies] +pulse-macros = { workspace = true } diff --git a/src/pc.rs b/pulse-wire/src/lib.rs similarity index 72% rename from src/pc.rs rename to pulse-wire/src/lib.rs index 8c5cb31..5a53401 100644 --- a/src/pc.rs +++ b/pulse-wire/src/lib.rs @@ -1,9 +1,18 @@ -pub trait PulseCom { +#[cfg(target_os = "macos")] +use std::path::PathBuf; + +pub mod terminal; + +pub fn server_path() -> PathBuf { + PathBuf::from("/tmp/pulse-engine.sock") +} + +pub trait PulseWire { fn to_com(&self) -> Vec; fn from_com(_com: &mut Vec) -> Self; } -impl PulseCom for Vec { +impl PulseWire for Vec { fn to_com(&self) -> Vec { let mut vec = Vec::new(); @@ -31,7 +40,7 @@ impl PulseCom for Vec { } } -impl PulseCom for String { +impl PulseWire for String { fn to_com(&self) -> Vec { let bytes = self.as_bytes(); let mut out = Vec::with_capacity(4 + bytes.len()); @@ -53,7 +62,7 @@ impl PulseCom for String { macro_rules! int_com { ($t:ty) => { - impl PulseCom for $t { + impl $crate::PulseWire for $t { fn to_com(&self) -> Vec { self.to_le_bytes().to_vec() } @@ -81,25 +90,3 @@ int_com!(usize); int_com!(f64); int_com!(f32); - -#[macro_export] -macro_rules! p_com { - (struct $name:ident { $($n:ident: $v:ty),* $(,)? }) => { - #[derive(Debug, Clone)] - pub struct $name { $(pub $n: $v),* } - - impl PulseCom for $name { - fn to_com(&self) -> Vec { - let mut vec = Vec::new(); - $(vec.extend(self.$n.to_com());)* - vec - } - - fn from_com(com: &mut Vec) -> Self { - Self { - $($n: <$v>::from_com(com),)* - } - } - } - }; -} diff --git a/src/ptc.rs b/pulse-wire/src/terminal.rs similarity index 82% rename from src/ptc.rs rename to pulse-wire/src/terminal.rs index 4dc117c..11a87de 100644 --- a/src/ptc.rs +++ b/pulse-wire/src/terminal.rs @@ -1,34 +1,65 @@ -use pulse_macros::p_com; +use crate::PulseWire; +use pulse_macros::pwp; -#[p_com] +#[pwp] +pub enum TerminalServerMessage { + // WatchList + WatchListUpdated(Vec), + + // Positions + PositionsUpdated(Vec), + + // Overview + OverviewUpdated(MarketOverview), + + // Signals + SignalsUpdated(Vec), + + // Inspector + Inspect(InspectTarget), + + // Status + SetStatus(Status), + + // Logs + SetLogs(Vec), + AddLog(EventLog), +} + +#[pwp] +pub enum TerminalClientMessage { + ExecuteCommand(String), +} + +#[pwp] pub struct WatchListItem { symbol: String, price: f64, trend: f64, } -#[p_com] +#[pwp] pub struct ActivePosition { symbol: String, profit: f64, amount: f64, } -#[p_com] +#[pwp] pub enum MarketTrend { Bullish, Bearish, Neutral, } -#[p_com] +#[pwp] pub enum Volatility { Low, Medium, High, } -#[p_com] +#[pwp] pub struct MarketOverview { trend: MarketTrend, volatility: Volatility, @@ -37,7 +68,7 @@ pub struct MarketOverview { alerts: Vec, } -#[p_com] +#[pwp] pub enum Feed { Connected, Disconnected, @@ -45,7 +76,7 @@ pub enum Feed { Failed, } -#[p_com] +#[pwp] pub struct Status { feed: Feed, exchange: String, @@ -53,13 +84,13 @@ pub struct Status { latency: u16, } -#[p_com] +#[pwp] pub enum SignalKind { Buy, Sell, } -#[p_com] +#[pwp] pub enum SignalParameter { Lim, Stl, @@ -67,7 +98,7 @@ pub enum SignalParameter { Chk, } -#[p_com] +#[pwp] pub struct Signal { kind: SignalKind, symbol: String, @@ -75,7 +106,7 @@ pub struct Signal { price: f64, } -#[p_com] +#[pwp] pub enum LogKind { Info, Warn, @@ -83,27 +114,27 @@ pub enum LogKind { Debug, } -#[p_com] +#[pwp] pub struct EventLog { kind: LogKind, name: String, message: String, } -#[p_com] +#[pwp] pub enum AlertLevel { High, Medium, Low, } -#[p_com] +#[pwp] pub struct Alert { level: AlertLevel, message: String, } -#[p_com] +#[pwp] pub enum InspectTarget { None, Symbol(WatchListItem), diff --git a/src/daemon/main.rs b/src/daemon/main.rs deleted file mode 100644 index 61b7dd0..0000000 --- a/src/daemon/main.rs +++ /dev/null @@ -1,22 +0,0 @@ -pub mod ptc; - -use std::time::Instant; - -use ptc::PulseCom; - -fn main() { - let input = ptc::EventLog { - kind: ptc::LogKind::Warn, - name: "Test".to_string(), - message: "Hello, WOrld".to_string(), - }; - - let start = Instant::now(); - - let mut val = input.to_com(); - let out = ptc::EventLog::from_com(&mut val); - - let elapsed = start.elapsed(); - - println!("{:?} {:?}", out, elapsed); -} diff --git a/src/daemon/ptc.rs b/src/daemon/ptc.rs deleted file mode 100644 index ab16651..0000000 --- a/src/daemon/ptc.rs +++ /dev/null @@ -1,2 +0,0 @@ -include!("../pc.rs"); -include!("../ptc.rs"); diff --git a/src/engine/main.rs b/src/engine/main.rs new file mode 100644 index 0000000..85c45ff --- /dev/null +++ b/src/engine/main.rs @@ -0,0 +1,114 @@ +use pulse_wire::terminal::{ActivePosition, Signal, WatchListItem}; + +pub mod terminal; + +#[tokio::main] +async fn main() -> tokio::io::Result<()> { + let terminal_server = terminal::TerminalServer::new(); + + { + let terminal_server = terminal_server.clone(); + + tokio::spawn(async move { + terminal_server + .run() + .await + .expect("Failed to run terminal server"); + }); + } + + loop { + tokio::time::sleep(tokio::time::Duration::from_millis(1500)).await; + + terminal_server + .broadcast( + pulse_wire::terminal::TerminalServerMessage::WatchListUpdated(vec![ + WatchListItem { + symbol: "BTC".to_string(), + price: 118_402.12, + trend: 0.82, + }, + WatchListItem { + symbol: "ETH".to_string(), + price: 3_912.48, + trend: -0.41, + }, + WatchListItem { + symbol: "SOL".to_string(), + price: 182.91, + trend: 2.18, + }, + WatchListItem { + symbol: "XRP".to_string(), + price: 2.84, + trend: 1.22, + }, + ]), + ) + .await?; + + terminal_server + .broadcast( + pulse_wire::terminal::TerminalServerMessage::PositionsUpdated(vec![ + ActivePosition { + symbol: "BTC".to_string(), + profit: 125.50, + amount: 0.25, + }, + ActivePosition { + symbol: "SOL".to_string(), + profit: 84.20, + amount: 5.0, + }, + ActivePosition { + symbol: "ETH".to_string(), + profit: -32.75, + amount: 1.0, + }, + ActivePosition { + symbol: "XRP".to_string(), + profit: 12.30, + amount: 0.75, + }, + ]), + ) + .await?; + + terminal_server + .broadcast(pulse_wire::terminal::TerminalServerMessage::SignalsUpdated( + vec![ + Signal { + kind: pulse_wire::terminal::SignalKind::Buy, + symbol: "BTC".to_string(), + param: pulse_wire::terminal::SignalParameter::Lim, + price: 118_800.0, + }, + Signal { + kind: pulse_wire::terminal::SignalKind::Buy, + symbol: "BTC".to_string(), + param: pulse_wire::terminal::SignalParameter::Tap, + price: 120_000.0, + }, + Signal { + kind: pulse_wire::terminal::SignalKind::Buy, + symbol: "BTC".to_string(), + param: pulse_wire::terminal::SignalParameter::Stl, + price: 118_000.0, + }, + ], + )) + .await?; + + terminal_server + .broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog( + pulse_wire::terminal::EventLog { + kind: pulse_wire::terminal::LogKind::Debug, + name: "Engine".to_string(), + message: "Hello, world!".to_string(), + }, + )) + .await?; + + println!("Ok"); + } +} diff --git a/src/engine/terminal.rs b/src/engine/terminal.rs new file mode 100644 index 0000000..9d2b2cd --- /dev/null +++ b/src/engine/terminal.rs @@ -0,0 +1,122 @@ +use std::sync::Arc; + +use pulse_wire::{PulseWire, terminal::TerminalClientMessage}; +use tokio::{ + io::{AsyncReadExt, AsyncWriteExt}, + net::{ + UnixListener, + unix::{OwnedReadHalf, OwnedWriteHalf}, + }, + sync::Mutex, +}; + +#[derive(Debug, Clone)] +pub struct TerminalServer { + clients: Arc>>, +} + +impl TerminalServer { + pub fn new() -> Self { + Self { + clients: Arc::new(Mutex::new(Vec::new())), + } + } + + pub async fn run(&self) -> tokio::io::Result<()> { + let path = pulse_wire::server_path(); + + if path.exists() { + tokio::fs::remove_file(&path).await?; + } + + let listener = UnixListener::bind(&path)?; + + println!("Terminal server listening on {:?}", path); + + loop { + let (stream, _) = listener.accept().await?; + + let (reader, writer) = stream.into_split(); + + self.clients.lock().await.push(writer); + + let s = self.clone(); + + tokio::spawn(async move { + if let Err(err) = s.handle_client(reader).await { + eprintln!("Terminal connection error: {err}"); + } + }); + } + } + + async fn handle_client(&self, mut reader: OwnedReadHalf) -> tokio::io::Result<()> { + loop { + let mut len_buf = [0u8; size_of::()]; + let size = reader.read_exact(&mut len_buf).await?; + + let len = usize::from_le_bytes(len_buf); + + if size == 0 || len == 0 { + break; + } + + let mut buffer = vec![0u8; len]; + + reader.read_exact(&mut buffer).await?; + + match TerminalClientMessage::from_com(&mut buffer) { + TerminalClientMessage::ExecuteCommand(command) => { + let command = command.as_str(); + + let (command, _args) = if let Some((command, args)) = command.split_once(" ") { + (command, args.split(" ").collect()) + } else { + (command, Vec::new()) + }; + + self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog( + pulse_wire::terminal::EventLog { + kind: pulse_wire::terminal::LogKind::Err, + name: "Command executor".to_string(), + message: format!("Command '{}' not found", command), + }, + )) + .await?; + } + } + } + + Ok(()) + } + + pub async fn broadcast( + &self, + message: pulse_wire::terminal::TerminalServerMessage, + ) -> tokio::io::Result<()> { + let msg = message.to_com(); + + let mut clients = self.clients.lock().await; + + for i in (0..clients.len()).rev() { + if let Err(e) = Self::send_to_client(&mut clients, i, &msg).await { + clients.remove(i); + println!("{e:?}"); + } + } + + Ok(()) + } + + pub async fn send_to_client( + clients: &mut tokio::sync::MutexGuard<'_, Vec>, + i: usize, + msg: &[u8], + ) -> tokio::io::Result<()> { + clients[i].write(&msg.len().to_le_bytes()).await?; + clients[i].write(msg).await?; + clients[i].flush().await?; + + Ok(()) + } +} diff --git a/src/terminal/command.rs b/src/terminal/command.rs index 8a8a4f6..b8849df 100644 --- a/src/terminal/command.rs +++ b/src/terminal/command.rs @@ -1,4 +1,4 @@ -use crate::{PulseTradeApp, ptc::EventLog}; +use crate::PulseTradeApp; impl PulseTradeApp { pub async fn execute_command(&mut self, ctx: &pulse_ui::state::Context, command: &str) { @@ -6,23 +6,23 @@ impl PulseTradeApp { return; } - let (command, _args) = if let Some((command, args)) = command.split_once(" ") { - (command, args.split(" ").collect()) - } else { - (command, Vec::new()) - }; - - match command { + match command.trim() { "exit" | "quit" | "leave" => { ctx.close().await; } _ => { - self.logs.lock().await.push(EventLog { - kind: crate::ptc::LogKind::Err, - name: "cmd".to_string(), - message: format!("Command '{}' not found", command), - }); + if let Err(e) = self + .sock + .as_mut() + .unwrap() + .send(pulse_wire::terminal::TerminalClientMessage::ExecuteCommand( + command.to_string(), + )) + .await + { + eprintln!("{e}"); + } } } } diff --git a/src/terminal/formatting.rs b/src/terminal/formatting.rs index 9b62766..dd431b2 100644 --- a/src/terminal/formatting.rs +++ b/src/terminal/formatting.rs @@ -1,4 +1,4 @@ -use crate::ptc::{ +use pulse_wire::terminal::{ ActivePosition, Alert, AlertLevel, EventLog, InspectTarget, LogKind, MarketOverview, Signal, SignalKind, Status, WatchListItem, }; @@ -137,7 +137,7 @@ struct Property(&'static str, String); impl Formatted for MarketOverview { fn get_formatted(&self) -> Vec { - vec![ + let mut o = vec![ Property("TREND", format!("{}", self.trend)), Property("VOLATILITY", format!("{}", self.volatility)), Property( @@ -149,7 +149,14 @@ impl Formatted for MarketOverview { }, ), ] - .get_formatted() + .get_formatted(); + + o.push(format!( + "\n{}", + apply_padding(self.alerts.get_formatted()).join("\n") + )); + + o } } diff --git a/src/terminal/main.rs b/src/terminal/main.rs index 1b90a56..c0d3faf 100644 --- a/src/terminal/main.rs +++ b/src/terminal/main.rs @@ -1,6 +1,6 @@ pub mod command; pub mod formatting; -pub mod ptc; +pub mod terminal; use std::any::Any; @@ -14,20 +14,22 @@ use pulse_ui::{ align::End, input::{Input, InputState}, outline::{Outline, VLine}, + scroll::{ScrollState, ScrollText}, spaced::SpacedColumns, }, }; -use crate::{ - formatting::{Formatted, apply_padding}, - ptc::{ - ActivePosition, Alert, EventLog, InspectTarget, MarketOverview, Signal, Status, - WatchListItem, - }, +use crate::formatting::{Formatted, apply_padding}; + +use pulse_wire::terminal::{ + ActivePosition, EventLog, InspectTarget, MarketOverview, Signal, Status, WatchListItem, }; pub struct PulseTradeApp { + sock: Option, + command: State, + scroll: State>, watch_list: State>, active_positions: State>, logs: State>, @@ -38,108 +40,12 @@ pub struct PulseTradeApp { } impl App for PulseTradeApp { - async fn init(&mut self, _ctx: &pulse_ui::state::Context) { - let mut watch_list = self.watch_list.lock().await; - - watch_list.push(WatchListItem { - symbol: "BTC".to_string(), - price: 118_402.12, - trend: 0.82, - }); - watch_list.push(WatchListItem { - symbol: "ETH".to_string(), - price: 3_912.48, - trend: -0.41, - }); - watch_list.push(WatchListItem { - symbol: "SOL".to_string(), - price: 182.91, - trend: 2.18, - }); - watch_list.push(WatchListItem { - symbol: "XRP".to_string(), - price: 2.84, - trend: 1.22, - }); - - let mut active_positions = self.active_positions.lock().await; - - active_positions.push(ActivePosition { - symbol: "BTC".to_string(), - profit: 125.50, - amount: 0.25, - }); - active_positions.push(ActivePosition { - symbol: "ETH".to_string(), - profit: -32.75, - amount: 1.0, - }); - active_positions.push(ActivePosition { - symbol: "SOL".to_string(), - profit: 84.20, - amount: 5.0, - }); - active_positions.push(ActivePosition { - symbol: "XRP".to_string(), - profit: 12.30, - amount: 0.75, - }); - - let mut signals = self.signals.lock().await; - - signals.push(Signal { - kind: ptc::SignalKind::Buy, - symbol: "BTC".to_string(), - param: ptc::SignalParameter::Lim, - price: 118_800.0, - }); - - signals.push(Signal { - kind: ptc::SignalKind::Buy, - symbol: "BTC".to_string(), - param: ptc::SignalParameter::Tap, - price: 120_000.0, - }); - - signals.push(Signal { - kind: ptc::SignalKind::Buy, - symbol: "BTC".to_string(), - param: ptc::SignalParameter::Stl, - price: 118_000.0, - }); - - let mut logs = self.logs.lock().await; - - logs.push(EventLog { - kind: ptc::LogKind::Warn, - name: "pulse.init".to_string(), - message: "We're still not done yet ;)".to_string(), - }); - - let mut market_overview = self.market_overview.lock().await; - - market_overview.alerts.push(Alert { - level: ptc::AlertLevel::High, - message: "BTC funding rate elevated".to_string(), - }); - - market_overview.alerts.push(Alert { - level: ptc::AlertLevel::Medium, - message: "Market volatility increasing".to_string(), - }); - - market_overview.alerts.push(Alert { - level: ptc::AlertLevel::Low, - message: "ETH volatility returning to normal".to_string(), - }); - } + async fn init(&mut self, _ctx: &pulse_ui::state::Context) {} async fn update(&mut self, ctx: &pulse_ui::state::Context, event: Box) { if let Some(Refresh) = event.downcast_ref() { return; - } - - if let Some(event) = event.downcast_ref() { + } else if let Some(event) = event.downcast_ref() { if self.command.value.lock().await.handle_event(event) { return; } else if let crossterm::event::Event::Key(key) = event { @@ -151,6 +57,14 @@ impl App for PulseTradeApp { drop(command); self.execute_command(ctx, command_text.trim()).await; + } else if key.code.is_up() { + self.scroll.value.lock().await.up(); + } else if key.code.is_down() { + self.scroll.value.lock().await.down(); + } else if key.code.is_back_tab() { + self.scroll.value.lock().await.back_tab(); + } else if key.code.is_tab() { + self.scroll.value.lock().await.tab(); } } } @@ -209,27 +123,22 @@ impl App for PulseTradeApp { SpacedColumns(vec![ ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " WATCHLIST\n{}", - apply_padding(self.watch_list.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 0, "WATCH LIST", &self.watch_list).await), ), ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " ACTIVE POSITIONS\n{}", - apply_padding(self.active_positions.lock().await.get_formatted()) - .join("\n") - )), + Box::new( + advanced_draw(&self.scroll, 1, "ACTIVE POSITIONS", &self.active_positions) + .await, + ), + ), + ( + LayoutItem::Widget(Size::Flex(1)), + Box::new( + advanced_draw(&self.scroll, 2, "MARKET OVERVIEW", &self.market_overview) + .await, + ), ), - (LayoutItem::Widget(Size::Flex(1)), { - let mo = self.market_overview.lock().await; - Box::new(format!( - " MARKET OVERVIEW\n{}\n\n{}", - apply_padding(mo.get_formatted()).join("\n"), - apply_padding(mo.alerts.get_formatted()).join("\n") - )) - }), ]), ); @@ -238,34 +147,22 @@ impl App for PulseTradeApp { SpacedColumns(vec![ ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " SIGNALS\n{}", - apply_padding(self.signals.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 3, "SIGNALS", &self.signals).await), ), ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " INSPECTOR\n{}", - apply_padding(self.inspect.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 4, "INSPECTOR", &self.inspect).await), ), ( LayoutItem::Widget(Size::Flex(1)), - Box::new(format!( - " STATUS\n{}", - apply_padding(self.status.lock().await.get_formatted()).join("\n") - )), + Box::new(advanced_draw(&self.scroll, 5, "STATUS", &self.status).await), ), ]), ); layout.draw( 5, - format!( - " EVENT LOGS\n{}", - apply_padding(self.logs.lock().await.get_formatted()).join("\n") - ), + advanced_draw(&self.scroll, 6, "EVENT LOGS", &self.logs).await, ); layout.draw(6, Input(" > ", &*self.command.lock().await)); @@ -273,26 +170,49 @@ impl App for PulseTradeApp { } #[tokio::main] -async fn main() { - pulse_ui::run(|ctx| PulseTradeApp { - command: ctx.use_state(InputState::new()), - watch_list: ctx.use_state(Vec::new()), - active_positions: ctx.use_state(Vec::new()), - signals: ctx.use_state(Vec::new()), - logs: ctx.use_state(Vec::new()), - inspect: ctx.use_state(InspectTarget::None), - market_overview: ctx.use_state(MarketOverview { - trend: ptc::MarketTrend::Bullish, - volatility: ptc::Volatility::High, - pressure: 0.324, - alerts: Vec::new(), - }), - status: ctx.use_state(Status { - feed: ptc::Feed::Connected, - exchange: "Binance".to_string(), - dex: "DEX SCREENER".to_string(), - latency: 18, - }), +async fn main() -> tokio::io::Result<()> { + let client = terminal::TerminalClient::new().await?; + + pulse_ui::run(|ctx| { + client.use_app(PulseTradeApp { + sock: None, + command: ctx.use_state(InputState::new()), + scroll: ctx.use_state(ScrollState(1, [0; 7])), + watch_list: ctx.use_state(Vec::new()), + active_positions: ctx.use_state(Vec::new()), + signals: ctx.use_state(Vec::new()), + logs: ctx.use_state(Vec::new()), + inspect: ctx.use_state(InspectTarget::None), + market_overview: ctx.use_state(MarketOverview { + trend: pulse_wire::terminal::MarketTrend::Bullish, + volatility: pulse_wire::terminal::Volatility::High, + pressure: 0.324, + alerts: Vec::new(), + }), + status: ctx.use_state(Status { + feed: pulse_wire::terminal::Feed::Connected, + exchange: "Binance".to_string(), + dex: "DEX SCREENER".to_string(), + latency: 18, + }), + }) }) .await; + + Ok(()) +} + +pub async fn advanced_draw( + scroll: &State>, + index: usize, + title: &'static str, + state: &State, +) -> ScrollText { + let scroll = scroll.lock().await; + + scroll.scroll( + index, + format!(" {}{title}\x1b[0m", scroll.get_selected(index)), + apply_padding(state.lock().await.get_formatted()).join("\n"), + ) } diff --git a/src/terminal/ptc.rs b/src/terminal/ptc.rs deleted file mode 100644 index ab16651..0000000 --- a/src/terminal/ptc.rs +++ /dev/null @@ -1,2 +0,0 @@ -include!("../pc.rs"); -include!("../ptc.rs"); diff --git a/src/terminal/terminal.rs b/src/terminal/terminal.rs new file mode 100644 index 0000000..617979d --- /dev/null +++ b/src/terminal/terminal.rs @@ -0,0 +1,111 @@ +use pulse_wire::{PulseWire, server_path, terminal::TerminalServerMessage}; +use tokio::{ + io::{AsyncReadExt, AsyncWriteExt}, + net::{ + UnixStream, + unix::{OwnedReadHalf, OwnedWriteHalf}, + }, +}; + +use crate::PulseTradeApp; + +pub struct TerminalClient { + writer: OwnedWriteHalf, + reader: Option, +} + +impl TerminalClient { + pub async fn new() -> tokio::io::Result { + let (reader, writer) = UnixStream::connect(server_path()).await?.into_split(); + + Ok(Self { + reader: Some(reader), + writer, + }) + } + + pub async fn send( + &mut self, + message: pulse_wire::terminal::TerminalClientMessage, + ) -> tokio::io::Result<()> { + let msg = message.to_com(); + self.writer.write(&msg.len().to_le_bytes()).await?; + self.writer.write(&msg).await?; + self.writer.flush().await?; + + Ok(()) + } + + pub fn use_app(mut self, mut app: PulseTradeApp) -> PulseTradeApp { + let mut reader = None; + + std::mem::swap(&mut self.reader, &mut reader); + + let mut reader = reader.expect("Reader failed to swap"); + + let watch_list = app.watch_list.clone(); + let active_positions = app.active_positions.clone(); + let logs = app.logs.clone(); + let signals = app.signals.clone(); + let market_overview = app.market_overview.clone(); + let status = app.status.clone(); + let inspect = app.inspect.clone(); + + tokio::spawn(async move { + loop { + let mut len_buf = [0u8; size_of::()]; + reader + .read_exact(&mut len_buf) + .await + .expect("Failed to get header length"); + + let len = usize::from_le_bytes(len_buf); + + let mut buffer = vec![0u8; len]; + + reader + .read_exact(&mut buffer) + .await + .expect("Failed to read socket"); + + match TerminalServerMessage::from_com(&mut buffer) { + TerminalServerMessage::WatchListUpdated(v) => { + *watch_list.lock().await = v; + } + + TerminalServerMessage::PositionsUpdated(v) => { + *active_positions.lock().await = v; + } + + TerminalServerMessage::OverviewUpdated(v) => { + *market_overview.lock().await = v; + } + + TerminalServerMessage::SignalsUpdated(v) => { + *signals.lock().await = v; + } + + TerminalServerMessage::Inspect(v) => { + *inspect.lock().await = v; + } + + TerminalServerMessage::SetStatus(v) => { + *status.lock().await = v; + } + + TerminalServerMessage::SetLogs(v) => { + *logs.lock().await = v; + } + + TerminalServerMessage::AddLog(v) => { + logs.lock().await.push(v); + } + } + } + }); + + app.sock = Some(self); + + app + } +}