Merge pull request #9 from selimaj-dev/engine-communication

Engine communication
This commit is contained in:
2026-07-22 16:04:40 +02:00
committed by GitHub
12 changed files with 485 additions and 156 deletions
Generated
+10
View File
@@ -304,12 +304,22 @@ dependencies = [
"unicode-ident",
]
[[package]]
name = "pulse-macros"
version = "0.1.0-alpha.0"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "pulse-trader"
version = "0.1.0-alpha.0"
dependencies = [
"chrono",
"crossterm",
"pulse-macros",
"pulse-ui",
"tokio",
]
+11 -1
View File
@@ -6,13 +6,23 @@ edition = "2024"
[dependencies]
chrono = "0.4.45"
pulse-ui = { workspace = true }
pulse-macros = { workspace = true }
tokio = { workspace = true, features = ["rt-multi-thread", "macros"] }
crossterm = { workspace = true }
[workspace]
members = ["pulse-ui"]
members = ["pulse-macros", "pulse-ui"]
[workspace.dependencies]
pulse-macros = { path = "pulse-macros", version = "0.1.0-alpha.0" }
pulse-ui = { path = "pulse-ui", version = "0.1.0-alpha.0" }
tokio = "1.52.3"
crossterm = "0.29.0"
[[bin]]
name = "pulse-trader"
path = "src/terminal/main.rs"
[[bin]]
name = "pulse-trader-daemon"
path = "src/daemon/main.rs"
+12
View File
@@ -0,0 +1,12 @@
[package]
name = "pulse-macros"
version = "0.1.0-alpha.0"
edition = "2024"
[lib]
proc-macro = true
[dependencies]
syn = { version = "2", features = ["full"] }
quote = "1"
proc-macro2 = "1"
+170
View File
@@ -0,0 +1,170 @@
use proc_macro::TokenStream;
use quote::quote;
use syn::{Fields, ItemEnum, ItemStruct, parse_macro_input};
#[proc_macro_attribute]
pub fn p_com(_: TokenStream, item: TokenStream) -> TokenStream {
let input = parse_macro_input!(item as syn::Item);
match input {
syn::Item::Struct(s) => expand_struct(s),
syn::Item::Enum(e) => expand_enum(e),
_ => {
return syn::Error::new_spanned(input, "p_com only supports structs and enums")
.to_compile_error()
.into();
}
}
}
fn expand_struct(mut input: ItemStruct) -> TokenStream {
let name = &input.ident;
match &mut input.fields {
Fields::Named(fields) => {
for field in fields.named.iter_mut() {
field.vis = syn::Visibility::Public(syn::token::Pub::default());
}
}
_ => {
return syn::Error::new_spanned(input, "p_com only supports structs with named fields")
.to_compile_error()
.into();
}
}
let fields = match &input.fields {
Fields::Named(fields) => &fields.named,
_ => {
return syn::Error::new_spanned(input, "p_com only supports structs with named fields")
.to_compile_error()
.into();
}
};
let field_names = fields.iter().map(|f| f.ident.as_ref().unwrap());
let field_names2 = fields.iter().map(|f| f.ident.as_ref().unwrap());
TokenStream::from(quote! {
#[derive(Debug, Clone)]
#input
impl PulseCom for #name {
fn to_com(&self) -> Vec<u8> {
let mut vec = Vec::new();
#(
vec.extend(self.#field_names.to_com());
)*
vec
}
fn from_com(com: &mut Vec<u8>) -> Self {
Self {
#(
#field_names2: PulseCom::from_com(com),
)*
}
}
}
})
}
fn expand_enum(input: ItemEnum) -> TokenStream {
let name = &input.ident;
let to_com = input.variants.iter().enumerate().map(|(i, variant)| {
let ident = &variant.ident;
let tag = i as u8;
match &variant.fields {
Fields::Unit => quote! {
Self::#ident => {
vec.push(#tag);
}
},
Fields::Unnamed(fields) if fields.unnamed.len() == 1 => quote! {
Self::#ident(v) => {
vec.push(#tag);
vec.extend(v.to_com());
}
},
Fields::Named(fields) => {
let names = fields.named.iter().map(|f| f.ident.as_ref().unwrap());
let names2 = fields.named.iter().map(|f| f.ident.as_ref().unwrap());
quote! {
Self::#ident { #( #names ),* } => {
vec.push(#tag);
#( vec.extend(#names2.to_com()); )*
}
}
}
_ => {
panic!("tuple variants with >1 field are not supported");
}
}
});
let from_com = input.variants.iter().enumerate().map(|(i, variant)| {
let ident = &variant.ident;
let tag = i as u8;
match &variant.fields {
Fields::Unit => quote! {
#tag => Self::#ident,
},
Fields::Unnamed(fields) if fields.unnamed.len() == 1 => {
quote! {
#tag => Self::#ident(PulseCom::from_com(com)),
}
}
Fields::Named(fields) => {
let names = fields.named.iter().map(|f| f.ident.as_ref().unwrap());
quote! {
#tag => Self::#ident {
#(
#names: PulseCom::from_com(com),
)*
},
}
}
_ => panic!("tuple variants with >1 field are not supported"),
}
});
TokenStream::from(quote! {
#[derive(Debug, Clone)]
#input
impl PulseCom for #name {
fn to_com(&self) -> Vec<u8> {
let mut vec = Vec::new();
match self {
#( #to_com )*
}
vec
}
fn from_com(com: &mut Vec<u8>) -> Self {
let kind = com.remove(0);
match kind {
#( #from_com )*
_ => panic!("invalid {} discriminant {}", stringify!(#name), kind),
}
}
}
})
}
+22
View File
@@ -0,0 +1,22 @@
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);
}
+2
View File
@@ -0,0 +1,2 @@
include!("../pc.rs");
include!("../ptc.rs");
+105
View File
@@ -0,0 +1,105 @@
pub trait PulseCom {
fn to_com(&self) -> Vec<u8>;
fn from_com(_com: &mut Vec<u8>) -> Self;
}
impl<T: PulseCom> PulseCom for Vec<T> {
fn to_com(&self) -> Vec<u8> {
let mut vec = Vec::new();
vec.extend_from_slice(&(self.len() as u32).to_le_bytes());
for item in self {
vec.extend(item.to_com());
}
vec
}
fn from_com(com: &mut Vec<u8>) -> Self {
let len_bytes: [u8; 4] = com.drain(..4).collect::<Vec<_>>().try_into().unwrap();
let len = u32::from_le_bytes(len_bytes) as usize;
let mut result = Vec::with_capacity(len);
for _ in 0..len {
result.push(T::from_com(com));
}
result
}
}
impl PulseCom for String {
fn to_com(&self) -> Vec<u8> {
let bytes = self.as_bytes();
let mut out = Vec::with_capacity(4 + bytes.len());
out.extend_from_slice(&(bytes.len() as u32).to_le_bytes());
out.extend_from_slice(bytes);
out
}
fn from_com(com: &mut Vec<u8>) -> Self {
let len = u32::from_le_bytes(com[..4].try_into().unwrap()) as usize;
com.drain(..4);
let bytes: Vec<u8> = com.drain(..len).collect();
String::from_utf8(bytes).unwrap()
}
}
macro_rules! int_com {
($t:ty) => {
impl PulseCom for $t {
fn to_com(&self) -> Vec<u8> {
self.to_le_bytes().to_vec()
}
fn from_com(com: &mut Vec<u8>) -> Self {
const N: usize = std::mem::size_of::<$t>();
let bytes: [u8; N] = com.drain(..N).collect::<Vec<_>>().try_into().unwrap();
<$t>::from_le_bytes(bytes)
}
}
};
}
int_com!(i8);
int_com!(i16);
int_com!(i32);
int_com!(i64);
int_com!(isize);
int_com!(u8);
int_com!(u16);
int_com!(u32);
int_com!(u64);
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),)*
}
}
}
};
}
+122 -119
View File
@@ -1,24 +1,117 @@
#[derive(Debug, Clone)]
use pulse_macros::p_com;
#[p_com]
pub struct WatchListItem {
pub symbol: String,
pub price: f64,
pub trend: f64,
symbol: String,
price: f64,
trend: f64,
}
#[derive(Debug, Clone)]
#[p_com]
pub struct ActivePosition {
pub symbol: String,
pub profit: f64,
pub amount: f64,
symbol: String,
profit: f64,
amount: f64,
}
#[derive(Debug, Clone, Copy)]
#[p_com]
pub enum MarketTrend {
Bullish,
Bearish,
Neutral,
}
#[p_com]
pub enum Volatility {
Low,
Medium,
High,
}
#[p_com]
pub struct MarketOverview {
trend: MarketTrend,
volatility: Volatility,
pressure: f64,
alerts: Vec<Alert>,
}
#[p_com]
pub enum Feed {
Connected,
Disconnected,
Connecting,
Failed,
}
#[p_com]
pub struct Status {
feed: Feed,
exchange: String,
dex: String,
latency: u16,
}
#[p_com]
pub enum SignalKind {
Buy,
Sell,
}
#[p_com]
pub enum SignalParameter {
Lim,
Stl,
Tap,
Chk,
}
#[p_com]
pub struct Signal {
kind: SignalKind,
symbol: String,
param: SignalParameter,
price: f64,
}
#[p_com]
pub enum LogKind {
Info,
Warn,
Err,
Debug,
}
#[p_com]
pub struct EventLog {
kind: LogKind,
name: String,
message: String,
}
#[p_com]
pub enum AlertLevel {
High,
Medium,
Low,
}
#[p_com]
pub struct Alert {
level: AlertLevel,
message: String,
}
#[p_com]
pub enum InspectTarget {
None,
Symbol(WatchListItem),
Position(ActivePosition),
Signal(Signal),
Alert(Alert),
}
impl std::fmt::Display for MarketTrend {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
@@ -29,11 +122,14 @@ impl std::fmt::Display for MarketTrend {
}
}
#[derive(Debug, Clone, Copy)]
pub enum Volatility {
Low,
Medium,
High,
impl std::fmt::Display for AlertLevel {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::High => write!(f, "H"),
Self::Medium => write!(f, "M"),
Self::Low => write!(f, "L"),
}
}
}
impl std::fmt::Display for Volatility {
@@ -46,80 +142,6 @@ impl std::fmt::Display for Volatility {
}
}
#[derive(Debug, Clone)]
pub struct MarketOverview {
pub trend: MarketTrend,
pub volatility: Volatility,
pub pressure: f64,
pub alerts: Vec<Alert>,
}
#[derive(Debug, Clone, Copy)]
pub enum Feed {
Connected,
Disconnected,
Connecting,
Failed,
}
pub struct Status {
pub feed: Feed,
pub exchange: String,
pub dex: String,
pub latency: u16,
}
#[derive(Debug, Clone, Copy)]
pub enum SignalKind {
Buy,
Sell,
}
impl std::fmt::Display for SignalKind {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Buy => write!(f, "BUY"),
Self::Sell => write!(f, "SELL"),
}
}
}
#[derive(Debug, Clone, Copy)]
pub enum SignalParameter {
Lim,
Stl,
Tap,
Chk,
}
impl std::fmt::Display for SignalParameter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Lim => write!(f, "LIM"),
Self::Stl => write!(f, "STL"),
Self::Tap => write!(f, "TAP"),
Self::Chk => write!(f, "CHK"),
}
}
}
#[derive(Debug, Clone)]
pub struct Signal {
pub kind: SignalKind,
pub symbol: String,
pub param: SignalParameter,
pub price: f64,
}
#[derive(Debug, Clone, Copy)]
pub enum LogKind {
Info,
Warn,
Err,
Debug,
}
impl std::fmt::Display for LogKind {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
@@ -131,41 +153,22 @@ impl std::fmt::Display for LogKind {
}
}
#[derive(Debug, Clone)]
pub struct EventLog {
pub kind: LogKind,
pub name: &'static str,
pub message: String,
}
#[derive(Debug, Clone, Copy)]
pub enum AlertLevel {
High,
Medium,
Low,
}
impl std::fmt::Display for AlertLevel {
impl std::fmt::Display for SignalKind {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::High => write!(f, "H"),
Self::Medium => write!(f, "M"),
Self::Low => write!(f, "L"),
Self::Buy => write!(f, "BUY"),
Self::Sell => write!(f, "SELL"),
}
}
}
#[derive(Debug, Clone)]
pub struct Alert {
pub level: AlertLevel,
pub message: String,
}
#[derive(Debug, Clone)]
pub enum InspectTarget {
None,
Symbol(Box<(WatchListItem, MarketOverview)>),
Position(ActivePosition),
Signal(Signal),
Alert(Alert),
impl std::fmt::Display for SignalParameter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Lim => write!(f, "LIM"),
Self::Stl => write!(f, "STL"),
Self::Tap => write!(f, "TAP"),
Self::Chk => write!(f, "CHK"),
}
}
}
+3 -3
View File
@@ -1,4 +1,4 @@
use crate::{PulseTradeApp, types::EventLog};
use crate::{PulseTradeApp, ptc::EventLog};
impl PulseTradeApp {
pub async fn execute_command(&mut self, ctx: &pulse_ui::state::Context, command: &str) {
@@ -19,8 +19,8 @@ impl PulseTradeApp {
_ => {
self.logs.lock().await.push(EventLog {
kind: crate::types::LogKind::Err,
name: "cmd",
kind: crate::ptc::LogKind::Err,
name: "cmd".to_string(),
message: format!("Command '{}' not found", command),
});
}
@@ -1,4 +1,4 @@
use crate::types::{
use crate::ptc::{
ActivePosition, Alert, AlertLevel, EventLog, InspectTarget, LogKind, MarketOverview, Signal,
SignalKind, Status, WatchListItem,
};
@@ -12,22 +12,15 @@ impl Formatted for InspectTarget {
match self {
Self::None => vec!["\x1b[2mnothing to inspect\x1b[0m".to_string()],
Self::Symbol(boxed) => {
let (watch, overview) = &**boxed;
vec![
Property("Symbol", format!("\x1b[35m{}\x1b[0m", watch.symbol)),
Property("Price", format!("\x1b[96m{}\x1b[0m", format_f64(watch.price))),
Property("Trend", format!("{}", format_f64(watch.trend))),
Property("Market", format!("{}", overview.trend)),
Property("Volatility", format!("{}", overview.volatility)),
Property(
"Pressure",
format!("{:+.2}%", overview.pressure * 100.0),
),
Property("Alerts", format!("{}", overview.alerts.len())),
]
.get_formatted()
}
Self::Symbol(watch) => vec![
Property("Symbol", format!("\x1b[35m{}\x1b[0m", watch.symbol)),
Property(
"Price",
format!("\x1b[96m{}\x1b[0m", format_f64(watch.price)),
),
Property("Trend", format!("{}", format_f64(watch.trend))),
]
.get_formatted(),
Self::Position(pos) => vec![
Property("Symbol", format!("\x1b[35m{}\x1b[0m", pos.symbol)),
+16 -16
View File
@@ -1,6 +1,6 @@
pub mod command;
pub mod formatting;
pub mod types;
pub mod ptc;
use std::any::Any;
@@ -20,7 +20,7 @@ use pulse_ui::{
use crate::{
formatting::{Formatted, apply_padding},
types::{
ptc::{
ActivePosition, Alert, EventLog, InspectTarget, MarketOverview, Signal, Status,
WatchListItem,
},
@@ -88,48 +88,48 @@ impl App for PulseTradeApp {
let mut signals = self.signals.lock().await;
signals.push(Signal {
kind: types::SignalKind::Buy,
kind: ptc::SignalKind::Buy,
symbol: "BTC".to_string(),
param: types::SignalParameter::Lim,
param: ptc::SignalParameter::Lim,
price: 118_800.0,
});
signals.push(Signal {
kind: types::SignalKind::Buy,
kind: ptc::SignalKind::Buy,
symbol: "BTC".to_string(),
param: types::SignalParameter::Tap,
param: ptc::SignalParameter::Tap,
price: 120_000.0,
});
signals.push(Signal {
kind: types::SignalKind::Buy,
kind: ptc::SignalKind::Buy,
symbol: "BTC".to_string(),
param: types::SignalParameter::Stl,
param: ptc::SignalParameter::Stl,
price: 118_000.0,
});
let mut logs = self.logs.lock().await;
logs.push(EventLog {
kind: types::LogKind::Warn,
name: "pulse.init",
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: types::AlertLevel::High,
level: ptc::AlertLevel::High,
message: "BTC funding rate elevated".to_string(),
});
market_overview.alerts.push(Alert {
level: types::AlertLevel::Medium,
level: ptc::AlertLevel::Medium,
message: "Market volatility increasing".to_string(),
});
market_overview.alerts.push(Alert {
level: types::AlertLevel::Low,
level: ptc::AlertLevel::Low,
message: "ETH volatility returning to normal".to_string(),
});
}
@@ -282,13 +282,13 @@ async fn main() {
logs: ctx.use_state(Vec::new()),
inspect: ctx.use_state(InspectTarget::None),
market_overview: ctx.use_state(MarketOverview {
trend: types::MarketTrend::Bullish,
volatility: types::Volatility::High,
trend: ptc::MarketTrend::Bullish,
volatility: ptc::Volatility::High,
pressure: 0.324,
alerts: Vec::new(),
}),
status: ctx.use_state(Status {
feed: types::Feed::Connected,
feed: ptc::Feed::Connected,
exchange: "Binance".to_string(),
dex: "DEX SCREENER".to_string(),
latency: 18,
+2
View File
@@ -0,0 +1,2 @@
include!("../pc.rs");
include!("../ptc.rs");