Working logger
This commit is contained in:
Generated
+1
@@ -3887,6 +3887,7 @@ dependencies = [
|
|||||||
"hypersdk",
|
"hypersdk",
|
||||||
"pulse-ui",
|
"pulse-ui",
|
||||||
"pulse-wire",
|
"pulse-wire",
|
||||||
|
"rand 0.8.7",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
"tokio",
|
"tokio",
|
||||||
]
|
]
|
||||||
|
|||||||
+9
-1
@@ -8,10 +8,18 @@ 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", "time"] }
|
tokio = { workspace = true, features = [
|
||||||
|
"rt-multi-thread",
|
||||||
|
"macros",
|
||||||
|
"net",
|
||||||
|
"fs",
|
||||||
|
"io-util",
|
||||||
|
"time",
|
||||||
|
] }
|
||||||
crossterm = { workspace = true }
|
crossterm = { workspace = true }
|
||||||
hypersdk = "0.2.14"
|
hypersdk = "0.2.14"
|
||||||
serde_json = "1"
|
serde_json = "1"
|
||||||
|
rand = "0.8.7"
|
||||||
|
|
||||||
[workspace]
|
[workspace]
|
||||||
members = ["pulse-macros", "pulse-ui", "pulse-wire"]
|
members = ["pulse-macros", "pulse-ui", "pulse-wire"]
|
||||||
|
|||||||
+68
-30
@@ -1,4 +1,4 @@
|
|||||||
use std::sync::Arc;
|
use std::{collections::HashMap, sync::Arc};
|
||||||
|
|
||||||
use pulse_wire::{
|
use pulse_wire::{
|
||||||
PulseWire,
|
PulseWire,
|
||||||
@@ -15,14 +15,14 @@ use tokio::{
|
|||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub struct TerminalServer {
|
pub struct TerminalServer {
|
||||||
clients: Mutex<Vec<OwnedWriteHalf>>,
|
clients: Mutex<HashMap<usize, OwnedWriteHalf>>,
|
||||||
logs: Mutex<Vec<EventLog>>,
|
logs: Mutex<Vec<EventLog>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TerminalServer {
|
impl TerminalServer {
|
||||||
pub fn new() -> Arc<Self> {
|
pub fn new() -> Arc<Self> {
|
||||||
Arc::new(Self {
|
Arc::new(Self {
|
||||||
clients: Mutex::new(Vec::new()),
|
clients: Mutex::new(HashMap::new()),
|
||||||
logs: Mutex::new(Vec::new()),
|
logs: Mutex::new(Vec::new()),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -43,19 +43,31 @@ impl TerminalServer {
|
|||||||
|
|
||||||
let (reader, writer) = stream.into_split();
|
let (reader, writer) = stream.into_split();
|
||||||
|
|
||||||
self.clients.lock().await.push(writer);
|
let id = rand::random();
|
||||||
|
|
||||||
|
self.clients.lock().await.insert(id, writer);
|
||||||
|
|
||||||
let s = self.clone();
|
let s = self.clone();
|
||||||
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
if let Err(err) = s.handle_client(reader).await {
|
if let Err(err) = s.handle_client(&id, reader).await {
|
||||||
eprintln!("Terminal connection error: {err}");
|
eprintln!("Terminal connection error: {err}");
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn handle_client(self: &Arc<Self>, mut reader: OwnedReadHalf) -> tokio::io::Result<()> {
|
async fn handle_client(
|
||||||
|
self: &Arc<Self>,
|
||||||
|
id: &usize,
|
||||||
|
mut reader: OwnedReadHalf,
|
||||||
|
) -> tokio::io::Result<()> {
|
||||||
|
self.send_to(
|
||||||
|
id,
|
||||||
|
pulse_wire::terminal::TerminalServerMessage::SetLogs(self.logs.lock().await.clone()),
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
let mut len_buf = [0u8; size_of::<usize>()];
|
let mut len_buf = [0u8; size_of::<usize>()];
|
||||||
let size = reader.read_exact(&mut len_buf).await?;
|
let size = reader.read_exact(&mut len_buf).await?;
|
||||||
@@ -82,13 +94,10 @@ impl TerminalServer {
|
|||||||
|
|
||||||
match command {
|
match command {
|
||||||
_ => {
|
_ => {
|
||||||
self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog(
|
self.error(
|
||||||
pulse_wire::terminal::EventLog {
|
"Command executor",
|
||||||
kind: pulse_wire::terminal::LogKind::Err,
|
&format!("Command '{}' not found", command),
|
||||||
name: "Command executor".to_string(),
|
)
|
||||||
message: format!("Command '{}' not found", command),
|
|
||||||
},
|
|
||||||
))
|
|
||||||
.await?;
|
.await?;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -107,49 +116,78 @@ impl TerminalServer {
|
|||||||
|
|
||||||
let mut clients = self.clients.lock().await;
|
let mut clients = self.clients.lock().await;
|
||||||
|
|
||||||
for i in (0..clients.len()).rev() {
|
let mut remove_clients = Vec::new();
|
||||||
if let Err(e) = Self::send_to_client(&mut clients, i, &msg).await {
|
|
||||||
clients.remove(i);
|
for (id, client) in clients.iter_mut() {
|
||||||
|
if let Err(e) = Self::send_to_client(client, &msg).await {
|
||||||
|
remove_clients.push(*id);
|
||||||
println!("{e:?}");
|
println!("{e:?}");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
for id in remove_clients {
|
||||||
|
clients.remove(&id);
|
||||||
|
}
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn send_to_client(
|
pub async fn send_to(
|
||||||
clients: &mut tokio::sync::MutexGuard<'_, Vec<OwnedWriteHalf>>,
|
self: &Arc<Self>,
|
||||||
i: usize,
|
id: &usize,
|
||||||
msg: &[u8],
|
message: pulse_wire::terminal::TerminalServerMessage,
|
||||||
) -> tokio::io::Result<()> {
|
) -> tokio::io::Result<()> {
|
||||||
clients[i].write(&msg.len().to_le_bytes()).await?;
|
Self::send_to_client(
|
||||||
clients[i].write(msg).await?;
|
self.clients.lock().await.get_mut(id).ok_or_else(|| {
|
||||||
clients[i].flush().await?;
|
tokio::io::Error::new(
|
||||||
|
std::io::ErrorKind::Other,
|
||||||
|
format!("Client({id}) does not exist"),
|
||||||
|
)
|
||||||
|
})?,
|
||||||
|
&message.to_com(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn send_to_client(client: &mut OwnedWriteHalf, msg: &[u8]) -> tokio::io::Result<()> {
|
||||||
|
client.write(&msg.len().to_le_bytes()).await?;
|
||||||
|
client.write(msg).await?;
|
||||||
|
client.flush().await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn log(self: &Arc<Self>, kind: LogKind, name: &str, message: &str) {
|
pub async fn log(
|
||||||
self.logs.lock().await.push(EventLog {
|
self: &Arc<Self>,
|
||||||
|
kind: LogKind,
|
||||||
|
name: &str,
|
||||||
|
message: &str,
|
||||||
|
) -> tokio::io::Result<()> {
|
||||||
|
let log = EventLog {
|
||||||
kind,
|
kind,
|
||||||
name: name.to_string(),
|
name: name.to_string(),
|
||||||
message: message.to_string(),
|
message: message.to_string(),
|
||||||
});
|
};
|
||||||
|
|
||||||
|
self.logs.lock().await.push(log.clone());
|
||||||
|
|
||||||
|
self.broadcast(pulse_wire::terminal::TerminalServerMessage::AddLog(log))
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn info(self: &Arc<Self>, name: &str, message: &str) {
|
pub async fn info(self: &Arc<Self>, name: &str, message: &str) -> tokio::io::Result<()> {
|
||||||
self.log(LogKind::Info, name, message).await
|
self.log(LogKind::Info, name, message).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn warn(self: &Arc<Self>, name: &str, message: &str) {
|
pub async fn warn(self: &Arc<Self>, name: &str, message: &str) -> tokio::io::Result<()> {
|
||||||
self.log(LogKind::Warn, name, message).await
|
self.log(LogKind::Warn, name, message).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn error(self: &Arc<Self>, name: &str, message: &str) {
|
pub async fn error(self: &Arc<Self>, name: &str, message: &str) -> tokio::io::Result<()> {
|
||||||
self.log(LogKind::Err, name, message).await
|
self.log(LogKind::Err, name, message).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn debug(self: &Arc<Self>, name: &str, message: &str) {
|
pub async fn debug(self: &Arc<Self>, name: &str, message: &str) -> tokio::io::Result<()> {
|
||||||
self.log(LogKind::Debug, name, message).await
|
self.log(LogKind::Debug, name, message).await
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user