use pulse_sdk::{map_postcard_err, prelude::*}; use pulse_ui::state::State; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::{ UnixStream, unix::{OwnedReadHalf, OwnedWriteHalf}, }, }; use crate::PulseTradeApp; pub struct TerminalClient { writer: OwnedWriteHalf, reader: Option, } impl TerminalClient { pub async fn new() -> tokio::io::Result { let (reader, writer) = UnixStream::connect(server_path()).await?.into_split(); Ok(Self { reader: Some(reader), writer, }) } pub async fn send( &mut self, message: pulse_sdk::terminal::TerminalClientMessage, ) -> tokio::io::Result<()> { let msg = map_postcard_err(postcard::to_allocvec(&message))?; 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 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.strategy.clone(); let status = app.status.clone(); let inspect = app.inspect.clone(); tokio::spawn(Self::run_client( reader, watch_list, active_positions, logs, signals, market_overview, status, inspect, )); app.sock = Some(self); app } pub async fn run_client( mut reader: OwnedReadHalf, watch_list: State>, active_positions: State>, logs: State>, signals: State>, market_overview: State>, status: State>, inspect: State, ) -> tokio::io::Result<()> { let mut len_buf = [0u8; size_of::()]; reader .read_exact(&mut len_buf) .await .expect("Failed to get header length"); let len = usize::from_le_bytes(len_buf); let mut buffer = vec![0u8; len]; reader .read_exact(&mut buffer) .await .expect("Failed to read socket"); match 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(()) } }