Merge pull request #10 from selimaj-dev/pwp-and-com
Pwp and communication between the engine and terminal
This commit is contained in:
Generated
+29
-1
@@ -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]]
|
||||
|
||||
+9
-6
@@ -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"
|
||||
|
||||
@@ -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<u8> {
|
||||
let mut vec = Vec::new();
|
||||
|
||||
@@ -63,7 +63,7 @@ fn expand_struct(mut input: ItemStruct) -> TokenStream {
|
||||
fn from_com(com: &mut Vec<u8>) -> 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<u8> {
|
||||
let mut vec = Vec::new();
|
||||
|
||||
|
||||
@@ -83,7 +83,15 @@ impl<'a, T> Drop for StateGuard<'a, T> {
|
||||
}
|
||||
|
||||
impl Context {
|
||||
pub async fn event<E: Any + Send + Sync>(&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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
use crate::widget::Widget;
|
||||
|
||||
pub struct ScrollState<const N: usize>(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<const N: usize> ScrollState<N> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,7 @@
|
||||
[package]
|
||||
name = "pulse-wire"
|
||||
version = "0.1.0-alpha.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
pulse-macros = { workspace = true }
|
||||
@@ -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<u8>;
|
||||
fn from_com(_com: &mut Vec<u8>) -> Self;
|
||||
}
|
||||
|
||||
impl<T: PulseCom> PulseCom for Vec<T> {
|
||||
impl<T: PulseWire> PulseWire for Vec<T> {
|
||||
fn to_com(&self) -> Vec<u8> {
|
||||
let mut vec = Vec::new();
|
||||
|
||||
@@ -31,7 +40,7 @@ impl<T: PulseCom> PulseCom for Vec<T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl PulseCom for String {
|
||||
impl PulseWire for String {
|
||||
fn to_com(&self) -> Vec<u8> {
|
||||
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<u8> {
|
||||
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<u8> {
|
||||
let mut vec = Vec::new();
|
||||
$(vec.extend(self.$n.to_com());)*
|
||||
vec
|
||||
}
|
||||
|
||||
fn from_com(com: &mut Vec<u8>) -> Self {
|
||||
Self {
|
||||
$($n: <$v>::from_com(com),)*
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
@@ -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<WatchListItem>),
|
||||
|
||||
// Positions
|
||||
PositionsUpdated(Vec<ActivePosition>),
|
||||
|
||||
// Overview
|
||||
OverviewUpdated(MarketOverview),
|
||||
|
||||
// Signals
|
||||
SignalsUpdated(Vec<Signal>),
|
||||
|
||||
// Inspector
|
||||
Inspect(InspectTarget),
|
||||
|
||||
// Status
|
||||
SetStatus(Status),
|
||||
|
||||
// Logs
|
||||
SetLogs(Vec<EventLog>),
|
||||
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<Alert>,
|
||||
}
|
||||
|
||||
#[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),
|
||||
@@ -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);
|
||||
}
|
||||
@@ -1,2 +0,0 @@
|
||||
include!("../pc.rs");
|
||||
include!("../ptc.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");
|
||||
}
|
||||
}
|
||||
@@ -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<Mutex<Vec<OwnedWriteHalf>>>,
|
||||
}
|
||||
|
||||
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::<usize>()];
|
||||
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<OwnedWriteHalf>>,
|
||||
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(())
|
||||
}
|
||||
}
|
||||
+13
-13
@@ -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}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<String> {
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+78
-158
@@ -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<terminal::TerminalClient>,
|
||||
|
||||
command: State<InputState>,
|
||||
scroll: State<ScrollState<7>>,
|
||||
watch_list: State<Vec<WatchListItem>>,
|
||||
active_positions: State<Vec<ActivePosition>>,
|
||||
logs: State<Vec<EventLog>>,
|
||||
@@ -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<dyn Any + Send + Sync>) {
|
||||
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<const N: usize, T: Formatted>(
|
||||
scroll: &State<ScrollState<N>>,
|
||||
index: usize,
|
||||
title: &'static str,
|
||||
state: &State<T>,
|
||||
) -> 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"),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -1,2 +0,0 @@
|
||||
include!("../pc.rs");
|
||||
include!("../ptc.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<OwnedReadHalf>,
|
||||
}
|
||||
|
||||
impl TerminalClient {
|
||||
pub async fn new() -> tokio::io::Result<Self> {
|
||||
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::<usize>()];
|
||||
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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user