[Feat] Split Daemon & TUI
This commit is contained in:
parent
f82500ea7d
commit
36a70e82a0
35 changed files with 970 additions and 239 deletions
15
iota-daemon-lib/Cargo.toml
Normal file
15
iota-daemon-lib/Cargo.toml
Normal file
|
|
@ -0,0 +1,15 @@
|
|||
[package]
|
||||
name = "iota-daemon-lib"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
iota-ipc = { path = "../iota-ipc" }
|
||||
iota-logger = { path = "../iota-logger" }
|
||||
iota-state = { path = "../iota-state" }
|
||||
iota-storage = { path = "../iota-storage" }
|
||||
iota-util = { path = "../iota-util" }
|
||||
omikron-connector = { path = "../omikron-connector" }
|
||||
mtp = { git = "https://git.methanium.net/Methanium/mtp.git" }
|
||||
sysinfo = "0.38.3"
|
||||
tokio = { version = "1.50.0", features = ["full"] }
|
||||
114
iota-daemon-lib/src/command_router.rs
Normal file
114
iota-daemon-lib/src/command_router.rs
Normal file
|
|
@ -0,0 +1,114 @@
|
|||
use crate::DaemonRuntime;
|
||||
use iota_ipc::DaemonMessage;
|
||||
use iota_logger::{log, log_command};
|
||||
use iota_storage::users::user_manager;
|
||||
use iota_storage::util::config_util::modify_config;
|
||||
use mtp::codec::{CommunicationType, CommunicationValue};
|
||||
use omikron_connector::omikron_connection::OMIKRON_CONNECTION;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct CommandRouter {
|
||||
runtime: Arc<DaemonRuntime>,
|
||||
}
|
||||
|
||||
impl CommandRouter {
|
||||
pub fn new(runtime: Arc<DaemonRuntime>) -> Self {
|
||||
Self { runtime }
|
||||
}
|
||||
|
||||
pub async fn route(&self, seq: u64, line: String) -> DaemonMessage {
|
||||
log_command!("{}", line);
|
||||
let result = self.execute(&line).await;
|
||||
DaemonMessage::CommandResult {
|
||||
seq,
|
||||
success: result.is_ok(),
|
||||
message: result.unwrap_or_else(|error| error),
|
||||
}
|
||||
}
|
||||
|
||||
async fn execute(&self, line: &str) -> Result<String, String> {
|
||||
let parts = line
|
||||
.trim_start_matches('/')
|
||||
.split_whitespace()
|
||||
.collect::<Vec<_>>();
|
||||
match parts.as_slice() {
|
||||
["tasks"] => Ok(self
|
||||
.runtime
|
||||
.state
|
||||
.active_tasks
|
||||
.iter()
|
||||
.map(|task| task.to_string())
|
||||
.collect::<Vec<_>>()
|
||||
.join(", ")),
|
||||
["help"] => Ok(
|
||||
"Available commands: tasks, ping, user, reconnect, regenerate, reload, shutdown"
|
||||
.into(),
|
||||
),
|
||||
["ping"] => self.ping(20).await,
|
||||
["ping", seconds] => self.ping(seconds.parse::<u64>().unwrap_or(20)).await,
|
||||
["user", "add", username] => {
|
||||
let (user, _) = omikron_connector::user_ops::create_user(username).await;
|
||||
user.map(|user| format!("Created user {}", user.user_id))
|
||||
.ok_or_else(|| "User creation failed".into())
|
||||
}
|
||||
["user", "remove", username] => {
|
||||
let user = user_manager::get_user_by_username(username)
|
||||
.ok_or_else(|| "Username does not exist".to_string())?;
|
||||
let message = CommunicationValue::new(CommunicationType::DeleteUser)
|
||||
.with_sender(user.user_id as u64);
|
||||
OMIKRON_CONNECTION
|
||||
.send_message(&message)
|
||||
.await
|
||||
.map_err(|error| error.to_string())?;
|
||||
user_manager::remove_user(user.user_id);
|
||||
Ok(format!("Removed user {}", user.user_id))
|
||||
}
|
||||
["user", "list"] => Ok(user_manager::get_users()
|
||||
.into_iter()
|
||||
.map(|user| format!("{} ({})", user.username, user.user_id))
|
||||
.collect::<Vec<_>>()
|
||||
.join("\n")),
|
||||
["reconnect"] => {
|
||||
OMIKRON_CONNECTION.reconnect().await;
|
||||
Ok("Reconnected to Omikron server".into())
|
||||
}
|
||||
["regenerate", "keys"] => {
|
||||
modify_config(|config| {
|
||||
config.public_key = None;
|
||||
config.private_key = None;
|
||||
config.iota_id = None;
|
||||
});
|
||||
OMIKRON_CONNECTION.reconnect().await;
|
||||
Ok("Key pair regenerated and Omikron reconnection requested".into())
|
||||
}
|
||||
["reload"] | ["restart"] => {
|
||||
*self.runtime.state.reload.write().await = true;
|
||||
*self.runtime.state.shutdown.write().await = true;
|
||||
Ok("Daemon restart requested".into())
|
||||
}
|
||||
["shutdown"] | ["stop"] => {
|
||||
*self.runtime.state.shutdown.write().await = true;
|
||||
Ok("Daemon shutdown requested".into())
|
||||
}
|
||||
_ => Err("Unknown command".into()),
|
||||
}
|
||||
}
|
||||
|
||||
async fn ping(&self, seconds: u64) -> Result<String, String> {
|
||||
let response = OMIKRON_CONNECTION
|
||||
.await_response(
|
||||
&CommunicationValue::new(CommunicationType::Ping),
|
||||
Some(Duration::from_secs(seconds)),
|
||||
)
|
||||
.await;
|
||||
match response {
|
||||
Ok(value) => {
|
||||
log!("{}", iota_logger::format_cv(&value));
|
||||
Ok("Ping response received".into())
|
||||
}
|
||||
Err(error) => Err(format!("Ping error: {error:?}")),
|
||||
}
|
||||
}
|
||||
}
|
||||
72
iota-daemon-lib/src/daemon_state.rs
Normal file
72
iota-daemon-lib/src/daemon_state.rs
Normal file
|
|
@ -0,0 +1,72 @@
|
|||
use iota_ipc::StateSnapshot;
|
||||
use iota_state::DaemonState;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use sysinfo::{RefreshKind, System};
|
||||
|
||||
/* This wrapper exposes daemon state as IPC-safe snapshots while preserving a
|
||||
* single owned state instance for all daemon subsystems. */
|
||||
#[derive(Clone, Default)]
|
||||
pub struct DaemonRuntime {
|
||||
pub state: Arc<DaemonState>,
|
||||
}
|
||||
|
||||
impl DaemonRuntime {
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
state: Arc::new(DaemonState::new()),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn snapshot(&self) -> StateSnapshot {
|
||||
let state = self
|
||||
.state
|
||||
.app
|
||||
.lock()
|
||||
.unwrap_or_else(|error| error.into_inner());
|
||||
StateSnapshot {
|
||||
cpu: state.cpu.clone(),
|
||||
ram: state.ram.clone(),
|
||||
ping: state.ping.clone(),
|
||||
net_up: state.net_up.clone(),
|
||||
net_down: state.net_down.clone(),
|
||||
sys_info: state.sys_info.clone(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn spawn_system_monitor(&self) {
|
||||
let runtime = self.clone();
|
||||
tokio::spawn(async move {
|
||||
runtime.state.active_tasks.insert("System monitor".into());
|
||||
let mut system = System::new_with_specifics(RefreshKind::everything());
|
||||
let mut counter = 0.0;
|
||||
loop {
|
||||
if *runtime.state.shutdown.read().await {
|
||||
break;
|
||||
}
|
||||
system.refresh_cpu_all();
|
||||
system.refresh_memory();
|
||||
let cpu = system.global_cpu_usage() as f64;
|
||||
let total_memory = system.total_memory();
|
||||
let ram = if total_memory == 0 {
|
||||
0.0
|
||||
} else {
|
||||
system.used_memory() as f64 / total_memory as f64 * 100.0
|
||||
};
|
||||
{
|
||||
let mut state = runtime
|
||||
.state
|
||||
.app
|
||||
.lock()
|
||||
.unwrap_or_else(|error| error.into_inner());
|
||||
state.push_cpu((counter, cpu));
|
||||
state.push_ram((counter, ram));
|
||||
state.sys_info = format!("CPU: {cpu:.1}% RAM: {ram:.1}%");
|
||||
}
|
||||
counter += 1.0;
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
}
|
||||
runtime.state.active_tasks.remove("System monitor");
|
||||
});
|
||||
}
|
||||
}
|
||||
113
iota-daemon-lib/src/ipc_server.rs
Normal file
113
iota-daemon-lib/src/ipc_server.rs
Normal file
|
|
@ -0,0 +1,113 @@
|
|||
use crate::{CommandRouter, DaemonRuntime};
|
||||
use iota_ipc::{ClientMessage, DaemonMessage, read_msg, write_msg};
|
||||
use std::io::Result;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use std::{env, os::fd::FromRawFd};
|
||||
use tokio::net::{UnixListener, UnixStream};
|
||||
use tokio::sync::broadcast;
|
||||
|
||||
pub struct IpcServer {
|
||||
path: PathBuf,
|
||||
runtime: Arc<DaemonRuntime>,
|
||||
messages: broadcast::Sender<DaemonMessage>,
|
||||
}
|
||||
|
||||
impl IpcServer {
|
||||
pub fn new(
|
||||
path: impl Into<PathBuf>,
|
||||
runtime: Arc<DaemonRuntime>,
|
||||
messages: broadcast::Sender<DaemonMessage>,
|
||||
) -> Self {
|
||||
Self {
|
||||
path: path.into(),
|
||||
runtime,
|
||||
messages,
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn run(self) -> Result<()> {
|
||||
if let Some(parent) = self.path.parent() {
|
||||
tokio::fs::create_dir_all(parent).await?;
|
||||
}
|
||||
let listener = match activated_listener()? {
|
||||
Some(listener) => listener,
|
||||
None => {
|
||||
remove_stale_socket(&self.path).await?;
|
||||
UnixListener::bind(&self.path)?
|
||||
}
|
||||
};
|
||||
loop {
|
||||
let (stream, _) = listener.accept().await?;
|
||||
let runtime = self.runtime.clone();
|
||||
let messages = self.messages.clone();
|
||||
tokio::spawn(async move {
|
||||
let _ = handle_client(stream, runtime, messages).await;
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/* systemd hands the first socket-activated file descriptor to the service as
|
||||
* descriptor 3. Manual launches continue to bind the configured socket path. */
|
||||
fn activated_listener() -> Result<Option<UnixListener>> {
|
||||
let listen_fds = env::var("LISTEN_FDS")
|
||||
.ok()
|
||||
.and_then(|value| value.parse::<u32>().ok());
|
||||
let listen_pid = env::var("LISTEN_PID")
|
||||
.ok()
|
||||
.and_then(|value| value.parse::<u32>().ok());
|
||||
if listen_fds != Some(1) || listen_pid != Some(std::process::id()) {
|
||||
return Ok(None);
|
||||
}
|
||||
let listener = unsafe { std::os::unix::net::UnixListener::from_raw_fd(3) };
|
||||
UnixListener::from_std(listener).map(Some)
|
||||
}
|
||||
|
||||
async fn remove_stale_socket(path: &Path) -> Result<()> {
|
||||
match tokio::fs::symlink_metadata(path).await {
|
||||
Ok(_) => tokio::fs::remove_file(path).await,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
|
||||
Err(error) => Err(error),
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_client(
|
||||
stream: UnixStream,
|
||||
runtime: Arc<DaemonRuntime>,
|
||||
messages: broadcast::Sender<DaemonMessage>,
|
||||
) -> Result<()> {
|
||||
let (mut reader, mut writer) = stream.into_split();
|
||||
let mut outgoing = messages.subscribe();
|
||||
let initial = DaemonMessage::StateUpdate(runtime.snapshot());
|
||||
write_msg(&mut writer, &initial).await?;
|
||||
let writer_task = tokio::spawn(async move {
|
||||
while let Ok(message) = outgoing.recv().await {
|
||||
if write_msg(&mut writer, &message).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
let router = CommandRouter::new(runtime.clone());
|
||||
loop {
|
||||
match read_msg::<_, ClientMessage>(&mut reader).await {
|
||||
Ok(ClientMessage::Command { seq, line }) => {
|
||||
let result = router.route(seq, line).await;
|
||||
let _ = messages.send(result);
|
||||
}
|
||||
Ok(ClientMessage::Subscribe) => {
|
||||
let _ = messages.send(DaemonMessage::StateUpdate(runtime.snapshot()));
|
||||
}
|
||||
Ok(ClientMessage::Ping { seq }) => {
|
||||
let _ = messages.send(DaemonMessage::Pong { seq });
|
||||
}
|
||||
Err(error) if error.kind() == std::io::ErrorKind::UnexpectedEof => break,
|
||||
Err(error) => {
|
||||
writer_task.abort();
|
||||
return Err(error);
|
||||
}
|
||||
}
|
||||
}
|
||||
writer_task.abort();
|
||||
Ok(())
|
||||
}
|
||||
8
iota-daemon-lib/src/lib.rs
Normal file
8
iota-daemon-lib/src/lib.rs
Normal file
|
|
@ -0,0 +1,8 @@
|
|||
pub mod command_router;
|
||||
pub mod daemon_state;
|
||||
pub mod ipc_server;
|
||||
pub mod log_broadcaster;
|
||||
|
||||
pub use command_router::CommandRouter;
|
||||
pub use daemon_state::DaemonRuntime;
|
||||
pub use ipc_server::IpcServer;
|
||||
21
iota-daemon-lib/src/log_broadcaster.rs
Normal file
21
iota-daemon-lib/src/log_broadcaster.rs
Normal file
|
|
@ -0,0 +1,21 @@
|
|||
use iota_ipc::{DaemonMessage, LogEntry};
|
||||
use iota_logger::subscribe;
|
||||
use tokio::sync::broadcast;
|
||||
|
||||
/* The daemon adapts logger output to the wire protocol so the logger stays
|
||||
* independent from both the socket implementation and TUI state. */
|
||||
pub fn spawn(message_tx: broadcast::Sender<DaemonMessage>) {
|
||||
let Some(mut logs) = subscribe() else {
|
||||
return;
|
||||
};
|
||||
tokio::spawn(async move {
|
||||
while let Ok(entry) = logs.recv().await {
|
||||
let _ = message_tx.send(DaemonMessage::LogEntry(LogEntry {
|
||||
timestamp_ms: entry.timestamp_ms,
|
||||
sender: entry.sender,
|
||||
message: entry.message,
|
||||
is_error: entry.is_error,
|
||||
}));
|
||||
}
|
||||
});
|
||||
}
|
||||
Loading…
Reference in a new issue