Merge pull request #16 from selimaj-dev/strategy-engine

Strategy engine
This commit is contained in:
2026-07-28 22:57:37 +02:00
committed by GitHub
26 changed files with 465 additions and 934 deletions
Generated
+88 -14
View File
@@ -1069,6 +1069,15 @@ dependencies = [
"syn 3.0.3", "syn 3.0.3",
] ]
[[package]]
name = "atomic-polyfill"
version = "1.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8cf2bce30dfe09ef0bfaef228b9d414faaf7e563035494d7fe092dba54b300f4"
dependencies = [
"critical-section",
]
[[package]] [[package]]
name = "atomic-waker" name = "atomic-waker"
version = "1.1.2" version = "1.1.2"
@@ -1765,6 +1774,15 @@ version = "0.5.4"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a" checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a"
[[package]]
name = "cobs"
version = "0.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0fa961b519f0b462e3a3b4a34b64d119eeaca1d59af726fe450bbba07a9fc0a1"
dependencies = [
"thiserror 2.0.19",
]
[[package]] [[package]]
name = "coins-ledger" name = "coins-ledger"
version = "0.13.0" version = "0.13.0"
@@ -1936,6 +1954,12 @@ dependencies = [
"cfg-if", "cfg-if",
] ]
[[package]]
name = "critical-section"
version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b"
[[package]] [[package]]
name = "crossbeam-utils" name = "crossbeam-utils"
version = "0.8.22" version = "0.8.22"
@@ -2240,6 +2264,18 @@ dependencies = [
"zeroize", "zeroize",
] ]
[[package]]
name = "embedded-io"
version = "0.4.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ef1a6892d9eef45c8fa6b9e0086428a2cca8491aca8f787c534a3d6d0bcb3ced"
[[package]]
name = "embedded-io"
version = "0.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "edd0f118536f44f5ccd48bcb8b111bdc3de888b58c74639dfb034a357d0f206d"
[[package]] [[package]]
name = "encoding_rs" name = "encoding_rs"
version = "0.8.35" version = "0.8.35"
@@ -2602,6 +2638,15 @@ dependencies = [
"tracing", "tracing",
] ]
[[package]]
name = "hash32"
version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b0c35f58762feb77d74ebe43bdbc3210f09be9fe6742234d573bacc26ed92b67"
dependencies = [
"byteorder",
]
[[package]] [[package]]
name = "hashbrown" name = "hashbrown"
version = "0.12.3" version = "0.12.3"
@@ -2639,6 +2684,20 @@ dependencies = [
"serde_core", "serde_core",
] ]
[[package]]
name = "heapless"
version = "0.7.17"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cdc6457c0eb62c71aac4bc17216026d8410337c4126773b9c5daba343f17964f"
dependencies = [
"atomic-polyfill",
"hash32",
"rustc_version 0.4.1",
"serde",
"spin",
"stable_deref_trait",
]
[[package]] [[package]]
name = "heck" name = "heck"
version = "0.5.0" version = "0.5.0"
@@ -3686,6 +3745,19 @@ version = "0.3.33"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e" checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e"
[[package]]
name = "postcard"
version = "1.1.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6764c3b5dd454e283a30e6dfe78e9b31096d9e32036b5d1eaac7a6119ccb9a24"
dependencies = [
"cobs",
"embedded-io 0.4.0",
"embedded-io 0.6.1",
"heapless",
"serde",
]
[[package]] [[package]]
name = "potential_utf" name = "potential_utf"
version = "0.1.5" version = "0.1.5"
@@ -3870,12 +3942,13 @@ dependencies = [
] ]
[[package]] [[package]]
name = "pulse-macros" name = "pulse-sdk"
version = "0.1.0-alpha.0" version = "0.1.0-alpha.0"
dependencies = [ dependencies = [
"proc-macro2", "hypersdk",
"quote", "postcard",
"syn 2.0.118", "serde",
"tokio",
] ]
[[package]] [[package]]
@@ -3886,8 +3959,9 @@ dependencies = [
"chrono", "chrono",
"crossterm", "crossterm",
"hypersdk", "hypersdk",
"postcard",
"pulse-sdk",
"pulse-ui", "pulse-ui",
"pulse-wire",
"rand 0.8.7", "rand 0.8.7",
"serde", "serde",
"serde_json", "serde_json",
@@ -3903,15 +3977,6 @@ dependencies = [
"tokio", "tokio",
] ]
[[package]]
name = "pulse-wire"
version = "0.1.0-alpha.0"
dependencies = [
"hypersdk",
"pulse-macros",
"serde",
]
[[package]] [[package]]
name = "quinn" name = "quinn"
version = "0.11.11" version = "0.11.11"
@@ -4962,6 +5027,15 @@ dependencies = [
"windows-sys 0.61.2", "windows-sys 0.61.2",
] ]
[[package]]
name = "spin"
version = "0.9.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3763264f6b73151db08c50ff20d7d8a0b8796e021cdea7ceedad07b80155fa0e"
dependencies = [
"lock_api",
]
[[package]] [[package]]
name = "spki" name = "spki"
version = "0.7.3" version = "0.7.3"
+5 -4
View File
@@ -5,11 +5,12 @@ edition = "2024"
[dependencies] [dependencies]
pulse-ui = { workspace = true } pulse-ui = { workspace = true }
pulse-wire = { workspace = true } pulse-sdk = { workspace = true }
crossterm = { workspace = true } crossterm = { workspace = true }
hypersdk = { workspace = true } hypersdk = { workspace = true }
serde = { workspace = true } serde = { workspace = true }
anyhow = { workspace = true } anyhow = { workspace = true }
postcard = { workspace = true }
tokio = { workspace = true, features = [ tokio = { workspace = true, features = [
"rt-multi-thread", "rt-multi-thread",
"macros", "macros",
@@ -25,12 +26,12 @@ rand = "0.8.7"
toml = "1.1.3" toml = "1.1.3"
[workspace] [workspace]
members = ["pulse-macros", "pulse-ui", "pulse-wire"] members = ["pulse-ui", "pulse-sdk"]
[workspace.dependencies] [workspace.dependencies]
pulse-macros = { path = "pulse-macros", version = "0.1.0-alpha.0" } postcard = { version = "1.1.3", features = ["alloc"] }
pulse-ui = { path = "pulse-ui", 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" } pulse-sdk = { path = "pulse-sdk", version = "0.1.0-alpha.0" }
hypersdk = "0.2.14" hypersdk = "0.2.14"
tokio = "1.52.3" tokio = "1.52.3"
crossterm = "0.29.0" crossterm = "0.29.0"
-12
View File
@@ -1,12 +0,0 @@
[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
@@ -1,170 +0,0 @@
use proc_macro::TokenStream;
use quote::quote;
use syn::{Fields, ItemEnum, ItemStruct, parse_macro_input};
#[proc_macro_attribute]
pub fn pwp(_: 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 PulseWire 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: PulseWire::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(PulseWire::from_com(com)),
}
}
Fields::Named(fields) => {
let names = fields.named.iter().map(|f| f.ident.as_ref().unwrap());
quote! {
#tag => Self::#ident {
#(
#names: PulseWire::from_com(com),
)*
},
}
}
_ => panic!("tuple variants with >1 field are not supported"),
}
});
TokenStream::from(quote! {
#[derive(Debug, Clone)]
#input
impl PulseWire 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),
}
}
}
})
}
@@ -1,9 +1,10 @@
[package] [package]
name = "pulse-wire" name = "pulse-sdk"
version = "0.1.0-alpha.0" version = "0.1.0-alpha.0"
edition = "2024" edition = "2024"
[dependencies] [dependencies]
pulse-macros = { workspace = true }
serde = { workspace = true } serde = { workspace = true }
hypersdk = { workspace = true } hypersdk = { workspace = true }
postcard = { workspace = true }
tokio = { workspace = true }
@@ -1,18 +1,13 @@
use pulse_macros::pwp; use crate::units::{Direction, Symbol, USD};
use crate::{ #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
PulseWire,
units::{Direction, Symbol, USD},
};
#[pwp]
pub enum MarketTrend { pub enum MarketTrend {
Bullish, Bullish,
Bearish, Bearish,
Neutral, Neutral,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum LogKind { pub enum LogKind {
Info, Info,
Warn, Warn,
@@ -20,7 +15,7 @@ pub enum LogKind {
Debug, Debug,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Signal { pub struct Signal {
pub symbol: String, pub symbol: String,
pub kind: Direction, pub kind: Direction,
@@ -31,14 +26,14 @@ pub struct Signal {
pub stop_loss: USD, pub stop_loss: USD,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct EventLog { pub struct EventLog {
pub kind: LogKind, pub kind: LogKind,
pub name: String, pub name: String,
pub message: String, pub message: String,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Position { pub struct Position {
pub symbol: Symbol, pub symbol: Symbol,
pub size: f64, pub size: f64,
+25
View File
@@ -0,0 +1,25 @@
#[cfg(target_os = "macos")]
use std::path::PathBuf;
pub mod general;
pub mod plugin;
pub mod terminal;
pub mod units;
pub use hypersdk;
pub mod prelude {
pub use crate::general::*;
pub use crate::plugin::*;
pub use crate::server_path;
pub use crate::terminal::*;
pub use crate::units::*;
pub use hypersdk;
}
pub fn server_path() -> PathBuf {
PathBuf::from("/tmp/pulse-engine.sock")
}
pub fn map_postcard_err<T>(res: postcard::Result<T>) -> tokio::io::Result<T> {
res.map_err(|e| tokio::io::Error::new(std::io::ErrorKind::Other, e))
}
+100
View File
@@ -0,0 +1,100 @@
use crate::{
general::{EventLog, Signal},
terminal::MarketItem,
units::Direction,
};
use hypersdk::hypercore::{Candle, CandleInterval, Subscription};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StrategyManifest {
pub name: String,
pub description: String,
pub author: String,
pub version: String,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct RiskManifest {
pub name: String,
pub description: String,
pub author: String,
pub version: String,
pub max_loss: u8,
pub cooldown: CandleInterval,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum StrategyMessage {
Log(EventLog),
GetWatchList,
GetCandlestick {
symbol: String,
interval: CandleInterval,
count: u32,
},
Subscribe(Subscription),
Unsubscribe(Subscription),
UnsubscribeAll,
Signal(StrategySignal),
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum RiskMessage {
Log(EventLog),
GetWatchList,
Approve(Signal),
Reject { reason: String },
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum StrategyEngineMessage {
Initialize,
WatchList(Vec<MarketItem>),
Command {
command: String,
args: Vec<String>,
},
CandleUpdate {
symbol: String,
interval: CandleInterval,
candle: Candle,
},
Candlestick {
symbol: String,
interval: CandleInterval,
candles: Vec<Candle>,
},
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum RiskEngineMessage {
Initialize,
WatchList(Vec<MarketItem>),
Command { command: String, args: Vec<String> },
Signal(StrategySignal),
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct StrategySignal {
pub symbol: String,
pub side: Direction,
pub confidence: f32,
pub price: Option<f64>,
}
@@ -1,13 +1,11 @@
use crate::{ use crate::{
PulseWire,
general::{EventLog, MarketTrend, Position, Signal}, general::{EventLog, MarketTrend, Position, Signal},
plugin::{RiskManifest, StrategyManifest}, plugin::{RiskManifest, StrategyManifest},
units::{Symbol, USD, Volatility}, units::{Symbol, USD, Volatility},
}; };
use hypersdk::hypercore::CandleInterval; use hypersdk::hypercore::CandleInterval;
use pulse_macros::pwp;
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum TerminalServerMessage { pub enum TerminalServerMessage {
// WatchList // WatchList
WatchListUpdated(Vec<MarketItem>), WatchListUpdated(Vec<MarketItem>),
@@ -32,12 +30,12 @@ pub enum TerminalServerMessage {
AddLog(EventLog), AddLog(EventLog),
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum TerminalClientMessage { pub enum TerminalClientMessage {
ExecuteCommand(String), ExecuteCommand(String),
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MarketItem { pub struct MarketItem {
pub symbol: Symbol, pub symbol: Symbol,
pub price: USD, pub price: USD,
@@ -45,7 +43,7 @@ pub struct MarketItem {
pub volume_24h: USD, pub volume_24h: USD,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MarketOverview { pub struct MarketOverview {
pub trend: MarketTrend, pub trend: MarketTrend,
pub volatility: Volatility, pub volatility: Volatility,
@@ -53,32 +51,32 @@ pub struct MarketOverview {
pub alerts: Vec<Alert>, pub alerts: Vec<Alert>,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum AlertLevel { pub enum AlertLevel {
High, High,
Medium, Medium,
Low, Low,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Alert { pub struct Alert {
level: AlertLevel, pub level: AlertLevel,
message: String, pub message: String,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Balance { pub struct Balance {
pub asset: String, pub asset: String,
pub amount: f64, pub amount: f64,
pub value: f64, pub value: f64,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum InspectTarget { pub enum InspectTarget {
None, None,
Some(Vec<InspectItem>), Some(Vec<InspectItem>),
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum InspectItem { pub enum InspectItem {
String(String), String(String),
Symbol(String), Symbol(String),
@@ -86,35 +84,35 @@ pub enum InspectItem {
F64(f64), F64(f64),
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Status { pub struct Status {
feed: Mode, pub feed: Mode,
exchange: String, pub exchange: String,
dex: String, pub dex: String,
latency: u16, pub latency: u16,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum Mode { pub enum Mode {
Auto, Auto,
Manual, Manual,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum ItemState { pub enum ItemState {
Running, Running,
Stopped, Stopped,
Error, Error,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Strategy { pub struct Strategy {
strategy: StrategyManifest, pub strategy: StrategyManifest,
risk: RiskManifest, pub risk: RiskManifest,
mode: Mode, pub mode: Mode,
state: ItemState, pub state: ItemState,
cooldown: CandleInterval, pub cooldown: CandleInterval,
} }
impl std::fmt::Display for AlertLevel { impl std::fmt::Display for AlertLevel {
@@ -1,20 +1,16 @@
use pulse_macros::pwp; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
use crate::PulseWire;
#[derive(Debug, Clone)]
pub struct Symbol(pub String); pub struct Symbol(pub String);
#[derive(Debug, Clone, Copy)] #[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
pub struct USD(pub f64); pub struct USD(pub f64);
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum Direction { pub enum Direction {
Buy, Buy,
Sell, Sell,
} }
#[pwp] #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum Volatility { pub enum Volatility {
Low, Low,
Medium, Medium,
@@ -37,26 +33,6 @@ impl std::fmt::Display for USD {
} }
} }
impl PulseWire for Symbol {
fn from_com(com: &mut Vec<u8>) -> Self {
Self(String::from_com(com))
}
fn to_com(&self) -> Vec<u8> {
self.0.to_com()
}
}
impl PulseWire for USD {
fn from_com(com: &mut Vec<u8>) -> Self {
Self(f64::from_com(com))
}
fn to_com(&self) -> Vec<u8> {
self.0.to_com()
}
}
pub fn format_f64(value: f64) -> String { pub fn format_f64(value: f64) -> String {
let abs = value.abs(); let abs = value.abs();
+50 -2
View File
@@ -16,9 +16,57 @@ pub struct RenderScope {
} }
impl RenderScope { impl RenderScope {
pub fn draw_text<P: Into<Point>, T: Display>(&mut self, at: P, text: T) { pub fn draw_text<P: Into<Point>, T: Display>(&mut self, at: P, text: T) -> u16 {
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;
while let Some(c) = chars.next() {
if c == '\x1b' && chars.peek() == Some(&'[') {
new_line.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;
}
}
new_line
})
.collect::<Vec<_>>()
.join("\n");
let lines = lines
.lines()
.take((self.rect.height - point.y) as usize)
.collect::<Vec<_>>();
let lines_len = lines.len();
self.draw_instructions self.draw_instructions
.push(Instr::DrawText(at.into(), text.to_string())); .push(Instr::DrawText(point, lines.join("\n")));
lines_len as u16
} }
} }
+9 -4
View File
@@ -14,21 +14,26 @@ impl Widget for ScrollText {
let title_lines = self.title.lines().count(); let title_lines = self.title.lines().count();
for (y, line) in self let mut y = title_lines as u16;
for line in self
.text .text
.lines() .lines()
.skip(self.scroll) .skip(self.scroll)
.take(scope.rect.height as usize - title_lines) .take(scope.rect.height as usize - title_lines)
.enumerate()
{ {
scope.draw_text((0, (y + title_lines) as u16), line); y += scope.draw_text((0, y), line);
} }
} }
} }
impl<const N: usize> ScrollState<N> { impl<const N: usize> ScrollState<N> {
pub fn get_selected(&self, index: usize) -> &'static str { pub fn get_selected(&self, index: usize) -> &'static str {
if index == self.0 { "\x1b[4m\x1b[1m" } else { "\x1b[1m" } if index == self.0 {
"\x1b[4m\x1b[1m"
} else {
"\x1b[1m"
}
} }
pub fn scroll(&self, index: usize, title: String, text: String) -> ScrollText { pub fn scroll(&self, index: usize, title: String, text: String) -> ScrollText {
-313
View File
@@ -1,313 +0,0 @@
use hypersdk::{
Address, Decimal,
hypercore::{Candle, CandleInterval, Subscription},
};
use crate::PulseWire;
impl PulseWire for Subscription {
fn from_com(com: &mut Vec<u8>) -> Self {
match u8::from_com(com) {
0 => Self::Bbo {
coin: PulseWire::from_com(com),
},
1 => Self::Trades {
coin: PulseWire::from_com(com),
},
2 => Self::L2Book {
coin: PulseWire::from_com(com),
n_sig_figs: PulseWire::from_com(com),
mantissa: PulseWire::from_com(com),
fast: PulseWire::from_com(com),
},
3 => Self::Candle {
coin: PulseWire::from_com(com),
interval: PulseWire::from_com(com),
},
4 => Self::AllMids {
dex: PulseWire::from_com(com),
},
5 => Self::OrderUpdates {
user: PulseWire::from_com(com),
},
6 => Self::UserFills {
user: PulseWire::from_com(com),
},
7 => Self::UserEvents {
user: PulseWire::from_com(com),
},
8 => Self::UserTwapSliceFills {
user: PulseWire::from_com(com),
},
9 => Self::UserTwapHistory {
user: PulseWire::from_com(com),
},
10 => Self::ActiveAssetCtx {
coin: PulseWire::from_com(com),
},
11 => Self::ActiveAssetData {
user: PulseWire::from_com(com),
coin: PulseWire::from_com(com),
},
12 => Self::WebData2 {
user: PulseWire::from_com(com),
dex: PulseWire::from_com(com),
},
13 => Self::ClearinghouseState {
user: PulseWire::from_com(com),
dex: PulseWire::from_com(com),
},
14 => Self::AllDexsClearinghouseState {
user: PulseWire::from_com(com),
},
15 => Self::OpenOrders {
user: PulseWire::from_com(com),
dex: PulseWire::from_com(com),
},
16 => Self::SpotState {
user: PulseWire::from_com(com),
is_portfolio_margin: PulseWire::from_com(com),
},
17 => Self::Notification {
user: PulseWire::from_com(com),
},
18 => Self::WebData3 {
user: PulseWire::from_com(com),
},
19 => Self::TwapStates {
user: PulseWire::from_com(com),
dex: Option::from_com(com),
},
20 => Self::UserFundings {
user: PulseWire::from_com(com),
},
21 => Self::UserNonFundingLedgerUpdates {
user: PulseWire::from_com(com),
},
22 => Self::AllDexsAssetCtxs,
23 => Self::FastAssetCtxs,
24 => Self::OutcomeMetaUpdates,
x => panic!("Invalid Subscription discriminant: {}", x),
}
}
fn to_com(&self) -> Vec<u8> {
let mut com = Vec::new();
match self {
Self::Bbo { coin } => {
com.push(0);
com.extend(coin.to_com());
}
Self::Trades { coin } => {
com.push(1);
com.extend(coin.to_com());
}
Self::L2Book {
coin,
n_sig_figs,
mantissa,
fast,
} => {
com.push(2);
com.extend(coin.to_com());
com.extend(n_sig_figs.to_com());
com.extend(mantissa.to_com());
com.extend(fast.to_com());
}
Self::Candle { coin, interval } => {
com.push(3);
com.extend(coin.to_com());
com.extend(interval.to_com());
}
Self::AllMids { dex } => {
com.push(4);
com.extend(dex.to_com());
}
Self::OrderUpdates { user } => {
com.push(5);
com.extend(user.to_com());
}
Self::UserFills { user } => {
com.push(6);
com.extend(user.to_com());
}
Self::UserEvents { user } => {
com.push(7);
com.extend(user.to_com());
}
Self::UserTwapSliceFills { user } => {
com.push(8);
com.extend(user.to_com());
}
Self::UserTwapHistory { user } => {
com.push(9);
com.extend(user.to_com());
}
Self::ActiveAssetCtx { coin } => {
com.push(10);
com.extend(coin.to_com());
}
Self::ActiveAssetData { user, coin } => {
com.push(11);
com.extend(user.to_com());
com.extend(coin.to_com());
}
Self::WebData2 { user, dex } => {
com.push(12);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::ClearinghouseState { user, dex } => {
com.push(13);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::AllDexsClearinghouseState { user } => {
com.push(14);
com.extend(user.to_com());
}
Self::OpenOrders { user, dex } => {
com.push(15);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::SpotState {
user,
is_portfolio_margin,
} => {
com.push(16);
com.extend(user.to_com());
com.extend(is_portfolio_margin.to_com());
}
Self::Notification { user } => {
com.push(17);
com.extend(user.to_com());
}
Self::WebData3 { user } => {
com.push(18);
com.extend(user.to_com());
}
Self::TwapStates { user, dex } => {
com.push(19);
com.extend(user.to_com());
com.extend(dex.to_com());
}
Self::UserFundings { user } => {
com.push(20);
com.extend(user.to_com());
}
Self::UserNonFundingLedgerUpdates { user } => {
com.push(21);
com.extend(user.to_com());
}
Self::AllDexsAssetCtxs => {
com.push(22);
}
Self::FastAssetCtxs => {
com.push(23);
}
Self::OutcomeMetaUpdates => {
com.push(24);
}
}
com
}
}
impl PulseWire for Address {
fn from_com(com: &mut Vec<u8>) -> Self {
let bytes: [u8; 20] = com.drain(..20).collect::<Vec<_>>().try_into().unwrap();
Self::from_slice(&bytes)
}
fn to_com(&self) -> Vec<u8> {
self.as_slice().to_vec()
}
}
impl PulseWire for CandleInterval {
fn from_com(com: &mut Vec<u8>) -> Self {
match u8::from_com(com) {
0 => Self::OneMinute,
1 => Self::ThreeMinutes,
2 => Self::FiveMinutes,
3 => Self::FifteenMinutes,
4 => Self::ThirtyMinutes,
5 => Self::OneHour,
6 => Self::TwoHours,
7 => Self::FourHours,
8 => Self::EightHours,
9 => Self::TwelveHours,
10 => Self::OneDay,
11 => Self::ThreeDays,
12 => Self::OneWeek,
13 => Self::OneMonth,
x => panic!("Invalid CandleInterval discriminant: {}", x),
}
}
fn to_com(&self) -> Vec<u8> {
vec![match self {
Self::OneMinute => 0,
Self::ThreeMinutes => 1,
Self::FiveMinutes => 2,
Self::FifteenMinutes => 3,
Self::ThirtyMinutes => 4,
Self::OneHour => 5,
Self::TwoHours => 6,
Self::FourHours => 7,
Self::EightHours => 8,
Self::TwelveHours => 9,
Self::OneDay => 10,
Self::ThreeDays => 11,
Self::OneWeek => 12,
Self::OneMonth => 13,
}]
}
}
impl PulseWire for Candle {
fn from_com(com: &mut Vec<u8>) -> Self {
Self {
open_time: u64::from_com(com),
close_time: u64::from_com(com),
coin: String::from_com(com),
interval: String::from_com(com),
open: Decimal::from_com(com),
high: Decimal::from_com(com),
low: Decimal::from_com(com),
close: Decimal::from_com(com),
volume: Decimal::from_com(com),
num_trades: u64::from_com(com),
}
}
fn to_com(&self) -> Vec<u8> {
let mut com = Vec::new();
com.extend(self.open_time.to_com());
com.extend(self.close_time.to_com());
com.extend(self.coin.to_com());
com.extend(self.interval.to_com());
com.extend(self.open.to_com());
com.extend(self.high.to_com());
com.extend(self.low.to_com());
com.extend(self.close.to_com());
com.extend(self.volume.to_com());
com.extend(self.num_trades.to_com());
com
}
}
impl PulseWire for Decimal {
fn from_com(com: &mut Vec<u8>) -> Self {
Self::deserialize(com[..16].try_into().unwrap())
}
fn to_com(&self) -> Vec<u8> {
self.serialize().to_vec()
}
}
-135
View File
@@ -1,135 +0,0 @@
#[cfg(target_os = "macos")]
use std::path::PathBuf;
pub mod general;
mod hyper_types;
pub mod plugin;
pub mod terminal;
pub mod units;
pub use hypersdk;
pub mod prelude {
pub use crate::PulseWire;
pub use crate::general::*;
pub use crate::plugin::*;
pub use crate::server_path;
pub use crate::terminal::*;
pub use crate::units::*;
pub use hypersdk;
}
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: PulseWire> PulseWire 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<T: PulseWire> PulseWire for Option<T> {
fn to_com(&self) -> Vec<u8> {
if let Some(v) = self {
vec![1].into_iter().chain(v.to_com()).collect()
} else {
vec![0]
}
}
fn from_com(com: &mut Vec<u8>) -> Self {
if com[0] > 0 {
Some(T::from_com(com))
} else {
None
}
}
}
impl PulseWire 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()
}
}
impl PulseWire for bool {
fn to_com(&self) -> Vec<u8> {
if *self { vec![1] } else { vec![0] }
}
fn from_com(com: &mut Vec<u8>) -> Self {
com[0] > 0
}
}
macro_rules! int_com {
($t:ty) => {
impl $crate::PulseWire 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);
-96
View File
@@ -1,96 +0,0 @@
use crate::{
PulseWire,
general::{EventLog, Signal},
terminal::MarketItem,
units::Direction,
};
use hypersdk::hypercore::{Candle, CandleInterval, Subscription};
use pulse_macros::pwp;
#[pwp]
#[derive(serde::Deserialize, serde::Serialize)]
pub struct StrategyManifest {
name: String,
description: String,
author: String,
version: String,
}
#[pwp]
#[derive(serde::Deserialize, serde::Serialize)]
pub struct RiskManifest {
name: String,
description: String,
author: String,
version: String,
max_loss: u8,
cooldown: CandleInterval,
}
#[pwp]
pub enum StrategyMessage {
RequestOHLC {
symbol: String,
interval: CandleInterval,
count: u32,
},
Subscribe(Subscription),
Unsubscribe(Subscription),
UnsubscribeAll,
Signal(StrategySignal),
Log(EventLog),
}
#[pwp]
pub enum StrategyEngineMessage {
Initialize {
watchlist: Vec<MarketItem>,
},
CandleUpdate {
symbol: String,
candle: String,
},
OHLC {
symbol: String,
interval: CandleInterval,
candles: Vec<Candle>,
},
Start,
Stop,
}
#[pwp]
pub enum RiskMessage {
Approve(Signal),
Reject { reason: String },
Log(EventLog),
}
#[pwp]
pub enum RiskEngineMessage {
Initialize { strategy: String },
Signal(StrategySignal),
MarketUpdate { symbol: String, price: f64 },
}
#[pwp]
pub struct StrategySignal {
pub symbol: String,
pub side: Direction,
pub confidence: f32,
pub price: Option<f64>,
}
+2 -2
View File
@@ -1,6 +1,6 @@
use crate::{engine::Engine, store::config::Config}; use crate::{engine::Engine, store::config::Config};
use pulse_wire::terminal::{ItemState, Mode, Strategy}; use pulse_sdk::terminal::{ItemState, Mode, Strategy};
use toml::Value; use toml::Value;
impl Engine { impl Engine {
@@ -100,7 +100,7 @@ impl Engine {
self.config.lock().await.cooldown = cooldown; self.config.lock().await.cooldown = cooldown;
self.terminal_server.broadcast( self.terminal_server.broadcast(
pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: self.strategy.strategy.manifest.lock().await.clone(), strategy: self.strategy.strategy.manifest.lock().await.clone(),
risk: self.strategy.risk.manifest.lock().await.clone(), risk: self.strategy.risk.manifest.lock().await.clone(),
mode: Mode::Auto, mode: Mode::Auto,
+8 -16
View File
@@ -1,12 +1,12 @@
pub mod plugin;
pub mod command; pub mod command;
pub mod plugin;
use crate::{ use crate::{
engine::plugin::StrategyEngine, engine::plugin::StrategyEngine,
store::{accounts::AccountList, config::Config}, store::{accounts::AccountList, config::Config},
terminal::TerminalServer, terminal::TerminalServer,
}; };
use pulse_wire::prelude::*; use pulse_sdk::prelude::*;
use std::sync::Arc; use std::sync::Arc;
use tokio::{sync::Mutex, task::JoinHandle}; use tokio::{sync::Mutex, task::JoinHandle};
@@ -16,6 +16,7 @@ pub struct Engine {
pub strategy: Arc<StrategyEngine>, pub strategy: Arc<StrategyEngine>,
pub config: Arc<Mutex<Config>>, pub config: Arc<Mutex<Config>>,
pub accounts: Arc<Mutex<AccountList>>, pub accounts: Arc<Mutex<AccountList>>,
pub watch_list: Arc<Mutex<Vec<MarketItem>>>,
} }
impl Engine { impl Engine {
@@ -32,15 +33,10 @@ impl Engine {
strategy: strategy.initialize(engine.clone()), strategy: strategy.initialize(engine.clone()),
config, config,
accounts, accounts,
watch_list: Arc::new(Mutex::new(Vec::new())),
})) }))
} }
pub async fn spawn_terminal_server(&self) -> JoinHandle<tokio::io::Result<()>> {
let terminal_server = self.terminal_server.clone();
tokio::spawn(async move { terminal_server.run().await })
}
pub async fn spawn_broadcaster(&self) -> JoinHandle<tokio::io::Result<()>> { pub async fn spawn_broadcaster(&self) -> JoinHandle<tokio::io::Result<()>> {
let s = self.clone(); let s = self.clone();
@@ -59,10 +55,12 @@ impl Engine {
match crate::fetch::fetch_watch_list(&client, watch_list).await { match crate::fetch::fetch_watch_list(&client, watch_list).await {
Ok(watch_list) => { Ok(watch_list) => {
*self.watch_list.lock().await = watch_list.clone();
if let Err(error) = self if let Err(error) = self
.terminal_server .terminal_server
.broadcast( .broadcast(
pulse_wire::terminal::TerminalServerMessage::WatchListUpdated( pulse_sdk::terminal::TerminalServerMessage::WatchListUpdated(
watch_list, watch_list,
), ),
) )
@@ -94,7 +92,7 @@ impl Engine {
Ok(state) => { Ok(state) => {
self.terminal_server self.terminal_server
.broadcast( .broadcast(
pulse_wire::terminal::TerminalServerMessage::PositionsUpdated( pulse_sdk::terminal::TerminalServerMessage::PositionsUpdated(
state state
.asset_positions .asset_positions
.into_iter() .into_iter()
@@ -128,12 +126,6 @@ impl Engine {
} }
} }
pub async fn run_engine(&self) -> tokio::io::Result<()> {
loop {
tokio::time::sleep(tokio::time::Duration::from_millis(5000)).await;
}
}
pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> { pub async fn invalid_command_usage(&self, name: &str) -> tokio::io::Result<()> {
self.terminal_server.error(name, "Invalid usage").await self.terminal_server.error(name, "Invalid usage").await
} }
+35 -10
View File
@@ -1,5 +1,5 @@
use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket}; use hypersdk::hypercore::{self, CandleInterval, Subscription, WebSocket};
use pulse_wire::prelude::*; use pulse_sdk::prelude::*;
use std::{ use std::{
collections::HashSet, collections::HashSet,
path::PathBuf, path::PathBuf,
@@ -65,20 +65,29 @@ impl StrategyEngine {
.upgrade() .upgrade()
.expect("Failed to upgrade engine (StrategyEngine)"); .expect("Failed to upgrade engine (StrategyEngine)");
let strategy = self.strategy.clone(); self.strategy
let risk = self.risk.clone(); .send(&StrategyEngineMessage::Initialize)
.await?;
loop { loop {
match strategy.recv().await? { match self.strategy.recv().await? {
None => {} None => {}
Some(StrategyMessage::GetWatchList) => {
self.strategy
.send(&StrategyEngineMessage::WatchList(
engine.watch_list.lock().await.clone(),
))
.await?;
}
Some(StrategyMessage::Log(mut log)) => { Some(StrategyMessage::Log(mut log)) => {
log.name.insert_str(0, "strategy::"); log.name.insert_str(0, "strategy::");
engine.terminal_server.log_raw(log).await?; engine.terminal_server.log_raw(log).await?;
} }
Some(StrategyMessage::Signal(signal)) => { Some(StrategyMessage::Signal(signal)) => {
risk.send(&RiskEngineMessage::Signal(signal)).await?; self.risk.send(&RiskEngineMessage::Signal(signal)).await?;
} }
Some(StrategyMessage::Subscribe(subscription)) => { Some(StrategyMessage::Subscribe(subscription)) => {
@@ -97,7 +106,7 @@ impl StrategyEngine {
} }
} }
Some(StrategyMessage::RequestOHLC { Some(StrategyMessage::GetCandlestick {
symbol, symbol,
interval, interval,
count, count,
@@ -109,6 +118,8 @@ impl StrategyEngine {
.unwrap() .unwrap()
.as_millis() as u64; .as_millis() as u64;
println!("getting candlestick");
let interval_ms = match interval { let interval_ms = match interval {
CandleInterval::OneMinute => 60_000, CandleInterval::OneMinute => 60_000,
CandleInterval::ThreeMinutes => 3 * 60_000, CandleInterval::ThreeMinutes => 3 * 60_000,
@@ -128,8 +139,14 @@ impl StrategyEngine {
let start_time = now.saturating_sub(interval_ms * count as u64); let start_time = now.saturating_sub(interval_ms * count as u64);
client self.strategy
.candle_snapshot(symbol, interval, start_time, now) .send(&StrategyEngineMessage::Candlestick {
candles: client
.candle_snapshot(&symbol, interval, start_time, now)
.await?,
symbol,
interval,
})
.await?; .await?;
} }
} }
@@ -142,12 +159,20 @@ impl StrategyEngine {
.upgrade() .upgrade()
.expect("Failed to upgrade engine (StrategyEngine)"); .expect("Failed to upgrade engine (StrategyEngine)");
let risk = self.risk.clone(); self.risk.send(&RiskEngineMessage::Initialize).await?;
loop { loop {
match risk.recv().await? { match self.risk.recv().await? {
None => {} None => {}
Some(RiskMessage::GetWatchList) => {
self.risk
.send(&&RiskEngineMessage::WatchList(
engine.watch_list.lock().await.clone(),
))
.await?;
}
Some(RiskMessage::Log(mut log)) => { Some(RiskMessage::Log(mut log)) => {
log.name.insert_str(0, "risk::"); log.name.insert_str(0, "risk::");
engine.terminal_server.log_raw(log).await?; engine.terminal_server.log_raw(log).await?;
+1 -1
View File
@@ -1,4 +1,4 @@
use pulse_wire::prelude::*; use pulse_sdk::prelude::*;
use serde_json::Value; use serde_json::Value;
use std::collections::HashMap; use std::collections::HashMap;
+2 -7
View File
@@ -7,18 +7,13 @@ pub mod terminal;
async fn main() -> anyhow::Result<()> { async fn main() -> anyhow::Result<()> {
let engine = engine::Engine::new().await?; let engine = engine::Engine::new().await?;
let terminal_server = engine.spawn_terminal_server().await;
let broadcaster = engine.spawn_broadcaster().await; let broadcaster = engine.spawn_broadcaster().await;
engine.strategy.spawn().await; engine.strategy.spawn().await;
engine.run_engine().await?; engine.terminal_server.run().await?;
let (terminal_server, broadcaster) = tokio::join!(terminal_server, broadcaster); broadcaster.await??;
terminal_server??;
broadcaster??;
Ok(()) Ok(())
} }
+9 -10
View File
@@ -1,16 +1,14 @@
use std::marker::PhantomData; use std::marker::PhantomData;
use pulse_wire::PulseWire; use pulse_sdk::map_postcard_err;
use serde::{Deserialize, Serialize};
use serde::Deserialize;
use tokio::{ use tokio::{
io::{AsyncReadExt, AsyncWriteExt}, io::{AsyncReadExt, AsyncWriteExt},
process::{Child, ChildStdout}, process::{Child, ChildStdout},
sync::Mutex, sync::Mutex,
}; };
#[derive(Debug)] #[derive(Debug)]
pub struct Plugin<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> { pub struct Plugin<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> {
pub manifest: Mutex<M>, pub manifest: Mutex<M>,
pub stdout: Mutex<ChildStdout>, pub stdout: Mutex<ChildStdout>,
pub process: Mutex<Child>, pub process: Mutex<Child>,
@@ -18,7 +16,7 @@ pub struct Plugin<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> {
pub _p: (PhantomData<S>, PhantomData<R>), pub _p: (PhantomData<S>, PhantomData<R>),
} }
impl<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> Plugin<S, R, M> { impl<S: Serialize, R: for<'de> Deserialize<'de>, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
pub fn new(mut child: Child, manifest: M) -> Self { pub fn new(mut child: Child, manifest: M) -> Self {
Self { Self {
stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")), stdout: Mutex::new(child.stdout.take().expect("Failed to obtain child stdout")),
@@ -45,19 +43,20 @@ impl<S: PulseWire, R: PulseWire, M: for<'de> Deserialize<'de>> Plugin<S, R, M> {
stdout.read_exact(&mut buffer).await?; stdout.read_exact(&mut buffer).await?;
Ok(Some(R::from_com(&mut buffer))) Ok(Some(map_postcard_err(postcard::from_bytes(&buffer))?))
} }
pub async fn send(&self, msg: &S) -> tokio::io::Result<()> { pub async fn send(&self, msg: &S) -> tokio::io::Result<()> {
self.send_raw(&msg.to_com()).await self.send_raw(&map_postcard_err(postcard::to_allocvec(msg))?)
.await
} }
pub async fn send_raw(&self, msg: &[u8]) -> tokio::io::Result<()> { pub async fn send_raw(&self, msg: &[u8]) -> tokio::io::Result<()> {
let mut process = self.process.lock().await; let mut process = self.process.lock().await;
let stdin = process.stdin.as_mut().unwrap(); let stdin = process.stdin.as_mut().unwrap();
stdin.write(&msg.len().to_le_bytes()).await?; stdin.write_all(&msg.len().to_le_bytes()).await?;
stdin.write(msg).await?; stdin.write_all(msg).await?;
stdin.flush().await?; stdin.flush().await?;
Ok(()) Ok(())
+11 -11
View File
@@ -1,5 +1,5 @@
use crate::engine::Engine; use crate::engine::Engine;
use pulse_wire::prelude::*; use pulse_sdk::{map_postcard_err, prelude::*};
use std::{ use std::{
collections::HashMap, collections::HashMap,
sync::{Arc, Weak}, sync::{Arc, Weak},
@@ -30,7 +30,7 @@ impl TerminalServer {
} }
pub async fn run(self: &Arc<Self>) -> tokio::io::Result<()> { pub async fn run(self: &Arc<Self>) -> tokio::io::Result<()> {
let path = pulse_wire::server_path(); let path = pulse_sdk::server_path();
if path.exists() { if path.exists() {
tokio::fs::remove_file(&path).await?; tokio::fs::remove_file(&path).await?;
@@ -64,7 +64,7 @@ impl TerminalServer {
self.send_to( self.send_to(
id, id,
pulse_wire::terminal::TerminalServerMessage::StrategyUpdated(Strategy { pulse_sdk::terminal::TerminalServerMessage::StrategyUpdated(Strategy {
strategy: engine.strategy.strategy.manifest.lock().await.clone(), strategy: engine.strategy.strategy.manifest.lock().await.clone(),
risk: engine.strategy.risk.manifest.lock().await.clone(), risk: engine.strategy.risk.manifest.lock().await.clone(),
mode: Mode::Auto, mode: Mode::Auto,
@@ -76,7 +76,7 @@ impl TerminalServer {
self.send_to( self.send_to(
id, id,
pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()), pulse_sdk::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()),
) )
.await?; .await?;
@@ -104,7 +104,7 @@ impl TerminalServer {
reader.read_exact(&mut buffer).await?; reader.read_exact(&mut buffer).await?;
match TerminalClientMessage::from_com(&mut buffer) { match map_postcard_err(postcard::from_bytes(&buffer))? {
TerminalClientMessage::ExecuteCommand(command) => { TerminalClientMessage::ExecuteCommand(command) => {
let command = command.as_str(); let command = command.as_str();
@@ -124,9 +124,9 @@ impl TerminalServer {
pub async fn broadcast( pub async fn broadcast(
self: &Arc<Self>, self: &Arc<Self>,
message: pulse_wire::terminal::TerminalServerMessage, message: pulse_sdk::terminal::TerminalServerMessage,
) -> tokio::io::Result<()> { ) -> tokio::io::Result<()> {
let msg = message.to_com(); let msg = map_postcard_err(postcard::to_allocvec(&message))?;
let mut clients = self.clients.lock().await; let mut clients = self.clients.lock().await;
@@ -149,7 +149,7 @@ impl TerminalServer {
pub async fn send_to( pub async fn send_to(
self: &Arc<Self>, self: &Arc<Self>,
id: &usize, id: &usize,
message: pulse_wire::terminal::TerminalServerMessage, message: pulse_sdk::terminal::TerminalServerMessage,
) -> tokio::io::Result<()> { ) -> tokio::io::Result<()> {
Self::send_to_client( Self::send_to_client(
self.clients.lock().await.get_mut(id).ok_or_else(|| { self.clients.lock().await.get_mut(id).ok_or_else(|| {
@@ -158,7 +158,7 @@ impl TerminalServer {
format!("Client({id}) does not exist"), format!("Client({id}) does not exist"),
) )
})?, })?,
&message.to_com(), &map_postcard_err(postcard::to_allocvec(&message))?,
) )
.await .await
} }
@@ -185,14 +185,14 @@ impl TerminalServer {
self.logs.lock().await.push(log.clone()); self.logs.lock().await.push(log.clone());
self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog(log)) self.broadcast(pulse_sdk::terminal::TerminalServerMessage::AddLog(log))
.await .await
} }
pub async fn log_raw(self: &Arc<Self>, log: EventLog) -> tokio::io::Result<()> { pub async fn log_raw(self: &Arc<Self>, log: EventLog) -> tokio::io::Result<()> {
self.logs.lock().await.push(log.clone()); self.logs.lock().await.push(log.clone());
self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog(log)) self.broadcast(pulse_sdk::terminal::TerminalServerMessage::AddLog(log))
.await .await
} }
+1 -1
View File
@@ -16,7 +16,7 @@ impl PulseTradeApp {
.sock .sock
.as_mut() .as_mut()
.unwrap() .unwrap()
.send(pulse_wire::terminal::TerminalClientMessage::ExecuteCommand( .send(pulse_sdk::terminal::TerminalClientMessage::ExecuteCommand(
command.to_string(), command.to_string(),
)) ))
.await .await
+1 -1
View File
@@ -1,4 +1,4 @@
use pulse_wire::prelude::*; use pulse_sdk::prelude::*;
pub trait Formatted { pub trait Formatted {
fn get_formatted(&self) -> Vec<String>; fn get_formatted(&self) -> Vec<String>;
+3 -2
View File
@@ -21,7 +21,7 @@ use pulse_ui::{
use crate::formatting::{Formatted, apply_padding}; use crate::formatting::{Formatted, apply_padding};
use pulse_wire::prelude::*; use pulse_sdk::prelude::*;
pub struct PulseTradeApp { pub struct PulseTradeApp {
sock: Option<terminal::TerminalClient>, sock: Option<terminal::TerminalClient>,
@@ -144,7 +144,8 @@ impl App for PulseTradeApp {
( (
LayoutItem::Widget(Size::Flex(1)), LayoutItem::Widget(Size::Flex(1)),
Box::new( Box::new(
advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.strategy).await, advanced_option_draw(&self.scroll, 3, "CONFIGURATION", &self.strategy)
.await,
), ),
), ),
( (
+78 -56
View File
@@ -1,4 +1,5 @@
use pulse_wire::prelude::*; use pulse_sdk::{map_postcard_err, prelude::*};
use pulse_ui::state::State;
use tokio::{ use tokio::{
io::{AsyncReadExt, AsyncWriteExt}, io::{AsyncReadExt, AsyncWriteExt},
@@ -27,9 +28,9 @@ impl TerminalClient {
pub async fn send( pub async fn send(
&mut self, &mut self,
message: pulse_wire::terminal::TerminalClientMessage, message: pulse_sdk::terminal::TerminalClientMessage,
) -> tokio::io::Result<()> { ) -> tokio::io::Result<()> {
let msg = message.to_com(); let msg = map_postcard_err(postcard::to_allocvec(&message))?;
self.writer.write(&msg.len().to_le_bytes()).await?; self.writer.write(&msg.len().to_le_bytes()).await?;
self.writer.write(&msg).await?; self.writer.write(&msg).await?;
self.writer.flush().await?; self.writer.flush().await?;
@@ -42,7 +43,7 @@ impl TerminalClient {
std::mem::swap(&mut self.reader, &mut reader); std::mem::swap(&mut self.reader, &mut reader);
let mut reader = reader.expect("Reader failed to swap"); let reader = reader.expect("Reader failed to swap");
let watch_list = app.watch_list.clone(); let watch_list = app.watch_list.clone();
let active_positions = app.active_positions.clone(); let active_positions = app.active_positions.clone();
@@ -52,61 +53,82 @@ impl TerminalClient {
let status = app.status.clone(); let status = app.status.clone();
let inspect = app.inspect.clone(); let inspect = app.inspect.clone();
tokio::spawn(async move { tokio::spawn(Self::run_client(
loop { reader,
let mut len_buf = [0u8; size_of::<usize>()]; watch_list,
reader active_positions,
.read_exact(&mut len_buf) logs,
.await signals,
.expect("Failed to get header length"); market_overview,
status,
let len = usize::from_le_bytes(len_buf); inspect,
));
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::StrategyUpdated(v) => {
*market_overview.lock().await = Some(v);
}
TerminalServerMessage::SignalsUpdated(v) => {
*signals.lock().await = v;
}
TerminalServerMessage::Inspect(v) => {
*inspect.lock().await = v;
}
TerminalServerMessage::StatusUpdated(v) => {
*status.lock().await = Some(v);
}
TerminalServerMessage::SetLogs(v) => {
*logs.lock().await = v;
}
TerminalServerMessage::AddLog(v) => {
logs.lock().await.push(v);
}
}
}
});
app.sock = Some(self); app.sock = Some(self);
app app
} }
pub async fn run_client(
mut reader: OwnedReadHalf,
watch_list: State<Vec<MarketItem>>,
active_positions: State<Vec<Position>>,
logs: State<Vec<EventLog>>,
signals: State<Vec<Signal>>,
market_overview: State<Option<Strategy>>,
status: State<Option<Status>>,
inspect: State<InspectTarget>,
) -> tokio::io::Result<()> {
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 map_postcard_err(postcard::from_bytes(&buffer))? {
TerminalServerMessage::WatchListUpdated(v) => {
*watch_list.lock().await = v;
}
TerminalServerMessage::PositionsUpdated(v) => {
*active_positions.lock().await = v;
}
TerminalServerMessage::StrategyUpdated(v) => {
*market_overview.lock().await = Some(v);
}
TerminalServerMessage::SignalsUpdated(v) => {
*signals.lock().await = v;
}
TerminalServerMessage::Inspect(v) => {
*inspect.lock().await = v;
}
TerminalServerMessage::StatusUpdated(v) => {
*status.lock().await = Some(v);
}
TerminalServerMessage::SetLogs(v) => {
*logs.lock().await = v;
}
TerminalServerMessage::AddLog(v) => {
logs.lock().await.push(v);
}
}
Ok(())
}
} }