use pulse_wire::prelude::*; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::{ UnixStream, unix::{OwnedReadHalf, OwnedWriteHalf}, }, }; use crate::PulseTradeApp; pub struct TerminalClient { writer: OwnedWriteHalf, reader: Option, } impl TerminalClient { pub async fn new() -> tokio::io::Result { let (reader, writer) = UnixStream::connect(server_path()).await?.into_split(); Ok(Self { reader: Some(reader), writer, }) } pub async fn send( &mut self, message: pulse_wire::terminal::TerminalClientMessage, ) -> tokio::io::Result<()> { let msg = message.to_com(); self.writer.write(&msg.len().to_le_bytes()).await?; self.writer.write(&msg).await?; self.writer.flush().await?; Ok(()) } pub fn use_app(mut self, mut app: PulseTradeApp) -> PulseTradeApp { let mut reader = None; std::mem::swap(&mut self.reader, &mut reader); let mut reader = reader.expect("Reader failed to swap"); let watch_list = app.watch_list.clone(); let active_positions = app.active_positions.clone(); let logs = app.logs.clone(); let signals = app.signals.clone(); let market_overview = app.market_overview.clone(); let status = app.status.clone(); let inspect = app.inspect.clone(); tokio::spawn(async move { loop { let mut len_buf = [0u8; size_of::()]; reader .read_exact(&mut len_buf) .await .expect("Failed to get header length"); let len = usize::from_le_bytes(len_buf); let mut buffer = vec![0u8; len]; reader .read_exact(&mut buffer) .await .expect("Failed to read socket"); match TerminalServerMessage::from_com(&mut buffer) { TerminalServerMessage::WatchListUpdated(v) => { *watch_list.lock().await = v; } TerminalServerMessage::PositionsUpdated(v) => { *active_positions.lock().await = v; } TerminalServerMessage::OverviewUpdated(v) => { *market_overview.lock().await = 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 } }