Terminal client
This commit is contained in:
+1
-1
@@ -8,7 +8,7 @@ pulse-ui = { workspace = true }
|
|||||||
pulse-wire = { workspace = true }
|
pulse-wire = { workspace = true }
|
||||||
|
|
||||||
chrono = "0.4.45"
|
chrono = "0.4.45"
|
||||||
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util"] }
|
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "net", "fs", "io-util", "time"] }
|
||||||
crossterm = { workspace = true }
|
crossterm = { workspace = true }
|
||||||
|
|
||||||
[workspace]
|
[workspace]
|
||||||
|
|||||||
+24
-3
@@ -2,9 +2,30 @@ pub mod terminal;
|
|||||||
|
|
||||||
#[tokio::main]
|
#[tokio::main]
|
||||||
async fn main() -> tokio::io::Result<()> {
|
async fn main() -> tokio::io::Result<()> {
|
||||||
let mut terminal_server = terminal::TerminalServer::new();
|
let terminal_server = terminal::TerminalServer::new();
|
||||||
|
|
||||||
terminal_server.run().await?;
|
{
|
||||||
|
let terminal_server = terminal_server.clone();
|
||||||
|
|
||||||
Ok(())
|
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::AddLog(
|
||||||
|
pulse_wire::terminal::EventLog {
|
||||||
|
kind: pulse_wire::terminal::LogKind::Debug,
|
||||||
|
name: "Engine".to_string(),
|
||||||
|
message: "Hello, world!".to_string(),
|
||||||
|
},
|
||||||
|
))
|
||||||
|
.await?;
|
||||||
|
println!("Ok");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+14
-10
@@ -1,3 +1,5 @@
|
|||||||
|
use std::sync::Arc;
|
||||||
|
|
||||||
use pulse_wire::PulseWire;
|
use pulse_wire::PulseWire;
|
||||||
use tokio::{
|
use tokio::{
|
||||||
io::{AsyncReadExt, AsyncWriteExt},
|
io::{AsyncReadExt, AsyncWriteExt},
|
||||||
@@ -5,20 +7,22 @@ use tokio::{
|
|||||||
UnixListener,
|
UnixListener,
|
||||||
unix::{OwnedReadHalf, OwnedWriteHalf},
|
unix::{OwnedReadHalf, OwnedWriteHalf},
|
||||||
},
|
},
|
||||||
|
sync::Mutex,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
#[derive(Debug, Clone)]
|
||||||
pub struct TerminalServer {
|
pub struct TerminalServer {
|
||||||
clients: Vec<OwnedWriteHalf>,
|
clients: Arc<Mutex<Vec<OwnedWriteHalf>>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TerminalServer {
|
impl TerminalServer {
|
||||||
pub fn new() -> Self {
|
pub fn new() -> Self {
|
||||||
Self {
|
Self {
|
||||||
clients: Vec::new(),
|
clients: Arc::new(Mutex::new(Vec::new())),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn run(&mut self) -> tokio::io::Result<()> {
|
pub async fn run(&self) -> tokio::io::Result<()> {
|
||||||
let path = pulse_wire::server_path();
|
let path = pulse_wire::server_path();
|
||||||
|
|
||||||
if path.exists() {
|
if path.exists() {
|
||||||
@@ -34,7 +38,7 @@ impl TerminalServer {
|
|||||||
|
|
||||||
let (reader, writer) = stream.into_split();
|
let (reader, writer) = stream.into_split();
|
||||||
|
|
||||||
self.clients.push(writer);
|
self.clients.lock().await.push(writer);
|
||||||
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
if let Err(err) = Self::handle_client(reader).await {
|
if let Err(err) = Self::handle_client(reader).await {
|
||||||
@@ -67,18 +71,18 @@ impl TerminalServer {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub async fn broadcast(
|
pub async fn broadcast(
|
||||||
&mut self,
|
&self,
|
||||||
message: pulse_wire::terminal::TerminalServerMessage,
|
message: pulse_wire::terminal::TerminalServerMessage,
|
||||||
) -> tokio::io::Result<()> {
|
) -> tokio::io::Result<()> {
|
||||||
let mut res = Ok(());
|
|
||||||
let msg = message.to_com();
|
let msg = message.to_com();
|
||||||
|
|
||||||
for client in &mut self.clients {
|
for i in (0..self.clients.lock().await.len()).rev() {
|
||||||
if let Err(e) = client.write(&msg).await {
|
if let Err(e) = self.clients.lock().await[i].write(&msg).await {
|
||||||
res = Err(e);
|
self.clients.lock().await.remove(i);
|
||||||
|
println!("{e:?}");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
res
|
Ok(())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+15
-2
@@ -1,5 +1,6 @@
|
|||||||
pub mod command;
|
pub mod command;
|
||||||
pub mod formatting;
|
pub mod formatting;
|
||||||
|
pub mod terminal;
|
||||||
|
|
||||||
use std::any::Any;
|
use std::any::Any;
|
||||||
|
|
||||||
@@ -270,8 +271,17 @@ impl App for PulseTradeApp {
|
|||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::main]
|
#[tokio::main]
|
||||||
async fn main() {
|
async fn main() -> tokio::io::Result<()> {
|
||||||
pulse_ui::run(|ctx| PulseTradeApp {
|
let mut client = terminal::TerminalClient::new().await?;
|
||||||
|
|
||||||
|
client
|
||||||
|
.send(pulse_wire::terminal::TerminalClientMessage::ExecuteCommand(
|
||||||
|
"Hello".to_string(),
|
||||||
|
))
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
pulse_ui::run(|ctx| {
|
||||||
|
client.use_app(PulseTradeApp {
|
||||||
command: ctx.use_state(InputState::new()),
|
command: ctx.use_state(InputState::new()),
|
||||||
watch_list: ctx.use_state(Vec::new()),
|
watch_list: ctx.use_state(Vec::new()),
|
||||||
active_positions: ctx.use_state(Vec::new()),
|
active_positions: ctx.use_state(Vec::new()),
|
||||||
@@ -291,5 +301,8 @@ async fn main() {
|
|||||||
latency: 18,
|
latency: 18,
|
||||||
}),
|
}),
|
||||||
})
|
})
|
||||||
|
})
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,104 @@
|
|||||||
|
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<()> {
|
||||||
|
self.writer.write(&message.to_com()).await?;
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn use_app(&mut self, 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 buffer = Vec::new();
|
||||||
|
|
||||||
|
let len = reader
|
||||||
|
.read(&mut buffer)
|
||||||
|
.await
|
||||||
|
.expect("Failed to read socket");
|
||||||
|
|
||||||
|
if len == 0 {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
buffer.truncate(len);
|
||||||
|
|
||||||
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user